pub struct InMemoryWorkflow<Args> { /* private fields */ }Expand description
In-memory queue that is based on channels
§Example
let mut backend = InMemoryWorkflow::create();§Features
| Feature | Status | Description |
|---|---|---|
Backend | ✅ | Basic Backend functionality |
TaskSink | ✅ | Ability to push new tasks |
Serialization | ❌ | Serialization support for arguments |
PipeExt | ⚠️ | Allow other backends to pipe to this backend |
BackendFactory | ❌ | Share the same storage across multiple workers |
Update | ❌ | Allow updating a task |
FetchById | ❌ | Allow fetching a task by its ID |
Reschedule | ❌ | Reschedule a task |
ResumeById | ❌ | Resume a task by its ID |
ResumeAbandoned | ❌ | Resume abandoned tasks |
Vacuum | ❌ | Vacuum the task storage |
Workflow | ⚠️ | Flexible enough to support workflows |
WaitForCompletion | ⚠️ | Wait for tasks to complete without blocking |
RegisterWorker | ❌ | Allow registering a worker with the backend |
ListWorkers | ❌ | List all workers registered with the backend |
ListTasks | ❌ | List 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>
impl<Args> InMemoryWorkflow<Args>
Source§impl<Args> InMemoryWorkflow<Args>
impl<Args> InMemoryWorkflow<Args>
Sourcepub fn new_with(
sender: MemorySink<Vec<u8>>,
receiver: BoxedReceiver<Vec<u8>>,
) -> Shared<Self>
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>
impl<Args> Backend for InMemoryWorkflow<Args>
Source§type Error = MemoryStorageError
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>>
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>>>
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>>
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>
impl<Args> BackendConfig for InMemoryWorkflow<Args>
Source§fn middleware(&mut self, _: &mut WorkerContext) -> Self::Layer
fn middleware(&mut self, _: &mut WorkerContext) -> Self::Layer
Returns the backend’s middleware layer.
Source§impl<Args> Debug for InMemoryWorkflow<Args>
impl<Args> Debug for InMemoryWorkflow<Args>
Source§impl<Args> Sink<Task<Vec<u8>>> for InMemoryWorkflow<Args>where
Args: Unpin,
impl<Args> Sink<Task<Vec<u8>>> for InMemoryWorkflow<Args>where
Args: Unpin,
Source§type Error = MemoryStorageError
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>>
fn poll_ready( self: Pin<&mut Self>, cx: &mut Context<'_>, ) -> Poll<Result<(), Self::Error>>
Attempts to prepare the
Sink to receive a value. Read moreSource§fn start_send(
self: Pin<&mut Self>,
item: Task<Vec<u8>>,
) -> Result<(), Self::Error>
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 moreSource§impl<Args, Output> WaitForCompletion<Output> for InMemoryWorkflow<Args>
impl<Args, Output> WaitForCompletion<Output> for InMemoryWorkflow<Args>
Source§type ResultStream = WaitForStream<Output>
type ResultStream = WaitForStream<Output>
The result stream type yielding task results
Source§fn wait_for(
&self,
task_ids: impl IntoIterator<Item = TaskId>,
) -> Self::ResultStream
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
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
fn wait_for_single(&self, task_id: TaskId) -> Self::ResultStream
Wait for a single task to complete, yielding its result
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<B> BackendExt for Bwhere
B: Backend,
impl<B> BackendExt for Bwhere
B: Backend,
Source§fn poll_next_args(
&mut self,
cx: &mut Context<'_>,
worker: &WorkerContext,
) -> Poll<Option<Result<Task<Self::Args>, PollNextArgsError<Self>>>>
fn poll_next_args( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Option<Result<Task<Self::Args>, PollNextArgsError<Self>>>>
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,
fn pipe_to<Dst>(self, backend: Dst) -> Pipe<Dst, Self>where
Self: Sized,
Pipes every task polled from this backend into
sink Read moreSource§fn inspect_err<F>(self, f: F) -> InspectErr<Self, F>
fn inspect_err<F>(self, f: F) -> InspectErr<Self, F>
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>
fn map_err<F, E2>(self, f: F) -> MapErr<Self, F>
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>
fn with_codec<NewCodec>(self, codec: NewCodec) -> WithCodec<Self, NewCodec>
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>>
fn poll_with_stream<S>(self, stream: S) -> PollWith<Self, StreamStrategy<S>>
Wake the worker when a stream receives a new item
Source§fn poll_with_interval(
self,
duration: Duration,
) -> PollWith<Self, IntervalStrategy>where
Self: Sized,
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,
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,
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,
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>
fn before_start<F, Fut>(self, f: F) -> BeforeStart<Self, Self::Error>
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>
fn before_stop<F, Fut>(self, f: F) -> BeforeStop<Self, Self::Error>
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>
fn after_start<F, Fut>(self, f: F) -> AfterStart<Self, Self::Error>
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>
fn after_stop<F, Fut>(self, f: F) -> AfterStop<Self, Self::Error>
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>
fn interleave<S>(self, stream: S) -> Interleave<Self, S>
Interleaves the external stream with the backend.
Source§fn wake_on_push(self) -> WakeOnPush<Self>where
Self: Sized,
fn wake_on_push(self) -> WakeOnPush<Self>where
Self: Sized,
Wakes the worker when a new item is pushed
Create a cloneable handle to the inner backend where all handles are clone.
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
Mutably borrows from an owned value. Read more
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, Item> SinkExt<Item> for T
impl<T, Item> SinkExt<Item> for T
Source§fn with<U, Fut, F, E>(self, f: F) -> With<Self, Item, U, Fut, F>
fn with<U, Fut, F, E>(self, f: F) -> With<Self, Item, U, Fut, F>
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>
fn with_flat_map<U, St, F>(self, f: F) -> WithFlatMap<Self, Item, U, St, F>
Composes a function in front of the sink. Read more
Source§fn sink_map_err<E, F>(self, f: F) -> SinkMapErr<Self, F>
fn sink_map_err<E, F>(self, f: F) -> SinkMapErr<Self, F>
Transforms the error returned by the sink.
Source§fn sink_err_into<E>(self) -> SinkErrInto<Self, Item, E>
fn sink_err_into<E>(self) -> SinkErrInto<Self, Item, E>
Map this sink’s error to a different error type using the
Into trait. Read moreSource§fn buffer(self, capacity: usize) -> Buffer<Self, Item>where
Self: Sized,
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 flush(&mut self) -> Flush<'_, Self, Item> ⓘwhere
Self: Unpin,
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,
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,
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> ⓘ
fn send_all<'a, St>(&'a mut self, stream: &'a mut St) -> SendAll<'a, Self, St> ⓘ
A future that completes after the given stream has been fully processed
into the sink, including flushing. Read more
Source§fn right_sink<Si1>(self) -> Either<Si1, Self> ⓘ
fn right_sink<Si1>(self) -> Either<Si1, Self> ⓘ
Source§fn poll_ready_unpin(
&mut self,
cx: &mut Context<'_>,
) -> Poll<Result<(), Self::Error>>where
Self: Unpin,
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,
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§impl<Args, S, E, C> TaskSink<Args, Durable> for Swhere
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,
impl<Args, S, E, C> TaskSink<Args, Durable> for Swhere
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>>
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>>>
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>>>
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>>>
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>>
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>>
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>>
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§impl<T> WithSubscriber for T
impl<T> WithSubscriber for T
Source§fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
Source§fn with_current_subscriber(self) -> WithDispatch<Self>
fn with_current_subscriber(self) -> WithDispatch<Self>
Source§impl<S, Args, Compact, Err> WorkflowSink<Args> for Swhere
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,
impl<S, Args, Compact, Err> WorkflowSink<Args> for Swhere
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>>
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>>
async fn start_fan_out( &mut self, args: Args, ) -> Result<(), TaskSinkError<<S as Backend>::Error>>
Push a single task into the workflow sink at the start