use crate::config::Compression;
use crate::fetch::{FetcherParams, ObjectEntry, run_fetcher};
use crate::framer::FramerFactory;
use crate::lane::S3Lane;
use crate::metrics::S3Metrics;
use crate::offset::Position;
use crate::split::SplitDescriptor;
use serde::{Deserialize, Serialize};
use spate_core::checkpoint::AckIssuer;
use spate_core::coordination::driver::{SplitOpening, SplitSource};
use spate_core::coordination::{SplitId, SplitProgress, SplitSpec};
use spate_core::error::{ErrorClass, SourceError};
use spate_core::source::LaneId;
use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32, Ordering};
use std::time::Duration;
use tokio::sync::mpsc;
const LANE_HANDOFF_CHUNKS: usize = 4;
const UNKNOWN: i64 = i64::MIN;
const PROGRESS_STATE_VERSION: u32 = 1;
#[derive(Debug)]
pub(crate) struct SplitTracker {
terminal: AtomicI64,
objects_done: AtomicU32,
}
impl SplitTracker {
pub(crate) fn new() -> SplitTracker {
SplitTracker {
terminal: AtomicI64::new(UNKNOWN),
objects_done: AtomicU32::new(0),
}
}
pub(crate) fn set_terminal(&self, watermark: i64) {
debug_assert_ne!(watermark, UNKNOWN, "terminal watermark is a real offset");
let prev = self.terminal.swap(watermark, Ordering::Release);
debug_assert!(
prev == UNKNOWN || prev == watermark,
"a lane decides end-of-input once; conflicting terminals {prev} vs {watermark}"
);
}
pub(crate) fn terminal(&self) -> Option<i64> {
match self.terminal.load(Ordering::Acquire) {
UNKNOWN => None,
w => Some(w),
}
}
pub(crate) fn object_done(&self) {
self.objects_done.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn objects_done(&self) -> u32 {
self.objects_done.load(Ordering::Relaxed)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum PoisonKind {
NotFound,
EtagDrift,
Undecodable,
RetriesExhausted,
}
impl PoisonKind {
pub(crate) fn reason_label(self) -> &'static str {
match self {
PoisonKind::NotFound => "not_found",
PoisonKind::EtagDrift => "etag_drift",
PoisonKind::Undecodable => "undecodable",
PoisonKind::RetriesExhausted => "retries_exhausted",
}
}
}
#[derive(Debug)]
pub(crate) struct PoisonReport {
pub(crate) split: SplitId,
pub(crate) kind: PoisonKind,
pub(crate) reason: String,
}
#[derive(Debug, PartialEq, Serialize, Deserialize)]
struct ProgressState {
v: u32,
#[serde(default)]
key: Option<String>,
#[serde(default)]
etag: Option<String>,
}
impl ProgressState {
fn at(objects: &[ObjectEntry], watermark: i64) -> ProgressState {
let ordinal = Position::decode(watermark).ordinal as usize;
let entry = objects.get(ordinal);
ProgressState {
v: PROGRESS_STATE_VERSION,
key: entry.map(|e| e.key.clone()),
etag: entry.and_then(|e| e.etag.clone()),
}
}
fn encode(&self) -> Vec<u8> {
serde_json::to_vec(self).expect("progress-state serialization is infallible")
}
fn decode(bytes: &[u8]) -> Result<ProgressState, String> {
if bytes.is_empty() {
return Ok(ProgressState {
v: PROGRESS_STATE_VERSION,
key: None,
etag: None,
});
}
let state: ProgressState = serde_json::from_slice(bytes)
.map_err(|e| format!("progress state failed to decode: {e}"))?;
if state.v != PROGRESS_STATE_VERSION {
return Err(format!(
"progress state version {} is not this release's version \
{PROGRESS_STATE_VERSION}",
state.v
));
}
Ok(state)
}
}
struct SplitState {
objects: Arc<Vec<ObjectEntry>>,
lane: LaneId,
pause: Arc<AtomicBool>,
stop: Arc<AtomicBool>,
tracker: Arc<SplitTracker>,
remaining_at_open: u64,
resume_watermark: i64,
last_committed: Option<i64>,
hinted: bool,
revocation: bool,
}
pub(crate) struct SplitCtx {
store: Arc<dyn object_store::ObjectStore>,
handle: tokio::runtime::Handle,
issuer: AckIssuer,
make_framer: FramerFactory,
compression: Compression,
chunk_bytes: usize,
range_bytes: usize,
retry_base: Duration,
metrics: Option<S3Metrics>,
splits: BTreeMap<SplitId, SplitState>,
poison_tx: std::sync::mpsc::Sender<PoisonReport>,
poison_rx: std::sync::mpsc::Receiver<PoisonReport>,
}
impl std::fmt::Debug for SplitCtx {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SplitCtx")
.field("splits", &self.splits.len())
.finish_non_exhaustive()
}
}
impl SplitCtx {
#[expect(
clippy::too_many_arguments,
reason = "assembled once by the source at open"
)]
pub(crate) fn new(
store: Arc<dyn object_store::ObjectStore>,
handle: tokio::runtime::Handle,
issuer: AckIssuer,
make_framer: FramerFactory,
compression: Compression,
chunk_bytes: usize,
range_bytes: usize,
retry_base: Duration,
metrics: Option<S3Metrics>,
) -> SplitCtx {
let (poison_tx, poison_rx) = std::sync::mpsc::channel();
SplitCtx {
store,
handle,
issuer,
make_framer,
compression,
chunk_bytes,
range_bytes,
retry_base,
metrics,
splits: BTreeMap::new(),
poison_tx,
poison_rx,
}
}
pub(crate) fn drain_poison(&mut self) -> Vec<PoisonReport> {
let mut reports = Vec::new();
while let Ok(report) = self.poison_rx.try_recv() {
if let Some(m) = &self.metrics {
m.objects_failed(report.kind).increment(1);
}
reports.push(report);
}
reports
}
pub(crate) fn set_paused(&self, lanes: &[LaneId], paused: bool) {
for state in self.splits.values() {
if lanes.contains(&state.lane) {
state.pause.store(paused, Ordering::Relaxed);
}
}
}
}
impl SplitSource for SplitCtx {
type Lane = S3Lane;
fn open_split(&mut self, opening: SplitOpening<'_>) -> Result<S3Lane, SourceError> {
let split = opening.split.id.clone();
let descriptor = SplitDescriptor::decode(&opening.split.descriptor).map_err(|e| {
SourceError::Client {
class: ErrorClass::Fatal,
reason: format!("split {split}: {}", e.reason),
}
})?;
let objects: Arc<Vec<ObjectEntry>> = Arc::new(descriptor.to_entries());
if let Some(p) = opening.resume
&& p.watermark < 0
{
return Err(SourceError::Client {
class: ErrorClass::Fatal,
reason: format!(
"split {split}: carried watermark {} is negative — no position in \
any descriptor",
p.watermark
),
});
}
let resume = opening.resume.map(|p| Position::decode(p.watermark));
let resume_watermark = opening.resume.map_or(0, |p| p.watermark);
let start_ordinal = resume.map_or(0, |p| p.ordinal);
let resume_etag = match opening.resume {
Some(p) => ProgressState::decode(&p.state)
.map_err(|reason| SourceError::Client {
class: ErrorClass::Fatal,
reason: format!("split {split}: {reason}"),
})?
.etag
.or_else(|| {
objects
.get(start_ordinal as usize)
.and_then(|e| e.etag.clone())
}),
None => None,
};
let tracker = Arc::new(SplitTracker::new());
let pause = Arc::new(AtomicBool::new(false));
let stop = Arc::new(AtomicBool::new(false));
let (tx, rx) = mpsc::channel(LANE_HANDOFF_CHUNKS);
drop(self.handle.spawn(run_fetcher(FetcherParams {
split: split.clone(),
store: Arc::clone(&self.store),
slice: Arc::clone(&objects),
start_ordinal,
resume_etag,
chunk_bytes: self.chunk_bytes,
range_bytes: self.range_bytes,
tx,
pause: Arc::clone(&pause),
stop: Arc::clone(&stop),
retry_base: self.retry_base,
retries: self.metrics.as_ref().map(|m| m.get_retries.clone()),
})));
let remaining_at_open = (objects.len() as u64).saturating_sub(u64::from(start_ordinal));
if let Some(m) = &self.metrics {
m.objects_remaining.increment(remaining_at_open as f64);
}
let lane = S3Lane::new(
opening.lane,
opening.partition,
rx,
self.handle.clone(),
self.issuer.clone(),
self.compression,
Arc::clone(&self.make_framer),
resume,
split.clone(),
Arc::clone(&tracker),
self.poison_tx.clone(),
opening.waker.clone(),
self.metrics.clone(),
);
self.splits.insert(
split,
SplitState {
objects,
lane: opening.lane,
pause,
stop,
tracker,
remaining_at_open,
resume_watermark,
last_committed: None,
hinted: false,
revocation: false,
},
);
Ok(lane)
}
fn validate_resume(
&self,
split: &SplitSpec,
progress: &SplitProgress,
) -> Result<(), SourceError> {
let refuse = |detail: String| SourceError::Client {
class: ErrorClass::Fatal,
reason: format!(
"split {}: carried progress no longer matches its descriptor ({detail}); \
this is unrecoverable divergence — requeue the split or start a fresh job",
split.id
),
};
let descriptor =
SplitDescriptor::decode(&split.descriptor).map_err(|e| SourceError::Client {
class: ErrorClass::Fatal,
reason: format!("split {}: {}", split.id, e.reason),
})?;
let state = ProgressState::decode(&progress.state).map_err(refuse)?;
if progress.watermark < 0 {
return Err(refuse(format!(
"watermark {} is negative — no position in any descriptor",
progress.watermark
)));
}
let pos = Position::decode(progress.watermark);
let ordinal = pos.ordinal as usize;
if descriptor.objects.is_empty() && progress.watermark == 0 {
return Ok(()); }
let Some(entry) = descriptor.objects.get(ordinal) else {
return Err(refuse(format!(
"watermark ordinal {ordinal} is outside the {}-object descriptor",
descriptor.objects.len()
)));
};
if let Some(key) = &state.key
&& key != &entry.key
{
return Err(refuse(format!(
"progress pins object \"{key}\" at ordinal {ordinal}, the descriptor has \
\"{}\" there",
entry.key
)));
}
if let (Some(pinned), Some(listed)) = (&state.etag, &entry.etag)
&& pinned != listed
{
return Err(refuse(format!(
"progress pins ETag {pinned} at ordinal {ordinal}, the descriptor has \
{listed} — progress and descriptor were minted against different content"
)));
}
if pos.record > 0 && state.etag.is_none() && entry.etag.is_none() {
return Err(refuse(format!(
"the watermark is mid-object (\"{}\", {} records committed) but neither \
the progress nor the descriptor carries an ETag to pin the re-read to",
entry.key, pos.record
)));
}
Ok(())
}
fn encode_commit(
&mut self,
split: &SplitId,
watermark: i64,
) -> Result<SplitProgress, SourceError> {
let Some(state) = self.splits.get_mut(split) else {
return Err(SourceError::Client {
class: ErrorClass::Fatal,
reason: format!("commit for unheld split {split} — driver/context wiring bug"),
});
};
let pos = Position::decode(watermark);
let members = state.objects.len();
let legal_empty = members == 0 && watermark == 0;
if pos.ordinal as usize >= members && !legal_empty {
return Err(SourceError::Client {
class: ErrorClass::Fatal,
reason: format!(
"split {split}: committed watermark {watermark} decodes to ordinal {} \
record {}, outside the {members}-member descriptor — offset accounting bug",
pos.ordinal, pos.record
),
});
}
state.last_committed = Some(watermark);
let payload = ProgressState::at(&state.objects, watermark).encode();
Ok(
if state.tracker.terminal() == Some(watermark) && !state.revocation {
SplitProgress::completed(watermark, payload)
} else {
SplitProgress::new(watermark, payload)
},
)
}
fn begin_revoke(&mut self, split: &SplitId) -> bool {
let Some(state) = self.splits.get_mut(split) else {
return false;
};
state.revocation = true;
state.stop.store(true, Ordering::Relaxed);
true
}
fn sweep(&mut self, split: &SplitId) -> Result<Option<SplitProgress>, SourceError> {
let Some(state) = self.splits.get(split) else {
return Ok(None);
};
if state.revocation {
return Ok(None);
}
let Some(terminal) = state.tracker.terminal() else {
return Ok(None); };
let acked_all = state.last_committed == Some(terminal)
|| (state.last_committed.is_none() && state.resume_watermark == terminal);
if !acked_all {
return Ok(None); }
let payload = ProgressState::at(&state.objects, terminal).encode();
Ok(Some(SplitProgress::completed(terminal, payload)))
}
fn drain_ready(&mut self, split: &SplitId) -> Result<Option<SplitProgress>, SourceError> {
let Some(state) = self.splits.get(split) else {
return Ok(None);
};
let Some(cut) = state.tracker.terminal() else {
return Ok(None);
};
let acked_all = state.last_committed == Some(cut)
|| (state.last_committed.is_none() && state.resume_watermark == cut);
if !acked_all {
return Ok(None); }
let payload = ProgressState::at(&state.objects, cut).encode();
Ok(Some(SplitProgress::new(cut, payload)))
}
fn close_split(&mut self, split: &SplitId) {
if let Some(state) = self.splits.remove(split)
&& let Some(m) = &self.metrics
{
let done = u64::from(state.tracker.objects_done());
m.objects_remaining
.decrement(state.remaining_at_open.saturating_sub(done) as f64);
}
}
fn take_finishing(&mut self) -> Vec<SplitId> {
self.splits
.iter_mut()
.filter_map(|(id, state)| {
if state.hinted || state.tracker.terminal().is_none() {
return None;
}
state.hinted = true;
Some(id.clone())
})
.collect()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn tracker_is_unknown_until_set_and_stable_after() {
let t = SplitTracker::new();
assert_eq!(t.terminal(), None);
t.set_terminal(42);
assert_eq!(t.terminal(), Some(42));
t.set_terminal(42); assert_eq!(t.terminal(), Some(42));
}
#[test]
fn tracker_carries_zero_watermarks_and_counts_objects() {
let t = SplitTracker::new();
t.set_terminal(0);
assert_eq!(t.terminal(), Some(0));
t.object_done();
t.object_done();
assert_eq!(t.objects_done(), 2);
}
#[test]
fn progress_state_round_trips_and_pins_the_encoding() {
let state = ProgressState {
v: PROGRESS_STATE_VERSION,
key: Some("exports/part-000.ndjson".into()),
etag: Some("\"9b2cf5\"".into()),
};
let bytes = state.encode();
assert_eq!(
String::from_utf8(bytes.clone()).unwrap(),
r#"{"v":1,"key":"exports/part-000.ndjson","etag":"\"9b2cf5\""}"#
);
assert_eq!(ProgressState::decode(&bytes).unwrap(), state);
}
#[test]
fn empty_progress_state_decodes_as_unpinned() {
let state = ProgressState::decode(b"").unwrap();
assert_eq!(state.key, None);
assert_eq!(state.etag, None);
}
#[test]
fn unknown_progress_state_version_is_rejected() {
let err = ProgressState::decode(br#"{"v":99}"#).unwrap_err();
assert!(err.contains("version 99"), "{err}");
}
fn test_ctx() -> (SplitCtx, tokio::runtime::Runtime) {
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()
.unwrap();
let ctx = SplitCtx::new(
Arc::new(object_store::memory::InMemory::new()),
rt.handle().clone(),
spate_core::checkpoint::Checkpointer::new().handle(),
Arc::new(|| Box::new(crate::testutil::TestLineFramer::new(1 << 20))),
Compression::Auto,
64,
1024,
Duration::from_millis(1),
None,
);
(ctx, rt)
}
fn spec_of(objects: &[(&str, Option<&str>)]) -> SplitSpec {
let entries: Vec<ObjectEntry> = objects
.iter()
.map(|(key, etag)| ObjectEntry {
key: (*key).to_string(),
size: 1,
etag: etag.map(str::to_owned),
last_modified_ms: 0,
})
.collect();
let id = crate::split::split_id_for(objects.iter().copied()).unwrap();
SplitSpec::new(
id,
SplitDescriptor::from_entries(&entries).encode().unwrap(),
)
}
fn progress(ordinal: u32, record: u64, key: Option<&str>, etag: Option<&str>) -> SplitProgress {
let watermark = Position { ordinal, record }.encode().unwrap();
let state = ProgressState {
v: PROGRESS_STATE_VERSION,
key: key.map(str::to_owned),
etag: etag.map(str::to_owned),
};
SplitProgress::new(watermark, state.encode())
}
#[test]
fn validate_resume_rejects_each_drift_kind() {
let (ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a")), ("b", Some("e-b")), ("c", Some("e-c"))]);
ctx.validate_resume(&spec, &progress(1, 0, Some("b"), Some("e-b")))
.unwrap();
ctx.validate_resume(&spec, &SplitProgress::new(0, Vec::new()))
.unwrap();
let err = ctx
.validate_resume(&spec, &progress(9, 0, Some("z"), None))
.unwrap_err();
assert!(err.to_string().contains("outside"), "{err}");
let err = ctx
.validate_resume(&spec, &progress(1, 0, Some("was-here"), Some("e-b")))
.unwrap_err();
assert!(err.to_string().contains("was-here"), "{err}");
let err = ctx
.validate_resume(&spec, &progress(1, 0, Some("b"), Some("old-etag")))
.unwrap_err();
assert!(err.to_string().contains("different content"), "{err}");
let err = ctx
.validate_resume(&spec, &SplitProgress::new(0, b"garbage".to_vec()))
.unwrap_err();
assert!(err.to_string().contains("decode"), "{err}");
let unpinned_spec = spec_of(&[("a", None), ("b", None)]);
let err = ctx
.validate_resume(&unpinned_spec, &progress(1, 3, Some("b"), None))
.unwrap_err();
assert!(err.to_string().contains("mid-object"), "{err}");
ctx.validate_resume(&spec, &progress(1, 3, Some("b"), None))
.unwrap();
}
#[test]
fn negative_watermark_from_the_store_is_refused_without_panicking() {
let (ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a"))]);
let err = ctx
.validate_resume(&spec, &SplitProgress::new(-1, Vec::new()))
.unwrap_err();
assert!(err.to_string().contains("negative"), "{err}");
}
#[test]
fn commit_beyond_the_descriptor_is_an_accounting_bug() {
use spate_core::source::LaneId;
let (mut ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a")), ("b", Some("e-b"))]);
let id = spec.id.clone();
let descriptor = SplitDescriptor::decode(&spec.descriptor).unwrap();
ctx.splits.insert(
id.clone(),
SplitState {
objects: Arc::new(descriptor.to_entries()),
lane: LaneId(0),
pause: Arc::new(AtomicBool::new(false)),
stop: Arc::new(AtomicBool::new(false)),
tracker: Arc::new(SplitTracker::new()),
remaining_at_open: 2,
resume_watermark: 0,
last_committed: None,
hinted: false,
revocation: false,
},
);
ctx.encode_commit(
&id,
Position {
ordinal: 1,
record: 3,
}
.encode()
.unwrap(),
)
.unwrap();
ctx.encode_commit(
&id,
Position {
ordinal: 1,
record: crate::offset::MAX_RECORD_INDEX + 1,
}
.encode()
.unwrap(),
)
.unwrap();
let err = ctx
.encode_commit(
&id,
Position {
ordinal: 2,
record: 0,
}
.encode()
.unwrap(),
)
.unwrap_err();
assert!(err.to_string().contains("accounting"), "{err}");
let err = ctx
.encode_commit(
&id,
Position {
ordinal: 5,
record: 0,
}
.encode()
.unwrap(),
)
.unwrap_err();
assert!(err.to_string().contains("accounting"), "{err}");
let err = ctx
.encode_commit(
&id,
Position {
ordinal: 2,
record: 1,
}
.encode()
.unwrap(),
)
.unwrap_err();
assert!(err.to_string().contains("accounting"), "{err}");
let empty_id = crate::split::split_id_for([("empty-placeholder", None)]).unwrap();
ctx.splits.insert(
empty_id.clone(),
SplitState {
objects: Arc::new(Vec::new()),
lane: LaneId(1),
pause: Arc::new(AtomicBool::new(false)),
stop: Arc::new(AtomicBool::new(false)),
tracker: Arc::new(SplitTracker::new()),
remaining_at_open: 0,
resume_watermark: 0,
last_committed: None,
hinted: false,
revocation: false,
},
);
ctx.encode_commit(&empty_id, 0).unwrap();
}
#[test]
fn empty_split_progress_is_legal_at_zero_only() {
let (ctx, _rt) = test_ctx();
let entries: Vec<ObjectEntry> = Vec::new();
let id = crate::split::split_id_for([("placeholder", None)]).unwrap();
let spec = SplitSpec::new(
id,
SplitDescriptor::from_entries(&entries).encode().unwrap(),
);
ctx.validate_resume(&spec, &SplitProgress::new(0, Vec::new()))
.unwrap();
let err = ctx
.validate_resume(&spec, &progress(0, 1, None, None))
.unwrap_err();
assert!(err.to_string().contains("outside"), "{err}");
}
fn seed_split(
ctx: &mut SplitCtx,
id: &SplitId,
objects: Vec<ObjectEntry>,
resume_watermark: i64,
) -> Arc<SplitTracker> {
let tracker = Arc::new(SplitTracker::new());
let remaining = objects.len() as u64;
ctx.splits.insert(
id.clone(),
SplitState {
objects: Arc::new(objects),
lane: LaneId(0),
pause: Arc::new(AtomicBool::new(false)),
stop: Arc::new(AtomicBool::new(false)),
tracker: Arc::clone(&tracker),
remaining_at_open: remaining,
resume_watermark,
last_committed: None,
hinted: false,
revocation: false,
},
);
tracker
}
fn objects_of(spec: &SplitSpec) -> Vec<ObjectEntry> {
SplitDescriptor::decode(&spec.descriptor)
.unwrap()
.to_entries()
}
#[test]
fn begin_revoke_declines_an_unknown_split() {
let (mut ctx, _rt) = test_ctx();
let id = crate::split::split_id_for([("never-opened", None)]).unwrap();
assert!(
!ctx.begin_revoke(&id),
"an unheld split must be declined so the driver never opens a phantom revocation"
);
}
#[test]
fn a_drain_commit_at_the_cut_is_never_completed() {
let (mut ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a"))]);
let id = spec.id.clone();
let tracker = seed_split(&mut ctx, &id, objects_of(&spec), 0);
assert!(ctx.begin_revoke(&id), "a held split accepts the revocation");
assert!(
ctx.splits[&id].stop.load(Ordering::Relaxed),
"begin_revoke releases the fetcher via the shared stop flag"
);
let cut = Position {
ordinal: 0,
record: 3,
}
.encode()
.unwrap();
tracker.set_terminal(cut);
let progress = ctx.encode_commit(&id, cut).unwrap();
assert_eq!(progress.watermark, cut);
assert!(
!progress.completed,
"the cut watermark of a handing-off split must never complete it"
);
}
#[test]
fn drain_ready_waits_for_the_acked_tail_then_hands_over_uncompleted() {
let (mut ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a"))]);
let id = spec.id.clone();
let tracker = seed_split(&mut ctx, &id, objects_of(&spec), 0);
assert!(ctx.begin_revoke(&id));
assert!(ctx.drain_ready(&id).unwrap().is_none());
let cut = Position {
ordinal: 0,
record: 2,
}
.encode()
.unwrap();
tracker.set_terminal(cut);
assert!(
ctx.drain_ready(&id).unwrap().is_none(),
"a revocation waits until every emitted record is acked and committed"
);
let _ = ctx.encode_commit(&id, cut).unwrap();
let progress = ctx
.drain_ready(&id)
.unwrap()
.expect("the acked tail is handed over");
assert_eq!(progress.watermark, cut);
assert!(
!progress.completed,
"a revocation hands the split off with completed: false"
);
}
#[test]
fn sweep_declines_a_draining_split_that_would_otherwise_complete() {
let (mut ctx, _rt) = test_ctx();
let spec = spec_of(&[("a", Some("e-a"))]);
let id = spec.id.clone();
let tracker = seed_split(&mut ctx, &id, objects_of(&spec), 0);
tracker.set_terminal(0);
assert!(
ctx.begin_revoke(&id),
"the split is held, so the revocation is accepted"
);
assert!(
ctx.sweep(&id).unwrap().is_none(),
"a handing-off split is given away via drain_ready, never completed by the sweep"
);
let spec2 = spec_of(&[("b", Some("e-b"))]);
let id2 = spec2.id.clone();
let tracker2 = seed_split(&mut ctx, &id2, objects_of(&spec2), 0);
tracker2.set_terminal(0);
assert!(
ctx.sweep(&id2).unwrap().is_some(),
"the same shape without a revocation completes via the sweep"
);
}
}