Skip to main content

MultiStream

Struct MultiStream 

Source
pub struct MultiStream<B: EventStreamBackend> {
    pub logger: Arc<ServerLogger>,
    /* private fields */
}
Expand description

Manages multiple streams with synchronization logic based on shared bindings.

This struct handles the creation and alignment of streams to ensure proper synchronization when bindings (e.g., buffers) are shared across different streams.

Fields§

§logger: Arc<ServerLogger>

The logger used by the server.

Implementations§

Source§

impl<B: EventStreamBackend> MultiStream<B>

Source

pub fn backend_mut(&mut self) -> &mut B

Mutable access to the stream-creation backend, e.g. to change the configuration new streams are created with. Already-created streams are unaffected.

Source

pub fn new(logger: Arc<ServerLogger>, backend: B, max_streams: u8) -> Self

Creates an empty multi-stream.

Source

pub fn stream_ids(&self) -> impl Iterator<Item = StreamId> + '_

Synthetic StreamIds, one per initialized stream (see StreamPool::stream_ids).

Source

pub fn gc(&mut self, gc: GcTask<B>)

Enqueue a task to be cleaned.

Source

pub fn try_stream_mut(&mut self, stream_id: &StreamId) -> Option<&mut B::Stream>

The backend stream on stream_id’s slot when that slot was ever initialized, mutably — a lookup that must not create a stream (see StreamPool::try_get_mut).

Source

pub fn resolve<'a>( &mut self, stream_id: StreamId, handles: impl Iterator<Item = &'a BufferBinding>, ) -> ResolvedStreams<'_, B>

Resolves and returns a mutable reference to the stream for the given ID, performing any necessary alignment based on the provided bindings.

This method ensures that the stream is synchronized with any shared bindings from other streams before returning the stream reference.

Trait Implementations§

Source§

impl<B: Debug + EventStreamBackend> Debug for MultiStream<B>

Source§

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

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

impl<B: EventStreamBackend> FailureStore for MultiStream<B>

Source§

type Factory = EventStreamBackendWrapper<B>

The factory the driver’s StreamPool was built from. Its streams expose the memory the taint is recorded on.
Source§

fn split(&mut self) -> (&mut StreamPool<Self::Factory>, &mut Failures)

The pool and the store, split-borrowed: nearly every operation here reaches the allocations through the pool while mutating the store.
Source§

fn parts(&self) -> (&StreamPool<Self::Factory>, &Failures)

split for the read-only questions.
Source§

fn ensure_written<'a>( &self, handles: impl Iterator<Item = &'a BufferBinding>, ) -> Result<(), ServerError>

Fails when the buffers handles name carry a failure, with the errors of the work that was supposed to write them. Read more
Source§

fn read_failure<'a>( &self, reads: impl Iterator<Item = &'a BufferBinding>, ) -> Option<ReadFailure>

The failure claiming bytes any of reads names, with its error — the check a launch makes before it runs. Read more
Source§

fn taint<'a>( &mut self, error: ServerError, written: impl Iterator<Item = &'a BufferBinding>, )

Taint every allocation in written with error: the work that was going to write those buffers did not run, so a read of any of them fails on this failure until something writes them again. Read more
Source§

fn written<'a>(&mut self, written: impl Iterator<Item = &'a BufferBinding>)

Release the failure on every allocation in written: work that writes them has been enqueued, so a read of one is no longer reading bytes nothing wrote.
Source§

fn propagate( &mut self, found: &ReadFailure, kernel: KernelId, written: Vec<BufferBinding>, )

A skipped launch’s outputs take the failure that stopped it: nothing wrote them, exactly as if the launch had failed, and the claim names the root cause rather than minting a new one. The skip is recorded on the failure, so a read of anything downstream can name the path back to the root. Read more
Source§

fn write_set(&mut self) -> Vec<BufferBinding>

An empty write set, pooled here so a launch allocates nothing for it. exit_write hands it back.
Source§

fn enter_write(&mut self, written: &[BufferBinding]) -> Option<FailureId>

Enter a write scope over written: taint every buffer the work is going to write with a provisional failure, minted here because the real one does not exist yet. Read more
Source§

fn exit_write( &mut self, provisional: Option<FailureId>, written: Vec<BufferBinding>, error: Option<&ServerError>, )

Settle the scope entered over written: release the provisional failure when the work was enqueued, and swap the real error in for it when the work was not. The taint is the whole answer — a read of one of these buffers fails on it, whoever asks — and the error is logged here, the backstop for the failure nobody ever reads. The staged vector goes back to the pool either way.

Auto Trait Implementations§

§

impl<B> !RefUnwindSafe for MultiStream<B>

§

impl<B> !UnwindSafe for MultiStream<B>

§

impl<B> Freeze for MultiStream<B>

§

impl<B> Send for MultiStream<B>
where StreamPool<EventStreamBackendWrapper<B>>: Send, GcThread<B>: Send,

§

impl<B> Sync for MultiStream<B>
where StreamPool<EventStreamBackendWrapper<B>>: Sync, GcThread<B>: Sync,

§

impl<B> Unpin for MultiStream<B>

§

impl<B> UnsafeUnpin for MultiStream<B>

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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<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> 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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
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