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 {
24 Self::default()
25 }
26
27 pub async fn inflight_len(&self) -> Result<usize> {
28 Ok(self.state.lock().await.inflight.len())
29 }
30}
31
32#[async_trait]
33impl FlowTaskQueue for InMemoryFlowTaskQueue {
34 async fn enqueue(&self, task: FlowTask) -> Result<()> {
35 self.state.lock().await.pending.push_back(task);
36 Ok(())
37 }
38
39 async fn lease(&self) -> Result<Option<FlowTaskLease>> {
40 let mut state = self.state.lock().await;
41 let Some(task) = state.pending.pop_front() else {
42 return Ok(None);
43 };
44 let lease_id = Uuid::new_v4().to_string();
45 state.inflight.push_back((lease_id.clone(), task.clone()));
46 Ok(Some(FlowTaskLease { lease_id, task }))
47 }
48
49 async fn heartbeat(&self, lease_id: &str) -> Result<String> {
50 let mut state = self.state.lock().await;
51 let Some((active_lease_id, _)) = state
52 .inflight
53 .iter_mut()
54 .find(|(active_lease_id, _)| active_lease_id == lease_id)
55 else {
56 return Err(FlowError::LeaseLost(lease_id.to_string()));
57 };
58 let renewed_lease_id = Uuid::new_v4().to_string();
59 *active_lease_id = renewed_lease_id.clone();
60 Ok(renewed_lease_id)
61 }
62
63 async fn ack(&self, lease_id: &str) -> Result<()> {
64 let mut state = self.state.lock().await;
65 let Some(position) = state
66 .inflight
67 .iter()
68 .position(|(active_lease_id, _)| active_lease_id == lease_id)
69 else {
70 return Err(FlowError::LeaseLost(lease_id.to_string()));
71 };
72 state.inflight.remove(position);
73 Ok(())
74 }
75
76 async fn requeue_inflight(&self) -> Result<usize> {
77 let mut state = self.state.lock().await;
78 let count = state.inflight.len();
79 while let Some((_, task)) = state.inflight.pop_front() {
80 state.pending.push_back(task);
81 }
82 Ok(count)
83 }
84
85 async fn len(&self) -> Result<usize> {
86 Ok(self.state.lock().await.pending.len())
87 }
88}