1use async_trait::async_trait;
2
3use crate::error::Result;
4use crate::runtime_build::RuntimeBuildId;
5
6use super::{FlowTask, FlowTaskLease};
7
8#[async_trait]
10pub trait FlowTaskDispatcher: Send + Sync {
11 async fn dispatch(&self, task: FlowTask) -> Result<()>;
13
14 fn has_runtime_build_route(&self, required_build_id: Option<&RuntimeBuildId>) -> bool {
16 required_build_id.is_none()
17 }
18
19 fn ensure_runtime_build_route(&self, required_build_id: Option<&RuntimeBuildId>) -> Result<()> {
21 if self.has_runtime_build_route(required_build_id) {
22 return Ok(());
23 }
24 Err(crate::FlowError::RuntimeBuildRouteNotFound {
25 required_build_id: required_build_id.cloned(),
26 })
27 }
28
29 async fn dispatch_for_runtime_build(
36 &self,
37 required_build_id: Option<&RuntimeBuildId>,
38 task: FlowTask,
39 ) -> Result<()> {
40 self.ensure_runtime_build_route(required_build_id)?;
41 self.dispatch(task).await
42 }
43}
44
45#[async_trait]
47pub trait FlowTaskQueue: Send + Sync {
48 async fn enqueue(&self, task: FlowTask) -> Result<()>;
50
51 async fn lease(&self) -> Result<Option<FlowTaskLease>>;
53
54 async fn heartbeat(&self, lease_id: &str) -> Result<String>;
59
60 async fn ack(&self, lease_id: &str) -> Result<()>;
66
67 async fn requeue_inflight(&self) -> Result<usize> {
69 Ok(0)
70 }
71
72 async fn dequeue(&self) -> Result<Option<FlowTask>> {
74 let Some(lease) = self.lease().await? else {
75 return Ok(None);
76 };
77 let task = lease.task.clone();
78 self.ack(&lease.lease_id).await?;
79 Ok(Some(task))
80 }
81
82 async fn len(&self) -> Result<usize>;
84
85 async fn is_empty(&self) -> Result<bool> {
87 Ok(self.len().await? == 0)
88 }
89}
90
91#[async_trait]
92impl<T> FlowTaskDispatcher for T
93where
94 T: FlowTaskQueue + ?Sized,
95{
96 async fn dispatch(&self, task: FlowTask) -> Result<()> {
97 FlowTaskQueue::enqueue(self, task).await
98 }
99}