Skip to main content

MultiStream

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:

  1. Materialise src. The id src might be the output of an op still pending on src_stream. We need its backing handle to actually exist before we can alias it. If HandleContainer::get_handle_ref returns None — meaning no op has produced a handle for src yet — we drain src_stream synchronously, forcing every pending op to run (and thus the handle to be registered). We also record src in shared_sources so that any next share of the same src can skip the drain: once registered, a handle stays put (see invariants below).

  2. Alias the handle under dst. HandleContainer::register_handle is called with dst and a clone() of src’s backend handle. Cubecl handles are Arc-style reference counters over a backing buffer, so clone() is cheap and the buffer survives until the last alias drops. After this call, handles[src] and handles[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_stream and removes handles[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,

pub fn tag_shared_view( &mut self, src_stream: StreamId, src: TensorId, dst: TensorId, handles: &mut HandleContainer<<R as FusionRuntime>::FusionHandle>, )

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, drain src_stream so the producing op runs and registers it. Skip the drain on subsequent shares of the same source by remembering it in shared_sources or by observing that the handle is already in the container (which covers the chained-share case where src is itself a previously-aliased view).
  • Then alias the backing handle under dst. register_handle clones 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>, )

Mark a tensor as read.

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§

§

impl<T> Adaptor<()> for T

§

fn adapt(&self)

Adapt the type to be passed to a metric.
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
§

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

§

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

§

impl<K, Q> Comparable<Q> for K
where K: Borrow<Q> + ?Sized, Q: Ord + ?Sized,

§

fn compare(&self, key: &Q) -> Ordering

Compare self to key and return their ordering.
§

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

§

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

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

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.

§

impl<T> Instrument for T

§

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

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

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
§

impl<T> Pointable for T

§

const ALIGN: usize

The alignment of pointer.
§

type Init = T

The type for initializers.
§

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

Initializes a with the given initializer. Read more
§

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

Dereferences the given pointer. Read more
§

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

Mutably dereferences the given pointer. Read more
§

unsafe fn drop(ptr: usize)

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

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

§

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
§

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
§

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 = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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.
§

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

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

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