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>
impl<B: EventStreamBackend> MultiStream<B>
Sourcepub fn backend_mut(&mut self) -> &mut B
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.
Sourcepub fn new(logger: Arc<ServerLogger>, backend: B, max_streams: u8) -> Self
pub fn new(logger: Arc<ServerLogger>, backend: B, max_streams: u8) -> Self
Creates an empty multi-stream.
Sourcepub fn stream_ids(&self) -> impl Iterator<Item = StreamId> + '_
pub fn stream_ids(&self) -> impl Iterator<Item = StreamId> + '_
Synthetic StreamIds, one per initialized stream (see StreamPool::stream_ids).
Sourcepub fn try_stream_mut(&mut self, stream_id: &StreamId) -> Option<&mut B::Stream>
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).
Sourcepub fn resolve<'a>(
&mut self,
stream_id: StreamId,
handles: impl Iterator<Item = &'a BufferBinding>,
) -> ResolvedStreams<'_, B>
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>
impl<B: Debug + EventStreamBackend> Debug for MultiStream<B>
Source§impl<B: EventStreamBackend> FailureStore for MultiStream<B>
impl<B: EventStreamBackend> FailureStore for MultiStream<B>
Source§type Factory = EventStreamBackendWrapper<B>
type Factory = EventStreamBackendWrapper<B>
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)
fn split(&mut self) -> (&mut StreamPool<Self::Factory>, &mut Failures)
Source§fn ensure_written<'a>(
&self,
handles: impl Iterator<Item = &'a BufferBinding>,
) -> Result<(), ServerError>
fn ensure_written<'a>( &self, handles: impl Iterator<Item = &'a BufferBinding>, ) -> Result<(), ServerError>
handles name carry a failure, with the errors
of the work that was supposed to write them. Read moreSource§fn read_failure<'a>(
&self,
reads: impl Iterator<Item = &'a BufferBinding>,
) -> Option<ReadFailure>
fn read_failure<'a>( &self, reads: impl Iterator<Item = &'a BufferBinding>, ) -> Option<ReadFailure>
reads names, with its error — the
check a launch makes before it runs. Read moreSource§fn taint<'a>(
&mut self,
error: ServerError,
written: impl Iterator<Item = &'a BufferBinding>,
)
fn taint<'a>( &mut self, error: ServerError, written: impl Iterator<Item = &'a BufferBinding>, )
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 moreSource§fn written<'a>(&mut self, written: impl Iterator<Item = &'a BufferBinding>)
fn written<'a>(&mut self, written: impl Iterator<Item = &'a BufferBinding>)
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>,
)
fn propagate( &mut self, found: &ReadFailure, kernel: KernelId, written: Vec<BufferBinding>, )
Source§fn write_set(&mut self) -> Vec<BufferBinding>
fn write_set(&mut self) -> Vec<BufferBinding>
exit_write hands it back.Source§fn enter_write(&mut self, written: &[BufferBinding]) -> Option<FailureId>
fn enter_write(&mut self, written: &[BufferBinding]) -> Option<FailureId>
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 moreSource§fn exit_write(
&mut self,
provisional: Option<FailureId>,
written: Vec<BufferBinding>,
error: Option<&ServerError>,
)
fn exit_write( &mut self, provisional: Option<FailureId>, written: Vec<BufferBinding>, error: Option<&ServerError>, )
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>
impl<B> Sync for MultiStream<B>
impl<B> Unpin for MultiStream<B>
impl<B> UnsafeUnpin for MultiStream<B>
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
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
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
impl<T> ErasedDestructor for Twhere
T: 'static,
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