Struct MultiStream
pub struct MultiStream<R>where
R: FusionRuntime,{ /* private fields */ }Expand description
Keep track of multiple concurrent lazy streams of operations.
§Why this exists
Each Stream holds a lazy queue of OperationIrs whose inputs are assumed
to live on that stream. That makes single-stream execution simple — every
TensorId in a queue is resolvable from the same handle map and the same
pending op chain. But a FusionTensor is Send + Clone,
so user code can move or clone a tensor from one thread (= one StreamId) to
another. The receiving thread will then submit ops whose inputs reference a
tensor whose home is a different stream’s queue. This struct is what makes
that case behave correctly without giving up the stream-locality invariant.
§Strategy: shared views
We never let a foreign-stream tensor id appear in another stream’s queue.
Instead, when FusionTensor::clone or
FusionTensor::into_ir detects that
self.stream != StreamId::current(), it allocates a fresh id (dst) and calls
tag_shared_view with (src_stream, src, dst). That call does two
things, in order:
-
Materialise
src. The idsrcmight be the output of an op still pending onsrc_stream. We need its backing handle to actually exist before we can alias it. IfHandleContainer::get_handle_refreturnsNone— meaning no op has produced a handle forsrcyet — we drainsrc_streamsynchronously, forcing every pending op to run (and thus the handle to be registered). We also recordsrcinshared_sourcesso that any next share of the samesrccan skip the drain: once registered, a handle stays put (see invariants below). -
Alias the handle under
dst.HandleContainer::register_handleis called withdstand aclone()ofsrc’s backend handle. Cubecl handles areArc-style reference counters over a backing buffer, soclone()is cheap and the buffer survives until the last alias drops. After this call,handles[src]andhandles[dst]are two distinct map entries that both point at the same allocation.
shared_view then returns a new FusionTensor carrying (id = dst, stream = current). Every subsequent op on that tensor enqueues on current like any
other local tensor — the rest of the fusion engine sees no special case.
§Freeing
Each FusionTensor::drop frees its id on its own stream field (the
home stream of that particular alias), not the calling thread’s stream: on
the home thread it enqueues an OperationIr::Drop(ir), from any other it goes
through foreign_drop, which never touches the queue. So:
- The original tensor’s drop targets
src_streamand removeshandles[src]. - The alias tensor’s drop targets the stream that minted it and removes
handles[dst].
Each removal decrements the backend handle’s Arc refcount; the underlying
buffer is freed only after the last side drops. No cross-stream coordination
is needed.
§Bounding shared_sources
Naively the set would grow forever, since tag_shared_view only ever inserts.
Cleanup happens when the last live FusionTensor with an id drops, without
waiting for the free to actually run. No future tag_shared_view can then
receive that id as a src, so removing the entry cannot trigger a redundant drain.
§The SSA-like invariant
Skipping the drain on subsequent shares of the same src relies on a
property of the fusion IR: every op output uses a fresh TensorId
allocated by crate::Client::create_empty_handle, never the id of an
input. Once handles[src] is set, no later op overwrites it; the data
behind src is effectively immutable from the IR’s point of view.
The cubecl-fusion engine does sometimes reuse the backing buffer of an
input for an output (in-place fusion), but that path is gated by
handle.can_mut(), which returns false the moment another reference exists.
Calling handle.clone() in step 2 above is precisely that extra reference,
so aliased sources are never eligible for in-place reuse — the engine
allocates a fresh output buffer instead.
§The chained-share fast path
When a share is itself re-shared (owner → peer → grandpeer), the second
tag_shared_view call has src = peer's id. That id was set up directly by
the previous call (via register_handle), not by an op enqueued on the peer
stream, so handles.get_handle_ref(&src) is already Some the moment we
look. The drain check therefore short-circuits — even though peer's id was
never added to shared_sources (only sources that required a
drain are tracked there), the handle-existence test alone is sufficient.
Implementations§
§impl<R> MultiStream<R>where
R: FusionRuntime,
impl<R> MultiStream<R>where
R: FusionRuntime,
Set up a cross-stream alias dst for the foreign tensor src that lives on
src_stream. Called when FusionTensor::clone
or FusionTensor::into_ir detects that the
tensor’s home stream is not the current stream.
See the MultiStream struct-level docs for the full strategy. In short:
- If
src’s handle isn’t materialised yet, drainsrc_streamso the producing op runs and registers it. Skip the drain on subsequent shares of the same source by remembering it inshared_sourcesor by observing that the handle is already in the container (which covers the chained-share case wheresrcis itself a previously-aliased view). - Then alias the backing handle under
dst.register_handleclones the cubecl handle (Arc-style), so both ids share refcount on the buffer until each side is freed.
pub fn mark_read(
&mut self,
id: StreamId,
ir: &TensorIr,
handles: &HandleContainer<<R as FusionRuntime>::FusionHandle>,
)
pub fn mark_read( &mut self, id: StreamId, ir: &TensorIr, handles: &HandleContainer<<R as FusionRuntime>::FusionHandle>, )
Mark a tensor as read.
pub fn drain(
&mut self,
handles: &mut HandleContainer<<R as FusionRuntime>::FusionHandle>,
id: StreamId,
)
pub fn drain( &mut self, handles: &mut HandleContainer<<R as FusionRuntime>::FusionHandle>, id: StreamId, )
Run id’s pending segment to completion.
Reports nothing. An operation that fails leaves its error on the
tensors it was going to write (see execution::set_output_errors), so
it is delivered by the read of one of those tensors — the point where a
caller is actually waiting for that data — and not to whoever happened
to drain the stream next. A drain that shares no tensor with the
failure has nothing to report and returns normally.
Auto Trait Implementations§
impl<R> !RefUnwindSafe for MultiStream<R>
impl<R> !Sync for MultiStream<R>
impl<R> !UnwindSafe for MultiStream<R>
impl<R> Freeze for MultiStream<R>
impl<R> Send for MultiStream<R>
impl<R> Unpin for MultiStream<R>
impl<R> UnsafeUnpin for MultiStream<R>
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
§impl<K, Q> Comparable<Q> for K
impl<K, Q> Comparable<Q> for K
§impl<K, Q> Equivalent<Q> for K
impl<K, Q> Equivalent<Q> for K
§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> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§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