Skip to main content

StreamOutcome

Enum StreamOutcome 

Source
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

§watermark: Watermark
§rpo: RpoStatus

R574-T4: staleness of the last durably-persisted watermark — still meaningful here, since “nothing new to stream” and “the caller’s scheduler stopped invoking us” look identical from the engine’s side and only this field tells them apart.

§

Streamed

Uploaded a contiguous range of frames and wrote a generation manifest.

Fields

§generation_key: String
§first_frame: u64
§last_frame: u64
§checkpoint_seq: u32
§frame_count: u64
§backpressure: BackpressureReport

R574-F2: backpressure activity during this call (policy, high-water spill-buffer occupancy, shed count, throttle retries).

§

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

§generation_key: String
§previous_generation: WalGeneration
§new_generation: WalGeneration
§first_frame: u64
§last_frame: u64
§frame_count: u64
§

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.

Fields

§checkpoint_seq: u32
§first_frame: u64
§last_frame: u64
§backpressure: BackpressureReport
§

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

§current_epoch: u64

The epoch recorded at the sink — the real owner’s token.

§our_epoch: u64

The (lower) epoch this writer tried to stream under.

§current_pointer_generation: u64

R736-T2: the pointer generation recorded at the sink.

§our_pointer_generation: u64

R736-T2: the (lower) generation this writer tried to stream under.

Implementations§

Source§

impl StreamOutcome

Source

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.

Source

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

Source§

fn clone(&self) -> StreamOutcome

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for StreamOutcome

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Eq for StreamOutcome

Source§

impl PartialEq for StreamOutcome

Source§

fn eq(&self, other: &StreamOutcome) -> bool

Equality operator ==. Read more
1.0.0 (const: unstable) · Source§

fn ne(&self, other: &Rhs) -> bool

Inequality operator !=. Read more
Source§

impl StructuralPartialEq for StreamOutcome

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> Downcast for T
where T: Any,

Source§

fn into_any(self: Box<T>) -> Box<dyn Any>

Converts 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>

Converts 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)

Converts &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)

Converts &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
where T: Any + Send,

Source§

fn into_any_send(self: Box<T>) -> Box<dyn Any + Send>

Converts Box<Trait> (where Trait: DowncastSend) to Box<dyn Any + Send>, which can then be downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> DowncastSync for T
where T: Any + Send + Sync,

Source§

fn into_any_sync(self: Box<T>) -> Box<dyn Any + Sync + Send>

Converts Box<Trait> (where Trait: DowncastSync) to Box<dyn Any + Send + Sync>, which can then be downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

fn into_any_arc(self: Arc<T>) -> Arc<dyn Any + Sync + Send> ⓘ

Converts Arc<Trait> (where Trait: DowncastSync) to Arc<Any>, which can then be downcast into Arc<ConcreteType> where ConcreteType implements Trait.
Source§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

Source§

fn equivalent(&self, key: &K) -> bool

Checks if this value is equivalent to the given key. Read more
Source§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

Source§

fn equivalent(&self, key: &K) -> bool

Compare self to key and return true if they are equal.
Source§

impl<K, Q> Equivalent<Q> for K
where K: Borrow<Q> + ?Sized, Q: Eq + ?Sized,

Source§

fn equivalent(&self, key: &Q) -> bool

Compare self to key and return true if they are equal.
Source§

impl<T> ErasedDestructor for T
where T: 'static,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Fruit for T
where T: Send + Downcast,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts 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
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more