Skip to main content

magi_code/sessions/
manager.rs

1use 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        // Do not create JSONL until first event append; avoids orphan session files.
41        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    // Shared only by clones of this admitted standalone owner, never by independent opens.
89    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        // A poisoned optimization registry must not make durable sessions unavailable.
207        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}