magi_code/sessions/
manager.rs1use super::metadata::{SessionMetadataRecord, SessionMetadataSummary};
2use super::read::validate_session_id;
3use super::store::{prepare_session_root, primary_path};
4use crate::persistence::CrossProcessFileLock;
5use std::{
6 collections::HashMap,
7 fmt,
8 path::{Path, PathBuf},
9 sync::{Arc, Mutex, OnceLock, Weak},
10};
11use uuid::Uuid;
12
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub(crate) struct SessionInternalDiagnostic {
15 pub(crate) session_id: Option<String>,
16 pub(crate) message: String,
17}
18
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub(crate) struct SessionListReport {
21 pub(crate) summaries: Vec<SessionMetadataSummary>,
22 pub(crate) diagnostics: Vec<SessionInternalDiagnostic>,
23}
24
25#[derive(Debug, Clone)]
26pub struct SessionManager {
27 pub(in crate::sessions) root: PathBuf,
28}
29
30impl SessionManager {
31 #[must_use]
32 pub fn new(root: PathBuf) -> Self {
33 Self { root }
34 }
35
36 pub(crate) fn create(&self) -> anyhow::Result<Session> {
37 prepare_session_root(&self.root)?;
38 let id = Uuid::new_v4().to_string();
39 let path = self.path_for_valid_id(&id)?;
40 Ok(Session::new(id, path))
42 }
43
44 pub fn open(&self, id: impl Into<String>) -> anyhow::Result<Session> {
45 let id = validate_session_id(id.into())?;
46 Ok(Session::new(id.clone(), self.path_for_valid_id(&id)?))
47 }
48
49 pub(crate) fn open_existing(&self, id: impl Into<String>) -> anyhow::Result<Session> {
50 let session = self.open(id)?;
51 super::store::open_existing_primary(&self.root, &session.id)?
52 .ok_or_else(|| anyhow::anyhow!("session JSONL is missing"))?;
53 Ok(session)
54 }
55
56 #[cfg(test)]
57 pub(crate) fn list(&self) -> anyhow::Result<Vec<Session>> {
58 Ok(self
59 .list_metadata_summaries()?
60 .into_iter()
61 .map(|summary| summary.session)
62 .collect())
63 }
64
65 pub(crate) fn most_recent(&self) -> anyhow::Result<Option<Session>> {
66 let report = self.list_metadata_report()?;
67 if !report.diagnostics.is_empty() {
68 anyhow::bail!("session discovery found unreadable or unsafe history");
69 }
70 Ok(report
71 .summaries
72 .into_iter()
73 .last()
74 .map(|summary| summary.session))
75 }
76
77 pub(crate) fn path_for_valid_id(&self, id: &str) -> anyhow::Result<PathBuf> {
78 primary_path(&self.root, id)
79 }
80}
81
82pub struct Session {
83 pub(in crate::sessions) id: String,
84 pub(in crate::sessions) path: PathBuf,
85 pub(in crate::sessions) metadata_cache: Arc<Mutex<Option<SessionMetadataRecord>>>,
86 replay_state: Arc<SessionReplayState>,
87 active_lease: Arc<Mutex<Option<Arc<ActiveLeaseSlot>>>>,
88 standalone_writer: Option<Arc<CrossProcessFileLock>>,
90}
91
92struct SessionReplayState {
93 generation: std::sync::atomic::AtomicU64,
94}
95
96type ActiveLeaseSlot = Mutex<Option<Arc<CrossProcessFileLock>>>;
97
98static ACTIVE_SESSION_LEASES: OnceLock<Mutex<HashMap<PathBuf, Weak<ActiveLeaseSlot>>>> =
99 OnceLock::new();
100static SESSION_REPLAY_GENERATIONS: OnceLock<Mutex<HashMap<PathBuf, Weak<SessionReplayState>>>> =
101 OnceLock::new();
102
103impl fmt::Debug for Session {
104 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
105 formatter
106 .debug_struct("Session")
107 .field("id", &self.id)
108 .field("path", &self.path)
109 .finish()
110 }
111}
112
113impl Clone for Session {
114 fn clone(&self) -> Self {
115 Self {
116 id: self.id.clone(),
117 path: self.path.clone(),
118 metadata_cache: Arc::clone(&self.metadata_cache),
119 replay_state: Arc::clone(&self.replay_state),
120 active_lease: Arc::clone(&self.active_lease),
121 standalone_writer: self.standalone_writer.clone(),
122 }
123 }
124}
125
126impl PartialEq for Session {
127 fn eq(&self, other: &Self) -> bool {
128 self.id == other.id && self.path == other.path
129 }
130}
131
132impl Eq for Session {}
133
134#[derive(Debug)]
135pub(crate) struct SessionWriterBusy;
136impl std::fmt::Display for SessionWriterBusy {
137 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
138 f.write_str("session is busy: another process owns the writer")
139 }
140}
141impl std::error::Error for SessionWriterBusy {}
142impl Session {
143 pub(in crate::sessions) fn new(id: String, path: PathBuf) -> Self {
144 let replay_state = replay_state_for_path(&path);
145 Self {
146 id,
147 path,
148 metadata_cache: Arc::new(Mutex::new(None)),
149 replay_state,
150 active_lease: Arc::new(Mutex::new(None)),
151 standalone_writer: None,
152 }
153 }
154
155 pub(crate) fn admit_standalone_writer(mut self) -> anyhow::Result<Self> {
156 if self.standalone_writer.is_none() {
157 let writer = self
158 .try_frontend_writer()?
159 .ok_or_else(|| anyhow::Error::new(SessionWriterBusy))?;
160 self.standalone_writer = Some(Arc::new(writer));
161 }
162 self.activate()
163 }
164
165 pub(crate) fn activate(self) -> anyhow::Result<Self> {
166 self.ensure_active_lease()?;
167 Ok(self)
168 }
169
170 pub(in crate::sessions) fn ensure_active_lease(&self) -> anyhow::Result<()> {
171 let slot = {
172 let mut local_slot = self
173 .active_lease
174 .lock()
175 .map_err(|_| anyhow::anyhow!("active session lease slot was poisoned"))?;
176 if let Some(slot) = local_slot.as_ref() {
177 Arc::clone(slot)
178 } else {
179 let slot = active_lease_slot(&self.path)?;
180 *local_slot = Some(Arc::clone(&slot));
181 slot
182 }
183 };
184 let mut lease = slot
185 .lock()
186 .map_err(|_| anyhow::anyhow!("active session lease was poisoned"))?;
187 if lease.is_none() {
188 let target = active_lease_target(&self.path);
189 *lease = Some(Arc::new(
190 CrossProcessFileLock::try_acquire(&target)?
191 .ok_or_else(|| anyhow::Error::new(SessionWriterBusy))?,
192 ));
193 }
194 Ok(())
195 }
196
197 pub(in crate::sessions) fn metadata_cache(&self) -> &Mutex<Option<SessionMetadataRecord>> {
198 &self.metadata_cache
199 }
200}
201
202fn replay_state_for_path(path: &Path) -> Arc<SessionReplayState> {
203 let key = normalize_active_lease_path(path);
204 let registry = SESSION_REPLAY_GENERATIONS.get_or_init(|| Mutex::new(HashMap::new()));
205 let Ok(mut states) = registry.lock() else {
206 return Arc::new(SessionReplayState {
208 generation: std::sync::atomic::AtomicU64::new(0),
209 });
210 };
211 if let Some(state) = states.get(&key).and_then(Weak::upgrade) {
212 return state;
213 }
214 states.retain(|_, state| state.strong_count() > 0);
215 let state = Arc::new(SessionReplayState {
216 generation: std::sync::atomic::AtomicU64::new(0),
217 });
218 states.insert(key, Arc::downgrade(&state));
219 state
220}
221
222fn active_lease_slot(path: &Path) -> anyhow::Result<Arc<ActiveLeaseSlot>> {
223 let key = normalize_active_lease_path(path);
224 let registry = ACTIVE_SESSION_LEASES.get_or_init(|| Mutex::new(HashMap::new()));
225 let mut slots = registry
226 .lock()
227 .map_err(|_| anyhow::anyhow!("active session lease registry was poisoned"))?;
228 if let Some(slot) = slots.get(&key).and_then(Weak::upgrade) {
229 return Ok(slot);
230 }
231 slots.retain(|_, slot| slot.strong_count() > 0);
232 let slot = Arc::new(Mutex::new(None));
233 slots.insert(key, Arc::downgrade(&slot));
234 Ok(slot)
235}
236
237fn normalize_active_lease_path(path: &Path) -> PathBuf {
238 if let Ok(canonical) = path.canonicalize() {
239 return canonical;
240 }
241 if let (Some(parent), Some(file_name)) = (path.parent(), path.file_name())
242 && let Ok(parent) = parent.canonicalize()
243 {
244 return parent.join(file_name);
245 }
246 path.to_path_buf()
247}
248
249pub(in crate::sessions) fn active_lease_target(path: &Path) -> PathBuf {
250 path.with_extension("active")
251}
252
253impl Session {
254 pub fn id(&self) -> &str {
255 &self.id
256 }
257
258 pub(crate) fn path(&self) -> &Path {
259 &self.path
260 }
261
262 pub(crate) fn replay_generation(&self) -> u64 {
263 self.replay_state
264 .generation
265 .load(std::sync::atomic::Ordering::Acquire)
266 }
267
268 pub(in crate::sessions) fn mark_replay_changed(&self) {
269 self.replay_state
270 .generation
271 .fetch_add(1, std::sync::atomic::Ordering::AcqRel);
272 }
273}