a3s_flow/worker/
memory.rs1use async_trait::async_trait;
2use std::collections::VecDeque;
3use tokio::sync::Mutex;
4use uuid::Uuid;
5
6use crate::error::{FlowError, Result};
7
8use super::{FlowTask, FlowTaskLease, FlowTaskQueue};
9
10#[derive(Debug, Default)]
12pub struct InMemoryFlowTaskQueue {
13 state: Mutex<InMemoryQueueState>,
14}
15
16#[derive(Debug, Default)]
17struct InMemoryQueueState {
18 pending: VecDeque<FlowTask>,
19 inflight: VecDeque<(String, FlowTask)>,
20}
21
22impl InMemoryFlowTaskQueue {
23 pub fn new() -> Self {
25 Self::default()
26 }
27
28 pub async fn inflight_len(&self) -> Result<usize> {
30 Ok(self.state.lock().await.inflight.len())
31 }
32}
33
34#[async_trait]
35impl FlowTaskQueue for InMemoryFlowTaskQueue {
36 async fn enqueue(&self, task: FlowTask) -> Result<()> {
37 self.state.lock().await.pending.push_back(task);
38 Ok(())
39 }
40
41 async fn lease(&self) -> Result<Option<FlowTaskLease>> {
42 let mut state = self.state.lock().await;
43 let Some(task) = state.pending.pop_front() else {
44 return Ok(None);
45 };
46 let lease_id = Uuid::new_v4().to_string();
47 state.inflight.push_back((lease_id.clone(), task.clone()));
48 Ok(Some(FlowTaskLease { lease_id, task }))
49 }
50
51 async fn heartbeat(&self, lease_id: &str) -> Result<String> {
52 let mut state = self.state.lock().await;
53 let Some((active_lease_id, _)) = state
54 .inflight
55 .iter_mut()
56 .find(|(active_lease_id, _)| active_lease_id == lease_id)
57 else {
58 return Err(FlowError::LeaseLost(lease_id.to_string()));
59 };
60 let renewed_lease_id = Uuid::new_v4().to_string();
61 *active_lease_id = renewed_lease_id.clone();
62 Ok(renewed_lease_id)
63 }
64
65 async fn ack(&self, lease_id: &str) -> Result<()> {
66 let mut state = self.state.lock().await;
67 let Some(position) = state
68 .inflight
69 .iter()
70 .position(|(active_lease_id, _)| active_lease_id == lease_id)
71 else {
72 return Err(FlowError::LeaseLost(lease_id.to_string()));
73 };
74 state.inflight.remove(position);
75 Ok(())
76 }
77
78 async fn requeue_inflight(&self) -> Result<usize> {
79 let mut state = self.state.lock().await;
80 let count = state.inflight.len();
81 while let Some((_, task)) = state.inflight.pop_front() {
82 state.pending.push_back(task);
83 }
84 Ok(count)
85 }
86
87 async fn len(&self) -> Result<usize> {
88 Ok(self.state.lock().await.pending.len())
89 }
90}