Skip to main content

a3s_flow/worker/
memory.rs

1use 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/// In-process FIFO queue for tests, embedded hosts, and local workers.
11#[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}