pub enum StreamOutcome {
Empty {
watermark: Watermark,
rpo: RpoStatus,
},
Streamed {
generation_key: String,
first_frame: u64,
last_frame: u64,
checkpoint_seq: u32,
frame_count: u64,
backpressure: BackpressureReport,
rpo: RpoStatus,
},
Restarted {
generation_key: String,
previous_generation: WalGeneration,
new_generation: WalGeneration,
first_frame: u64,
last_frame: u64,
frame_count: u64,
backpressure: BackpressureReport,
rpo: RpoStatus,
},
Shed {
checkpoint_seq: u32,
first_frame: u64,
last_frame: u64,
backpressure: BackpressureReport,
rpo: RpoStatus,
},
Fenced {
current_epoch: u64,
our_epoch: u64,
current_pointer_generation: u64,
our_pointer_generation: u64,
},
}Expand description
What a tail_frames call did.
Variants§
Empty
No new frames since the last tail — sink is already current.
Fields
Streamed
Uploaded a contiguous range of frames and wrote a generation manifest.
Fields
backpressure: BackpressureReportR574-F2: backpressure activity during this call (policy, high-water spill-buffer occupancy, shed count, throttle retries).
rpo: RpoStatusR574-T4: see StreamOutcome::Empty::rpo.
Restarted
The live WAL is not provably the one the sidecar watermark was taken
from, so this call re-uploaded frames 1..N from the top rather than
resuming.
R858-B19 widened this from “checkpoint_seq advanced” to “the
WalGeneration did not prove itself unchanged”, which is why both
fields are now generations rather than bare sequence numbers: the
motivating case is a writer restart where the sequence reads 0 -> 0
and only the salt moved. previous_generation.salt == None names the
third case — a sidecar written before the salt existed, restarted
because it cannot be checked, not because it was seen to change.
Fields
previous_generation: WalGenerationnew_generation: WalGenerationbackpressure: BackpressureReportR574-F2: see StreamOutcome::Streamed::backpressure.
rpo: RpoStatusR574-T4: see StreamOutcome::Empty::rpo.
Shed
R574-F2: BackpressurePolicy::Shed dropped every buffered frame in
this call before any of them persisted (R2 was throttling harder
than the spill buffer + backoff could absorb). No manifest/watermark
was written — the next tail_frames call re-attempts the same
range from the unchanged prior watermark. Distinct from Empty,
which means the engine itself had nothing new.
Fenced
R732-F2 (W245) / R736-T2 (W250): this writer is a stale owner and
wrote nothing. Either the sink’s watermark is stamped with an epoch
higher than StreamConfig::epoch (ownership moved within the cell),
or with a pointer generation higher than
StreamConfig::pointer_generation (ownership moved to a different
cell) — the two-level fence bounces on either. Detected before the
first frame upload, so a fenced call is a pure read — no frames, no
manifest, no watermark write.
This is the outcome the whole fencing design exists to produce. Without it a partitioned old master and a freshly-promoted new master both stream into the same prefix and silently corrupt each other; with it the loser finds out on its very next tail and can stop.
Deliberately carries no RpoStatus: a fenced writer’s view of
watermark staleness is not its stream’s RPO any more, and reporting one
here would page the wrong operator about the wrong node.
Fields
Implementations§
Source§impl StreamOutcome
impl StreamOutcome
Sourcepub fn rpo(&self) -> Option<&RpoStatus>
pub fn rpo(&self) -> Option<&RpoStatus>
This call’s RPO snapshot, or None for StreamOutcome::Fenced —
see that variant’s doc for why it deliberately carries none.
R782: the accessor a caller (tenant-streamer’s tail loop) uses to
push watermark_age onward without re-deriving this match on every
call site that needs it.
Sourcepub fn backpressure(&self) -> Option<&BackpressureReport>
pub fn backpressure(&self) -> Option<&BackpressureReport>
This call’s backpressure activity, or None for the two outcomes that
never reached the drain loop (StreamOutcome::Empty had nothing to
send, StreamOutcome::Fenced was refused before the first frame).
R760-B8: the accessor a multi-tenant caller needs. Streamed is not
the same thing as “the sink took everything” — under
BackpressurePolicy::Shed a call that persisted a partial prefix and
dropped the rest reports Streamed with a nonzero
BackpressureReport::frames_shed, and a caller that only matches the
variant cannot tell that apart from a clean tail. roadcase’s shard
flusher reads this on every arm so a struggling cell is nameable from
its own metrics rather than by bisecting tenants.
Trait Implementations§
Source§impl Clone for StreamOutcome
impl Clone for StreamOutcome
Source§fn clone(&self) -> StreamOutcome
fn clone(&self) -> StreamOutcome
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for StreamOutcome
impl Debug for StreamOutcome
impl Eq for StreamOutcome
Source§impl PartialEq for StreamOutcome
impl PartialEq for StreamOutcome
impl StructuralPartialEq for StreamOutcome
Auto Trait Implementations§
impl Freeze for StreamOutcome
impl RefUnwindSafe for StreamOutcome
impl Send for StreamOutcome
impl Sync for StreamOutcome
impl Unpin for StreamOutcome
impl UnsafeUnpin for StreamOutcome
impl UnwindSafe for StreamOutcome
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Downcast for Twhere
T: Any,
impl<T> Downcast for Twhere
T: Any,
Source§fn into_any(self: Box<T>) -> Box<dyn Any>
fn into_any(self: Box<T>) -> Box<dyn Any>
Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>, which can then be
downcast into Box<dyn ConcreteType> where ConcreteType implements Trait.Source§fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
Rc<Trait> (where Trait: Downcast) to Rc<Any>, which can then be further
downcast into Rc<ConcreteType> where ConcreteType implements Trait.Source§fn as_any(&self) -> &(dyn Any + 'static)
fn as_any(&self) -> &(dyn Any + 'static)
&Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &Any’s vtable from &Trait’s.Source§fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
&mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &mut Any’s vtable from &mut Trait’s.Source§impl<T> DowncastSend for T
impl<T> DowncastSend for T
Source§impl<T> DowncastSync for T
impl<T> DowncastSync for T
Source§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
Source§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
Source§fn equivalent(&self, key: &K) -> bool
fn equivalent(&self, key: &K) -> bool
key and return true if they are equal.Source§impl<K, Q> Equivalent<Q> for K
impl<K, Q> Equivalent<Q> for K
Source§fn equivalent(&self, key: &Q) -> bool
fn equivalent(&self, key: &Q) -> bool
key and return true if they are equal.impl<T> ErasedDestructor for Twhere
T: 'static,
impl<T> Fruit for T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more