Skip to main content

InMemoryWorkflow

Struct InMemoryWorkflow 

Source
pub struct InMemoryWorkflow<Args> { /* private fields */ }
Expand description

In-memory queue that is based on channels

§Example

let mut backend = InMemoryWorkflow::create();

§Features

FeatureStatusDescription
BackendBasic Backend functionality
TaskSinkAbility to push new tasks
SerializationSerialization support for arguments
PipeExt⚠️Allow other backends to pipe to this backend
BackendFactoryShare the same storage across multiple workers
UpdateAllow updating a task
FetchByIdAllow fetching a task by its ID
RescheduleReschedule a task
ResumeByIdResume a task by its ID
ResumeAbandonedResume abandoned tasks
VacuumVacuum the task storage
Workflow⚠️Flexible enough to support workflows
WaitForCompletion⚠️Wait for tasks to complete without blocking
RegisterWorkerAllow registering a worker with the backend
ListWorkersList all workers registered with the backend
ListTasksList all tasks in the backend

Key: ✅ : Supported | ⚠️ : Not implemented | ❌ : Not Supported | ❗ Limited support

Tests:
§Backend
#[tokio::main]
async fn main() {
    // let mut backend = /* snip */;
    
    async fn task(task: u32, worker: WorkerContext) {
        // Do some work 
    }
    let worker = WorkerBuilder::new("backend-test")
        .backend(backend)
        .build(task);
    let _ = worker.stream().take(1).collect::<Vec<_>>().await;
}
§TaskSink
#[tokio::main]
async fn main() {
    // let mut backend = /* snip */;
    
    backend.push(42).await.unwrap();

    async fn task(task: u32, worker: WorkerContext) {
        worker.stop().unwrap();
    }
    let worker = WorkerBuilder::new("task-sink-test")
        .backend(backend)
        .build(task);
    worker.run().await.unwrap();
}

Implementations§

Source§

impl<Args> InMemoryWorkflow<Args>

Source

pub fn create() -> Shared<Self>

Create a new in-memory storage

Source§

impl<Args> InMemoryWorkflow<Args>

Source

pub fn new_with( sender: MemorySink<Vec<u8>>, receiver: BoxedReceiver<Vec<u8>>, ) -> Shared<Self>

Create a storage given a sender and receiver

Trait Implementations§

Source§

impl<Args> Backend for InMemoryWorkflow<Args>

Source§

type Task = Task<Vec<u8>>

The type of task the backend emits.
Source§

type Error = MemoryStorageError

The error type returned by backend operations
Source§

fn poll_ready( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Result<(), Self::Error>>

Polls whether the backend is ready for the worker to request more work. Read more
Source§

fn poll_next( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Option<Result<Self::Task, Self::Error>>>

Polls the backend for the next available task for this worker. Read more
Source§

fn poll_close( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Result<(), Self::Error>>

Flushes/releases any resources the backend holds (pending acks, open subscriptions, connections) before the worker fully shuts down.
Source§

impl<Args> BackendConfig for InMemoryWorkflow<Args>

Source§

type Id = RandomId

The internal type of this Backend’s TaskId
Source§

type Args = Args

The type of argument this backend emits
Source§

type Kind = Durable

This defines the kind of Backend
Source§

type Config = ()

The config for the backend
Source§

type Layer = StoreResultsLayer

The type representing backend middleware layer.
Source§

fn config(&self) -> &Self::Config

Returns the config associated with the backend.
Source§

fn middleware(&mut self, _: &mut WorkerContext) -> Self::Layer

Returns the backend’s middleware layer.
Source§

impl<Args> Debug for InMemoryWorkflow<Args>

Source§

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

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

impl<Args> Sink<Task<Vec<u8>>> for InMemoryWorkflow<Args>
where Args: Unpin,

Source§

type Error = MemoryStorageError

The type of value produced by the sink when an error occurs.
Source§

fn poll_ready( self: Pin<&mut Self>, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>

Attempts to prepare the Sink to receive a value. Read more
Source§

fn start_send( self: Pin<&mut Self>, item: Task<Vec<u8>>, ) -> Result<(), Self::Error>

Begin the process of sending a value to the sink. Each call to this function must be preceded by a successful call to poll_ready which returned Poll::Ready(Ok(())). Read more
Source§

fn poll_flush( self: Pin<&mut Self>, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>

Flush any remaining output from this sink. Read more
Source§

fn poll_close( self: Pin<&mut Self>, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>

Flush any remaining output and close this sink, if necessary. Read more
Source§

impl<Args, Output> WaitForCompletion<Output> for InMemoryWorkflow<Args>
where Output: DeserializeOwned + Unpin + Send + 'static,

Source§

type ResultStream = WaitForStream<Output>

The result stream type yielding task results
Source§

fn wait_for( &self, task_ids: impl IntoIterator<Item = TaskId>, ) -> Self::ResultStream

Wait for multiple tasks to complete, yielding results as they become available
Source§

fn check_status( &self, task_ids: impl IntoIterator<Item = TaskId> + Send, ) -> impl Future<Output = Result<Vec<TaskResult<Output>>, Self::Error>> + Send

Check current status of tasks without waiting
Source§

fn wait_for_single(&self, task_id: TaskId) -> Self::ResultStream

Wait for a single task to complete, yielding its result
Source§

impl<Args> WireFormatBackend for InMemoryWorkflow<Args>

Source§

type Codec = JsonCodec

The codec used to encode and decode tasks.
Source§

type Compact = Vec<u8>

The compact representation of task arguments.
Source§

fn codec(&self) -> &Self::Codec

Returns a reference to the codec used by the backend.

Auto Trait Implementations§

§

impl<Args> !Freeze for InMemoryWorkflow<Args>

§

impl<Args> !RefUnwindSafe for InMemoryWorkflow<Args>

§

impl<Args> !UnwindSafe for InMemoryWorkflow<Args>

§

impl<Args> Send for InMemoryWorkflow<Args>
where PhantomData<Args>: Send,

§

impl<Args> Sync for InMemoryWorkflow<Args>
where PhantomData<Args>: Sync,

§

impl<Args> Unpin for InMemoryWorkflow<Args>
where PhantomData<Args>: Unpin,

§

impl<Args> UnsafeUnpin for InMemoryWorkflow<Args>
where PhantomData<Args>: UnsafeUnpin,

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<B> BackendExt for B
where B: Backend,

Source§

fn poll_next_args( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Option<Result<Task<Self::Args>, PollNextArgsError<Self>>>>
where Self: Sized + BackendConfig + WireFormatBackend + Backend<Task = Task<Self::Compact>>, Self::Codec: Codec<Self::Args, Compact = Self::Compact>, <Self::Codec as Codec<Self::Args>>::Error: Error + Send + Sync + 'static,

A convenience method for calling poll_next and decoding the Args in one step, returning a Task<Self::Args, ..> instead of Task<Self::Compact, ..>.
Source§

fn pipe_to<Dst>(self, backend: Dst) -> Pipe<Dst, Self>
where Self: Sized,

Pipes every task polled from this backend into sink Read more
Source§

fn inspect_err<F>(self, f: F) -> InspectErr<Self, F>
where Self: Sized, F: Fn(&Self::Error),

Attaches a callback F to be run on each error produced while polling the backend.
Source§

fn map_err<F, E2>(self, f: F) -> MapErr<Self, F>
where Self: Sized, F: Fn(Self::Error) -> E2,

Maps errors produced by the backend from Self::Error into E2, useful for heterogeneous composed backends.
Source§

fn with_codec<NewCodec>(self, codec: NewCodec) -> WithCodec<Self, NewCodec>
where Self: Sized + BackendConfig, NewCodec: Codec<Self::Args>,

Swaps out the backend’s serialization codec entirely (JSON, MessagePack, Protobuf, …) without touching storage logic.
Source§

fn poll_with_stream<S>(self, stream: S) -> PollWith<Self, StreamStrategy<S>>
where Self: Sized, S: Stream + Unpin + Send + 'static,

Wake the worker when a stream receives a new item
Source§

fn poll_with_interval( self, duration: Duration, ) -> PollWith<Self, IntervalStrategy>
where Self: Sized,

Available on crate feature sleep only.
Wake the worker periodically
Source§

fn poll_with_backoff( self, interval: Duration, config: BackoffConfig, ) -> PollWith<Self, BackoffStrategy>
where Self: Sized,

Available on crate feature sleep only.
Wake the worker periodically with a backoff
Source§

fn poll_with_strategy<S>(self, strategy: S) -> PollWith<Self, S>
where Self: Sized, S: PollStrategy,

Wake the worker with a custom strategy
Source§

fn instrumented(self, span: Span) -> Instrumented<Self>
where Self: Sized,

Available on crate feature tracing only.
Provides a span to decorate emitted events
Source§

fn before_start<F, Fut>(self, f: F) -> BeforeStart<Self, Self::Error>
where Self: Sized, F: Fn(&mut Self) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,

Runs an async callback once, before the backend’s first poll_ready is delegated.
Source§

fn before_stop<F, Fut>(self, f: F) -> BeforeStop<Self, Self::Error>
where Self: Sized, F: Fn(&mut Self) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,

Runs an async callback once, before the backend’s poll_close is called.
Source§

fn after_start<F, Fut>(self, f: F) -> AfterStart<Self, Self::Error>
where Self: Sized, F: for<'c> Fn(&mut Self) -> Fut + for<'c> Send + for<'c> Sync + 'static, Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,

Runs an async callback once, after the backend’s first poll_ready is successful.
Source§

fn after_stop<F, Fut>(self, f: F) -> AfterStop<Self, Self::Error>
where Self: Sized, F: Fn(&mut Self) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,

Runs an async callback once, after the worker has stopped and backend has cleaned up.
Source§

fn interleave<S>(self, stream: S) -> Interleave<Self, S>
where Self: Sized, S: Stream<Item = Result<Self::Task, Self::Error>> + Unpin,

Interleaves the external stream with the backend.
Source§

fn wake_on_push(self) -> WakeOnPush<Self>
where Self: Sized,

Wakes the worker when a new item is pushed
Source§

fn shared(self) -> Shared<Self>
where Self: WireFormatBackend + Send, Self::Codec: Clone,

Create a cloneable handle to the inner backend where all handles are clone.
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> 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, Item> SinkExt<Item> for T
where T: Sink<Item> + ?Sized,

Source§

fn with<U, Fut, F, E>(self, f: F) -> With<Self, Item, U, Fut, F>
where F: FnMut(U) -> Fut, Fut: Future<Output = Result<Item, E>>, E: From<Self::Error>, Self: Sized,

Composes a function in front of the sink. Read more
Source§

fn with_flat_map<U, St, F>(self, f: F) -> WithFlatMap<Self, Item, U, St, F>
where F: FnMut(U) -> St, St: Stream<Item = Result<Item, Self::Error>>, Self: Sized,

Composes a function in front of the sink. Read more
Source§

fn sink_map_err<E, F>(self, f: F) -> SinkMapErr<Self, F>
where F: FnOnce(Self::Error) -> E, Self: Sized,

Transforms the error returned by the sink.
Source§

fn sink_err_into<E>(self) -> SinkErrInto<Self, Item, E>
where Self: Sized, Self::Error: Into<E>,

Map this sink’s error to a different error type using the Into trait. Read more
Source§

fn buffer(self, capacity: usize) -> Buffer<Self, Item>
where Self: Sized,

Available on crate feature alloc only.
Adds a fixed-size buffer to the current sink. Read more
Source§

fn close(&mut self) -> Close<'_, Self, Item>
where Self: Unpin,

Close the sink.
Source§

fn fanout<Si>(self, other: Si) -> Fanout<Self, Si>
where Self: Sized, Item: Clone, Si: Sink<Item, Error = Self::Error>,

Fanout items to multiple sinks. Read more
Source§

fn flush(&mut self) -> Flush<'_, Self, Item>
where Self: Unpin,

Flush the sink, processing all pending items. Read more
Source§

fn send(&mut self, item: Item) -> Send<'_, Self, Item>
where Self: Unpin,

A future that completes after the given item has been fully processed into the sink, including flushing. Read more
Source§

fn feed(&mut self, item: Item) -> Feed<'_, Self, Item>
where Self: Unpin,

A future that completes after the given item has been received by the sink. Read more
Source§

fn send_all<'a, St>(&'a mut self, stream: &'a mut St) -> SendAll<'a, Self, St>
where St: TryStream<Ok = Item, Error = Self::Error> + Stream + Unpin + ?Sized, Self: Unpin,

A future that completes after the given stream has been fully processed into the sink, including flushing. Read more
Source§

fn left_sink<Si2>(self) -> Either<Self, Si2>
where Si2: Sink<Item, Error = Self::Error>, Self: Sized,

Wrap this sink in an Either sink, making it the left-hand variant of that Either. Read more
Source§

fn right_sink<Si1>(self) -> Either<Si1, Self>
where Si1: Sink<Item, Error = Self::Error>, Self: Sized,

Wrap this stream in an Either stream, making it the right-hand variant of that Either. Read more
Source§

fn poll_ready_unpin( &mut self, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>
where Self: Unpin,

A convenience method for calling Sink::poll_ready on Unpin sink types.
Source§

fn start_send_unpin(&mut self, item: Item) -> Result<(), Self::Error>
where Self: Unpin,

A convenience method for calling Sink::start_send on Unpin sink types.
Source§

fn poll_flush_unpin( &mut self, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>
where Self: Unpin,

A convenience method for calling Sink::poll_flush on Unpin sink types.
Source§

fn poll_close_unpin( &mut self, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>
where Self: Unpin,

A convenience method for calling Sink::poll_close on Unpin sink types.
Source§

impl<Args, S, E, C> TaskSink<Args, Durable> for S
where S: Sink<Task<<C as Codec<Args>>::Compact>, Error = E> + Unpin + Backend<Error = E> + WireFormatBackend<Codec = C> + BackendConfig<Args = Args, Kind = Durable> + Send, Args: Send, <C as Codec<Args>>::Compact: Send, C: Codec<Args> + Clone + Send + Sync, E: Send, <C as Codec<Args>>::Error: Error + Send + Sync + 'static,

Source§

fn start_send( self: Pin<&mut S>, item: Task<Args>, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Begins sending a task to the sink. Read more
Source§

fn poll_ready( self: Pin<&mut S>, cx: &mut Context<'_>, ) -> Poll<Result<(), TaskSinkError<<S as Backend>::Error>>>

Polls the sink until it is ready to accept another task. Read more
Source§

fn poll_flush( self: Pin<&mut S>, cx: &mut Context<'_>, ) -> Poll<Result<(), TaskSinkError<<S as Backend>::Error>>>

Polls the sink until all previously submitted tasks have been flushed. Read more
Source§

fn poll_close( self: Pin<&mut S>, cx: &mut Context<'_>, ) -> Poll<Result<(), TaskSinkError<<S as Backend>::Error>>>

Polls the sink until it has been closed. Read more
Source§

async fn push( &mut self, task: Args, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Pushes a single task into the backend. Read more
Source§

async fn push_bulk( &mut self, tasks: Vec<Args>, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Pushes multiple tasks into the backend. Read more
Source§

async fn push_stream( &mut self, tasks: impl Stream<Item = Args> + Unpin + Send, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Pushes tasks from a stream into the backend. Read more
Source§

async fn push_task( &mut self, task: Task<Args>, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Pushes a fully constructed task into the backend. Read more
Source§

async fn push_all( &mut self, tasks: impl Stream<Item = Task<Args>> + Unpin + Send, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Pushes fully constructed tasks from a stream into the backend. 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<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
Source§

impl<S, Args, Compact, Err> WorkflowSink<Args> for S
where S: Send + Sink<Task<Compact>, Error = Err> + Backend<Error = Err> + WireFormatBackend<Compact = Compact> + BackendConfig + Unpin, Args: Send, <S as BackendConfig>::Id: GenerateId + Send + Sync + FromStr + Display, <S as WireFormatBackend>::Codec: Codec<Args, Compact = Compact>, Err: Error + Send + Sync + 'static, <<S as WireFormatBackend>::Codec as Codec<Args>>::Error: Into<Box<dyn Error + Sync + Send>> + Send + Sync + 'static, Compact: Send + 'static, <<S as BackendConfig>::Id as FromStr>::Err: Error + Send + Sync + 'static,

Source§

async fn push_start( &mut self, args: Args, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Push a single task into the workflow sink at the start
Source§

async fn start_fan_out( &mut self, args: Args, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>
where Args: GraphCodec<S>, <Args as GraphCodec<S>>::Error: Error + Send + Sync + 'static,

Push a single task into the workflow sink at the start
Source§

async fn push_step( &mut self, step: Args, index: usize, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Push a step into the workflow sink at the specified index Read more
Source§

async fn push_node( &mut self, node: Args, index: NodeIndex, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>

Push a node into the workflow sink at the specified index Read more