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    /// Creates an empty in-process queue.
24    pub fn new() -> Self {
25        Self::default()
26    }
27
28    /// Returns the number of currently leased tasks.
29    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}