pub struct SqliteStorage<Args> { /* private fields */ }Expand description
SqliteStorage is a storage backend for apalis using sqlite as the database.
It supports both standard polling and event-driven (hooked) storage mechanisms.
§Features
| Feature | Status | Description |
|---|---|---|
Backend | ✅ | Supports storage and retrieval of tasks |
TaskSink | ✅ | Ability to push new tasks |
Serialization | ✅ | Serialization support for arguments |
Workflow | ✅ | Flexible enough to support workflows |
Web Interface | ✅ | Expose a web interface for monitoring tasks |
FetchById | ✅ | Allow fetching a task by its ID |
RegisterWorker | ✅ | Allow registering a worker with the backend |
BackendFactory | ✅ | Share one connection across multiple workers via SqliteStorageFactory |
WaitForCompletion | ✅ | Wait for tasks to complete without blocking |
ResumeById | ✅ | Resume a task by its ID |
ResumeAbandoned | ✅ | Resume abandoned tasks |
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();
}
§Serialization
#[tokio::main]
async fn main() {
// let mut backend = /* snip */;
fn assert_codec<B: BackendConfig + WireFormatBackend>(backend: B)
where
B::Codec: Codec<B::Args, Compact=Vec<u8>>,
{
}
assert_codec(backend);
}
§Workflow
#[tokio::main]
async fn main() {
// let mut backend = /* snip */;
backend.push(42).await.unwrap();
async fn task1(task: u32, worker: WorkerContext) -> u32 {
task + 99
}
async fn task2(task: u32, worker: WorkerContext) -> u32 {
task + 1
}
async fn task3(task: u32, worker: WorkerContext) {
assert_eq!(task, 142);
worker.stop().unwrap();
}
let workflow = SteppedFlow::new("test-workflow")
.and_then(task1)
.and_then(task2)
.and_then(task3);
let worker = WorkerBuilder::new("workflow-test")
.backend(backend)
.build(workflow);
worker.run().await.unwrap();
}
§WebUI
#[tokio::main]
async fn main() {
// let mut backend = /* snip */;
fn assert_web_ui<B: Expose<u32, Durable>>(backend: B) {};
assert_web_ui(backend);
}
§WaitForCompletion
#[tokio::main]
async fn main() {
// let mut backend = /* snip */;
fn assert_wait_for_completion<B: WaitForCompletion<u32> + Backend>(backend: B) {};
assert_wait_for_completion(backend);
}
Implementations§
Source§impl SqliteStorage<()>
impl SqliteStorage<()>
Sourcepub fn connect_with_callback(
url: &str,
) -> Result<(SqlitePool, HookCallbackListener), Error>
pub fn connect_with_callback( url: &str, ) -> Result<(SqlitePool, HookCallbackListener), Error>
Connects to a database returning a pool and listener
Sourcepub fn migrations() -> Migrator
pub fn migrations() -> Migrator
Get sqlite migrations without running them
Source§impl<T> SqliteStorage<T>
impl<T> SqliteStorage<T>
Sourcepub fn new(pool: &SqlitePool) -> SqliteStorage<T>
pub fn new(pool: &SqlitePool) -> SqliteStorage<T>
Create a new SqliteStorage
Sourcepub fn with_config(self, config: Config) -> SqliteStorage<T>
pub fn with_config(self, config: Config) -> SqliteStorage<T>
Create a new SqliteStorage with a custom configuration
Sourcepub fn with_callback(
self,
callback: HookCallbackListener,
) -> PollWith<Self, StreamStrategy<HookCallbackListener>>
pub fn with_callback( self, callback: HookCallbackListener, ) -> PollWith<Self, StreamStrategy<HookCallbackListener>>
Attach a callback to an instance
Sourcepub fn pool(&self) -> &SqlitePool
pub fn pool(&self) -> &SqlitePool
Get the underlying pool
Trait Implementations§
Source§impl<Args> Backend for SqliteStorage<Args>
impl<Args> Backend for SqliteStorage<Args>
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<SqliteTask, Self::Error>>>
fn poll_next( &mut self, cx: &mut Context<'_>, worker: &WorkerContext, ) -> Poll<Option<Result<SqliteTask, 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 SqliteStorage<Args>
impl<Args> BackendConfig for SqliteStorage<Args>
Source§type Layer = TaskPersistLayer<JsonCodec<Value>, Value>
type Layer = TaskPersistLayer<JsonCodec<Value>, Value>
The type representing backend middleware layer.
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<T> Clone for SqliteStorage<T>
impl<T> Clone for SqliteStorage<T>
Source§impl<Args: Debug> Debug for SqliteStorage<Args>
impl<Args: Debug> Debug for SqliteStorage<Args>
Source§impl<Args> FetchById for SqliteStorage<Args>
impl<Args> FetchById for SqliteStorage<Args>
Source§fn fetch_by_id(
&mut self,
id: &TaskId,
) -> impl Future<Output = Result<Option<SqliteTask>, Self::Error>> + Send
fn fetch_by_id( &mut self, id: &TaskId, ) -> impl Future<Output = Result<Option<SqliteTask>, Self::Error>> + Send
Fetch a task by its unique identifier
Source§impl<Args> ListAllTasks for SqliteStorage<Args>
impl<Args> ListAllTasks for SqliteStorage<Args>
Source§fn list_all_tasks(
&self,
filter: &Filter,
) -> impl Future<Output = Result<Vec<SqliteTask>, Self::Error>> + Send
fn list_all_tasks( &self, filter: &Filter, ) -> impl Future<Output = Result<Vec<SqliteTask>, Self::Error>> + Send
List tasks matching the given filter in all queues
Source§impl<Args> ListQueues for SqliteStorage<Args>
impl<Args> ListQueues for SqliteStorage<Args>
Source§impl<Args> ListTasks for SqliteStorage<Args>
impl<Args> ListTasks for SqliteStorage<Args>
Source§fn list_tasks(
&self,
filter: &Filter,
) -> impl Future<Output = Result<Vec<SqliteTask>, Self::Error>> + Send
fn list_tasks( &self, filter: &Filter, ) -> impl Future<Output = Result<Vec<SqliteTask>, Self::Error>> + Send
List tasks matching the given filter in the current queue
Source§impl<Args: Sync> ListWorkers for SqliteStorage<Args>
impl<Args: Sync> ListWorkers for SqliteStorage<Args>
Source§fn list_workers(
&self,
) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send
fn list_workers( &self, ) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send
List all registered workers in the current queue
Source§fn list_all_workers(
&self,
) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send
fn list_all_workers( &self, ) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send
List all registered workers in all queues
Source§impl<Args> Metrics for SqliteStorage<Args>
impl<Args> Metrics for SqliteStorage<Args>
Source§impl<Args> Sink<Task<Vec<u8>>> for SqliteStorage<Args>
impl<Args> Sink<Task<Vec<u8>>> for SqliteStorage<Args>
Source§fn poll_ready(
self: Pin<&mut Self>,
_: &mut Context<'_>,
) -> Poll<Result<(), Self::Error>>
fn poll_ready( self: Pin<&mut Self>, _: &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: SqliteTask) -> Result<(), Self::Error>
fn start_send(self: Pin<&mut Self>, item: SqliteTask) -> 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<T> TryNewBackend for SqliteStorage<T>
impl<T> TryNewBackend for SqliteStorage<T>
impl<'pin, Args> Unpin for SqliteStorage<Args>where
PinnedFieldsOf<__SqliteStorage<'pin, Args>>: Unpin,
Source§impl<Args> Vacuum for SqliteStorage<Args>
impl<Args> Vacuum for SqliteStorage<Args>
Source§impl<Args, O> WaitForCompletion<O> for SqliteStorage<Args>
impl<Args, O> WaitForCompletion<O> for SqliteStorage<Args>
Source§type ResultStream = Pin<Box<dyn Stream<Item = Result<TaskResult<O>, <SqliteStorage<Args> as Backend>::Error>> + Send>>
type ResultStream = Pin<Box<dyn Stream<Item = Result<TaskResult<O>, <SqliteStorage<Args> as Backend>::Error>> + Send>>
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<O>>, Self::Error>> + Send
fn check_status( &self, task_ids: impl IntoIterator<Item = TaskId> + Send, ) -> impl Future<Output = Result<Vec<TaskResult<O>>, 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 SqliteStorage<Args>
impl<Args> !RefUnwindSafe for SqliteStorage<Args>
impl<Args> !UnwindSafe for SqliteStorage<Args>
impl<Args> Send for SqliteStorage<Args>where
PhantomData<Args>: Send,
impl<Args> Sync for SqliteStorage<Args>where
PhantomData<Args>: Sync,
impl<Args> UnsafeUnpin for SqliteStorage<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,
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,
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,
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> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
impl<B, Args, Kind> Expose<Args, Kind> for Bwhere
B: Backend + Metrics + ListWorkers + ListQueues + ListAllTasks + ListTasks + TaskSink<Args, Kind>,
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> ⓘ
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 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> ⓘ
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 moreSource§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,
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