Skip to main content

mlua_swarm/store/operator_session/
inmemory.rs

1//! `InMemoryOperatorSessionStore` — a process-volatile
2//! [`OperatorSessionStore`] used as the default when no store path is
3//! configured. Byte-for-byte the pre-persistence behaviour: sessions die
4//! with the process.
5
6use super::{
7    Inner, OperatorSessionRecord, OperatorSessionStore, OperatorSessionStoreError, SessionId,
8    SharedInner,
9};
10use async_trait::async_trait;
11use std::sync::Mutex;
12
13/// Process-volatile [`OperatorSessionStore`] default backend.
14#[derive(Default)]
15pub struct InMemoryOperatorSessionStore {
16    inner: SharedInner,
17}
18
19impl InMemoryOperatorSessionStore {
20    /// Create an empty store.
21    pub fn new() -> Self {
22        Self {
23            inner: Mutex::new(Inner::default()),
24        }
25    }
26}
27
28#[async_trait]
29impl OperatorSessionStore for InMemoryOperatorSessionStore {
30    fn name(&self) -> &str {
31        "in-memory"
32    }
33
34    async fn put(&self, record: OperatorSessionRecord) -> Result<(), OperatorSessionStoreError> {
35        let mut inner = self.inner.lock().unwrap();
36        if !inner.records.contains_key(&record.sid) {
37            inner.order.push(record.sid.clone());
38        }
39        inner.records.insert(record.sid.clone(), record);
40        Ok(())
41    }
42
43    async fn delete(&self, sid: &SessionId) -> Result<(), OperatorSessionStoreError> {
44        let mut inner = self.inner.lock().unwrap();
45        if inner.records.remove(sid).is_none() {
46            return Err(OperatorSessionStoreError::NotFound(sid.clone()));
47        }
48        inner.order.retain(|s| s != sid);
49        Ok(())
50    }
51
52    async fn list(&self) -> Result<Vec<OperatorSessionRecord>, OperatorSessionStoreError> {
53        let inner = self.inner.lock().unwrap();
54        let mut records: Vec<OperatorSessionRecord> = inner
55            .order
56            .iter()
57            .filter_map(|sid| inner.records.get(sid).cloned())
58            .collect();
59        records.sort_by_key(|r| r.joined_at_secs);
60        Ok(records)
61    }
62}
63
64// ──────────────────────────────────────────────────────────────────────────
65// tests
66// ──────────────────────────────────────────────────────────────────────────
67
68#[cfg(test)]
69mod tests {
70    use super::*;
71
72    fn mk(sid: &str, joined_at_secs: u64) -> OperatorSessionRecord {
73        OperatorSessionRecord {
74            sid: SessionId::parse(sid).unwrap(),
75            token_digest: OperatorSessionRecord::digest_of(&format!("bearer-{sid}")),
76            capability_manifest: None,
77            joined_at_secs,
78            desc: None,
79            observed: Vec::new(),
80            observed_total: 0,
81        }
82    }
83
84    /// The in-memory backend holds live records, so the 記名 an `Assign`
85    /// wrote is what the next `list()` reports — no encode/decode in
86    /// between to lose it.
87    #[tokio::test]
88    async fn the_kimei_round_trips() {
89        let s = InMemoryOperatorSessionStore::new();
90        let mut rec = mk("S-1", 100);
91        rec.desc = Some("rewriting the seat resolver in mlua-swarm-server".to_string());
92        rec.record_observed(super::super::ObservedAssignment::new(
93            "R-1".to_string(),
94            "phase-a-op".to_string(),
95            Some("resolve issue #10".to_string()),
96            Some("/repo".to_string()),
97            None,
98            None,
99            140,
100        ));
101        s.put(rec).await.unwrap();
102
103        let list = s.list().await.unwrap();
104        assert_eq!(
105            list[0].desc.as_deref(),
106            Some("rewriting the seat resolver in mlua-swarm-server")
107        );
108        assert_eq!(list[0].observed.len(), 1);
109        assert_eq!(list[0].observed_total, 1);
110        assert_eq!(list[0].last_activity_secs(), 140);
111    }
112
113    #[tokio::test]
114    async fn put_then_list() {
115        let s = InMemoryOperatorSessionStore::new();
116        s.put(mk("S-1", 100)).await.unwrap();
117        s.put(mk("S-2", 50)).await.unwrap();
118        let list = s.list().await.unwrap();
119        let sids: Vec<_> = list.iter().map(|r| r.sid.to_string()).collect();
120        assert_eq!(sids, vec!["S-2", "S-1"], "ascending by joined_at_secs");
121    }
122
123    #[tokio::test]
124    async fn put_is_upsert() {
125        let s = InMemoryOperatorSessionStore::new();
126        s.put(mk("S-1", 100)).await.unwrap();
127        let mut updated = mk("S-1", 100);
128        updated.desc = Some("the same session, re-put".to_string());
129        s.put(updated).await.unwrap();
130        let list = s.list().await.unwrap();
131        assert_eq!(list.len(), 1);
132        assert_eq!(list[0].desc.as_deref(), Some("the same session, re-put"));
133    }
134
135    #[tokio::test]
136    async fn verify_bearer_accepts_only_the_minted_bearer() {
137        let record = mk("S-1", 100);
138        assert!(record.verify_bearer("bearer-S-1"));
139        assert!(!record.verify_bearer("bearer-S-2"));
140        assert!(!record.verify_bearer(""));
141        // The stored value is the digest, not the bearer.
142        assert_ne!(record.token_digest, "bearer-S-1");
143        assert_eq!(record.token_digest.len(), 64, "hex SHA-256");
144    }
145
146    #[tokio::test]
147    async fn delete_removes_and_missing_is_not_found() {
148        let s = InMemoryOperatorSessionStore::new();
149        s.put(mk("S-1", 100)).await.unwrap();
150        s.delete(&SessionId::parse("S-1").unwrap()).await.unwrap();
151        assert!(s.list().await.unwrap().is_empty());
152        let err = s
153            .delete(&SessionId::parse("S-1").unwrap())
154            .await
155            .unwrap_err();
156        assert!(matches!(err, OperatorSessionStoreError::NotFound(_)));
157    }
158
159    #[tokio::test]
160    async fn name_is_in_memory() {
161        assert_eq!(InMemoryOperatorSessionStore::new().name(), "in-memory");
162    }
163}