apalis_core/backend/
results.rs1use futures_core::Stream;
2
3use crate::{
4 backend::Backend,
5 task::{status::Status, task_id::TaskId},
6};
7
8#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
10#[derive(Debug, Clone)]
11pub struct TaskResult<T> {
12 pub task_id: TaskId,
14 pub attempt: usize,
16 pub status: Status,
18 pub result: Result<T, String>,
20}
21
22impl<T> TaskResult<T> {
23 pub fn task_id(&self) -> &TaskId {
25 &self.task_id
26 }
27
28 pub fn status(&self) -> &Status {
30 &self.status
31 }
32
33 pub fn result(&self) -> &Result<T, String> {
35 &self.result
36 }
37
38 pub fn take(self) -> Result<T, String> {
40 self.result
41 }
42}
43
44pub trait WaitForCompletion<Output>: Backend {
46 type ResultStream: Stream<Item = Result<TaskResult<Output>, Self::Error>> + Send + 'static;
48
49 fn wait_for(&self, task_ids: impl IntoIterator<Item = TaskId>) -> Self::ResultStream;
51
52 fn wait_for_single(&self, task_id: TaskId) -> Self::ResultStream {
54 self.wait_for(std::iter::once(task_id))
55 }
56
57 fn check_status(
59 &self,
60 task_ids: impl IntoIterator<Item = TaskId> + Send,
61 ) -> impl Future<Output = Result<Vec<TaskResult<Output>>, Self::Error>> + Send;
62}