mlua_swarm/store/operator_session/
inmemory.rs1use super::{
7 Inner, OperatorSessionRecord, OperatorSessionStore, OperatorSessionStoreError, SessionId,
8 SharedInner,
9};
10use async_trait::async_trait;
11use std::sync::Mutex;
12
13#[derive(Default)]
15pub struct InMemoryOperatorSessionStore {
16 inner: SharedInner,
17}
18
19impl InMemoryOperatorSessionStore {
20 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#[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 #[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 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}