use super::fakes::*;
use super::*;
use crate::config::PipelineConfig;
use crate::error::{ErrorClass, SourceError};
use crate::ops::RunnableChain;
use crate::pipeline::runtime::PipelineRuntime;
use crate::record::PartitionId;
use crate::sink::shard_queues;
use crate::source::LaneId;
use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
fn test_sink() -> (SinkRuntime, Arc<AtomicBool>) {
let (queues, receivers) = shard_queues(1, 8);
let drained = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&drained);
let drain: super::SinkDrainFn = Box::new(move |_budget| {
Box::pin(async move {
let _receivers = receivers;
flag.store(true, Ordering::Relaxed);
DrainReport::default()
})
});
(
SinkRuntime {
queues: vec![queues],
drain,
probe: None,
},
drained,
)
}
struct Harness {
shared: Arc<Mutex<SourceLog>>,
script: Arc<Mutex<VecDeque<Script>>>,
chain: Arc<ChainShared>,
drained: Arc<AtomicBool>,
shutdown: super::runtime::ShutdownHandle,
join: std::thread::JoinHandle<Result<ExitReport, super::runtime::StartError>>,
}
fn start(
mode_factory: impl Fn(Arc<ChainShared>, Arc<Mutex<SourceLog>>) -> FakeChain + Send + 'static,
) -> Harness {
start_with_config(test_config(1), mode_factory)
}
fn start_with_config(
config: PipelineConfig,
mode_factory: impl Fn(Arc<ChainShared>, Arc<Mutex<SourceLog>>) -> FakeChain + Send + 'static,
) -> Harness {
start_with_options(config, test_options(), mode_factory)
}
fn start_with_options(
config: PipelineConfig,
options: RuntimeOptions,
mode_factory: impl Fn(Arc<ChainShared>, Arc<Mutex<SourceLog>>) -> FakeChain + Send + 'static,
) -> Harness {
let (source, shared, script) = FakeSource::new();
let chain_shared = Arc::new(ChainShared::default());
let (sink, drained) = test_sink();
let budget = Arc::new(crate::backpressure::InflightBudget::new());
let cs = Arc::clone(&chain_shared);
let log = Arc::clone(&shared);
let runtime = PipelineRuntime::new(
config,
source,
move |_thread| {
Box::new(mode_factory(Arc::clone(&cs), Arc::clone(&log))) as Box<dyn RunnableChain>
},
sink,
budget,
)
.with_options(options);
let shutdown = runtime.shutdown_handle();
let join = std::thread::spawn(move || runtime.run());
Harness {
shared,
script,
chain: chain_shared,
drained,
shutdown,
join,
}
}
fn assign_one_lane(h: &Harness, ranges: &[std::ops::Range<i64>]) {
h.script
.lock()
.unwrap()
.push_back(Script::Assign(vec![LaneSpec {
id: LaneId(0),
partition: PartitionId(0),
batches: batches(ranges),
}]));
}
#[test]
fn commit_tick_flushes_chains_without_an_idle_lull() {
let options = RuntimeOptions {
idle_flush: Duration::from_secs(60),
..test_options()
};
let h = start_with_options(test_config(1), options, |shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20]);
wait_for("batches consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) >= 20
});
let flushes_before = h.chain.flushes.load(Ordering::Relaxed);
wait_for("commit-tick flushes", Duration::from_secs(5), || {
h.chain.flushes.load(Ordering::Relaxed) >= flushes_before + 3
});
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn happy_path_consumes_and_commits() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30]);
wait_for("all payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 30
});
wait_for("watermark committed", Duration::from_secs(5), || {
h.shared.lock().unwrap().committed.get(&PartitionId(0)) == Some(&30)
});
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert_eq!(report.final_watermarks, vec![(PartitionId(0), 30)]);
assert!(h.drained.load(Ordering::Relaxed), "sink drain must run");
let log = h.shared.lock().unwrap();
assert!(log.opened);
assert!(log.flush_commits >= 1, "shutdown must flush commits");
}
#[test]
fn revocation_drains_flushes_and_commits_in_order() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30]);
wait_for("first batch consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) >= 10
});
h.script
.lock()
.unwrap()
.push_back(Script::Revoke(vec![LaneId(0)]));
wait_for("revocation processed", Duration::from_secs(5), || {
let log = h.shared.lock().unwrap();
log.log.iter().any(|e| e == "revoke-delivered") && log.flush_commits >= 1
});
let consumed_at_revoke = h.chain.consumed.load(Ordering::Relaxed);
std::thread::sleep(Duration::from_millis(80));
assert_eq!(
h.chain.consumed.load(Ordering::Relaxed),
consumed_at_revoke,
"no consumption after lanes were revoked"
);
{
let log = h.shared.lock().unwrap();
let revoke_at = log
.log
.iter()
.position(|e| e == "revoke-delivered")
.unwrap();
let flush_after = log.log[revoke_at..].iter().position(|e| e == "flush");
let commit_after = log.log[revoke_at..].iter().position(|e| e == "commit");
assert!(
flush_after.is_some(),
"chain must flush during revocation drain"
);
assert!(
commit_after.is_some(),
"revocation must commit acknowledged offsets"
);
assert_eq!(
log.committed.get(&PartitionId(0)),
Some(&(consumed_at_revoke as i64))
);
}
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn revocation_flushes_parked_acks_before_committing() {
let mut config = test_config(1);
config.checkpoint.interval = Duration::from_secs(60);
let options = RuntimeOptions {
idle_flush: Duration::from_secs(60),
..test_options()
};
let held = Arc::new(Mutex::new(Vec::new()));
let h = start_with_options(config, options, move |shared, log| FakeChain {
shared,
log,
mode: ChainMode::HoldAcksUntilFlush {
held: Arc::clone(&held),
},
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20]);
wait_for("batches consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) >= 20
});
assert!(
h.shared.lock().unwrap().committed.is_empty(),
"acks are parked and no commit tick is due; nothing may commit \
before the revocation drain"
);
h.script
.lock()
.unwrap()
.push_back(Script::Revoke(vec![LaneId(0)]));
wait_for(
"revocation committed the drained acks",
Duration::from_secs(5),
|| h.shared.lock().unwrap().committed.get(&PartitionId(0)) == Some(&20),
);
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn blocked_chain_pauses_then_resumes() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::BlockOnce(AtomicBool::new(false)),
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20]);
wait_for("all payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 20
});
wait_for("pause and resume observed", Duration::from_secs(5), || {
let log = h.shared.lock().unwrap();
!log.pauses.is_empty() && !log.resumes.is_empty()
});
{
let log = h.shared.lock().unwrap();
assert_eq!(log.pauses[0], vec![LaneId(0)]);
assert_eq!(log.resumes[0], vec![LaneId(0)]);
}
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn shutdown_during_permanently_blocked_batch_exits_promptly_and_fails_the_batch() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::BlockForever,
batches_seen: 0,
});
assign_one_lane(&h, std::slice::from_ref(&(0..10)));
std::thread::sleep(Duration::from_millis(100));
let begun = std::time::Instant::now();
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert!(
begun.elapsed() < Duration::from_secs(5),
"shutdown must not wait for the barrier deadline"
);
assert_eq!(h.chain.consumed.load(Ordering::Relaxed), 0);
let log = h.shared.lock().unwrap();
assert!(
log.committed.is_empty(),
"a blocked, abandoned batch must not commit: {:?}",
log.committed
);
}
#[test]
fn fatal_chain_fails_pipeline_and_stalls_watermark() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::FatalAtBatch(2),
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30]);
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("pipeline must fail");
};
assert_eq!(failure.component, "fake-chain");
let log = h.shared.lock().unwrap();
let committed = log.committed.get(&PartitionId(0)).copied();
assert!(
committed == Some(10) || committed.is_none(),
"watermark must not pass the failed batch (got {committed:?})"
);
assert!(
h.drained.load(Ordering::Relaxed),
"failure path drains sink"
);
}
#[test]
fn panicking_chain_fails_pipeline() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::PanicAtBatch(2),
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20]);
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("pipeline must fail");
};
assert!(failure.reason.contains("panicked"), "{}", failure.reason);
let log = h.shared.lock().unwrap();
let committed = log.committed.get(&PartitionId(0)).copied();
assert!(committed == Some(10) || committed.is_none());
}
#[test]
fn permanent_watermark_stall_fails_pipeline_as_checkpoint() {
let mut cfg = test_config(1);
cfg.checkpoint.stalled_fail_after = Duration::from_millis(50);
let h = start_with_config(cfg, |shared, log| FakeChain {
shared,
log,
mode: ChainMode::FailAckAtBatch(1),
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30]);
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("a permanent stall must fail the pipeline");
};
assert_eq!(failure.component, "checkpoint");
assert!(failure.reason.contains("stalled"), "{}", failure.reason);
}
#[test]
fn pending_batch_limit_pauses_then_resumes_lanes() {
let held: Arc<Mutex<Vec<crate::checkpoint::AckRef>>> = Arc::new(Mutex::new(Vec::new()));
let release = Arc::new(AtomicBool::new(false));
let held_c = Arc::clone(&held);
let release_c = Arc::clone(&release);
let mut cfg = test_config(1);
cfg.checkpoint.max_pending_batches = 3;
let h = start_with_config(cfg, move |shared, log| FakeChain {
shared,
log,
mode: ChainMode::HoldAcks {
held: Arc::clone(&held_c),
release: Arc::clone(&release_c),
},
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30, 30..40, 40..50, 50..60]);
wait_for(
"controller pauses under pending pressure",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.pauses
.iter()
.any(|p| p.contains(&LaneId(0)))
},
);
release.store(true, Ordering::Relaxed);
held.lock().unwrap().clear();
wait_for(
"controller resumes after pending clears",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.resumes
.iter()
.any(|r| r.contains(&LaneId(0)))
},
);
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn a_lane_added_under_pending_pressure_starts_paused() {
let held: Arc<Mutex<Vec<crate::checkpoint::AckRef>>> = Arc::new(Mutex::new(Vec::new()));
let release = Arc::new(AtomicBool::new(false));
let held_c = Arc::clone(&held);
let release_c = Arc::clone(&release);
let mut cfg = test_config(1);
cfg.checkpoint.max_pending_batches = 3;
let h = start_with_config(cfg, move |shared, log| FakeChain {
shared,
log,
mode: ChainMode::HoldAcks {
held: Arc::clone(&held_c),
release: Arc::clone(&release_c),
},
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30, 30..40, 40..50, 50..60]);
wait_for(
"pending pressure engages on the original lane",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.pauses
.iter()
.any(|p| p.contains(&LaneId(0)))
},
);
wait_for(
"all six batches are in flight",
Duration::from_secs(5),
|| held.lock().unwrap().len() == 6,
);
held.lock().unwrap().drain(..4);
wait_for(
"the resolved batches are committed, leaving pending in the band",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.committed
.get(&PartitionId(0))
.is_some_and(|&w| w >= 40)
},
);
h.script
.lock()
.unwrap()
.push_back(Script::Add(vec![LaneSpec {
id: LaneId(1),
partition: PartitionId(1),
batches: batches(std::slice::from_ref(&(100..110))),
}]));
wait_for(
"the added lane is paused too",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.pauses
.iter()
.any(|p| p.contains(&LaneId(1)))
},
);
release.store(true, Ordering::Relaxed);
held.lock().unwrap().clear();
wait_for(
"both lanes resume once pending drains",
Duration::from_secs(5),
|| {
h.shared
.lock()
.unwrap()
.resumes
.iter()
.any(|r| r.contains(&LaneId(1)))
},
);
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
}
#[test]
fn shutdown_flushes_chain_before_exit() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, std::slice::from_ref(&(0..5)));
wait_for("payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 5
});
let flushes_before = h.chain.flushes.load(Ordering::Relaxed);
h.shutdown.trigger();
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert!(
h.chain.flushes.load(Ordering::Relaxed) > flushes_before,
"shutdown must flush the chain"
);
assert!(h.drained.load(Ordering::Relaxed));
assert_eq!(report.final_watermarks, vec![(PartitionId(0), 5)]);
}
#[test]
fn drained_source_completes_pipeline_without_external_trigger() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20, 20..30]);
wait_for("all payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 30
});
h.script.lock().unwrap().push_back(Script::Drained);
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert_eq!(report.final_watermarks, vec![(PartitionId(0), 30)]);
assert!(h.drained.load(Ordering::Relaxed), "sink drain must run");
let log = h.shared.lock().unwrap();
assert!(
log.flush_commits >= 1,
"drained completion must flush commits synchronously"
);
}
#[test]
fn drained_before_any_data_completes_with_empty_watermarks() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
h.script.lock().unwrap().push_back(Script::Drained);
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert!(report.final_watermarks.is_empty());
assert!(h.drained.load(Ordering::Relaxed), "sink drain must run");
}
#[test]
fn drained_with_unacknowledged_batches_fails_instead_of_completing() {
let held: Arc<Mutex<Vec<crate::checkpoint::AckRef>>> = Arc::new(Mutex::new(Vec::new()));
let release = Arc::new(AtomicBool::new(false));
let held_c = Arc::clone(&held);
let release_c = Arc::clone(&release);
let h = start(move |shared, log| FakeChain {
shared,
log,
mode: ChainMode::HoldAcks {
held: Arc::clone(&held_c),
release: Arc::clone(&release_c),
},
batches_seen: 0,
});
assign_one_lane(&h, &[0..10, 10..20]);
wait_for("all payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 20
});
h.script.lock().unwrap().push_back(Script::Drained);
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("drained exit with pending acks must fail, not complete");
};
assert_eq!(failure.component, "source");
assert!(
failure.reason.contains("unacknowledged"),
"{}",
failure.reason
);
assert!(report.final_watermarks.is_empty(), "nothing was committed");
drop(held);
}
#[test]
fn drained_with_failing_commit_fails_instead_of_completing() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
h.shared.lock().unwrap().fail_commits = true;
assign_one_lane(&h, &[0..10, 10..20]);
wait_for("all payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 20
});
h.script.lock().unwrap().push_back(Script::Drained);
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("drained exit with an unpersisted commit must fail, not complete");
};
assert_eq!(failure.component, "source");
assert!(
failure.reason.contains("did not persist"),
"{}",
failure.reason
);
assert!(
report.final_watermarks.is_empty(),
"nothing was durably committed"
);
}
#[test]
fn source_error_classification() {
let retryable = SourceError::Client {
class: ErrorClass::Retryable,
reason: "hiccup".into(),
};
let fatal = SourceError::Client {
class: ErrorClass::Fatal,
reason: "broken".into(),
};
assert!(matches!(
retryable,
SourceError::Client {
class: ErrorClass::Retryable,
..
}
));
assert!(matches!(
fatal,
SourceError::Client {
class: ErrorClass::Fatal,
..
}
));
}
#[test]
fn controller_panic_stops_drivers_and_fails_instead_of_hanging() {
let h = start(|shared, log| FakeChain {
shared,
log,
mode: ChainMode::BlockForever,
batches_seen: 0,
});
{
let mut script = h.script.lock().unwrap();
script.push_back(Script::Assign(vec![LaneSpec {
id: LaneId(0),
partition: PartitionId(0),
batches: batches(std::slice::from_ref(&(0..10))),
}]));
script.push_back(Script::PanicPoll);
}
let deadline = Instant::now() + Duration::from_secs(30);
while !h.join.is_finished() {
assert!(
Instant::now() < deadline,
"run() hung after the controller thread panicked"
);
std::thread::sleep(Duration::from_millis(20));
}
let report = h.join.join().unwrap().unwrap();
let ExitState::Failed(failure) = report.state else {
panic!("a controller panic must fail the run");
};
assert_eq!(failure.component, "controller");
}
#[test]
fn caller_owned_io_runtime_is_used_and_shut_down_by_run() {
let io = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.thread_name("custom-io")
.enable_all()
.build()
.expect("runtime");
struct SetOnDrop(Arc<AtomicBool>);
impl Drop for SetOnDrop {
fn drop(&mut self) {
self.0.store(true, Ordering::Relaxed);
}
}
let runtime_shut_down = Arc::new(AtomicBool::new(false));
let guard = SetOnDrop(Arc::clone(&runtime_shut_down));
io.spawn(async move {
let _guard = guard;
std::future::pending::<()>().await;
});
let (queues, receivers) = shard_queues(1, 8);
let drain_worker: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
let seen = Arc::clone(&drain_worker);
let drain: super::SinkDrainFn = Box::new(move |_budget| {
Box::pin(async move {
let _receivers = receivers;
let name = tokio::spawn(async { std::thread::current().name().map(String::from) })
.await
.expect("probe task");
*seen.lock().unwrap() = name;
DrainReport::default()
})
});
let sink = SinkRuntime {
queues: vec![queues],
drain,
probe: None,
};
let (source, _shared, _script) = FakeSource::new();
let chain_shared = Arc::new(ChainShared::default());
let log = Arc::new(Mutex::new(SourceLog::default()));
let cs = Arc::clone(&chain_shared);
let runtime = PipelineRuntime::new(
test_config(1),
source,
move |_thread| {
Box::new(FakeChain {
shared: Arc::clone(&cs),
log: Arc::clone(&log),
mode: ChainMode::Ok,
batches_seen: 0,
}) as Box<dyn RunnableChain>
},
sink,
Arc::new(crate::backpressure::InflightBudget::new()),
)
.with_options(test_options())
.with_io_runtime(io);
let shutdown = runtime.shutdown_handle();
let join = std::thread::spawn(move || runtime.run());
std::thread::sleep(Duration::from_millis(50));
shutdown.trigger();
let report = join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
assert_eq!(
drain_worker.lock().unwrap().as_deref(),
Some("custom-io"),
"drain-spawned work must run on the caller's runtime"
);
assert!(
runtime_shut_down.load(Ordering::Relaxed),
"run() must shut the caller-owned runtime down on exit"
);
}
#[test]
fn startup_error_after_driver_spawn_stops_drivers_and_returns_err() {
let occupied = std::net::TcpListener::bind("127.0.0.1:0").expect("bind probe port");
let addr = occupied.local_addr().unwrap();
let mut cfg = test_config(1);
cfg.metrics.listen = addr;
let h = start_with_config(cfg, |shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
let deadline = Instant::now() + Duration::from_secs(10);
while !h.join.is_finished() {
assert!(
Instant::now() < deadline,
"run() did not return promptly on a bind failure"
);
std::thread::sleep(Duration::from_millis(20));
}
let result = h.join.join().unwrap();
assert!(
matches!(result, Err(StartError::Io(_))),
"expected StartError::Io, got {result:?}"
);
drop(occupied);
}
#[test]
fn the_controller_owns_the_source_series_not_the_per_thread_shadows() {
let mut cfg = test_config(4); cfg.metrics.exporter = crate::config::MetricsExporter::Prometheus;
let pipeline_name = cfg.pipeline.name.clone();
let handle = crate::metrics::install(&crate::metrics::MetricsSettings {
exporter: crate::metrics::Exporter::Prometheus,
listen: cfg.metrics.listen,
per_partition_detail: cfg.metrics.per_partition_detail,
e2e_basis: crate::metrics::E2eBasis::Ingest,
})
.expect("install the exporter");
let h = start_with_config(cfg, |shared, log| FakeChain {
shared,
log,
mode: ChainMode::Ok,
batches_seen: 0,
});
assign_one_lane(&h, std::slice::from_ref(&(0..10)));
wait_for("payloads consumed", Duration::from_secs(5), || {
h.chain.consumed.load(Ordering::Relaxed) == 10
});
h.script.lock().unwrap().push_back(Script::Drained);
let report = h.join.join().unwrap().unwrap();
assert_eq!(report.state, ExitState::Completed);
let rendered = handle.render();
let needle = format!(
r#"spate_source_lag_records{{pipeline="{pipeline_name}",component="source",component_type="source",partition="0"}} {FAKE_SOURCE_LAG}"#
);
assert!(
rendered.contains(&needle),
"the source's lag must reach the exposition — the handle set the \
source was given owns the series:\n{rendered}"
);
}