Skip to main content

magi_code/sessions/
write.rs

1use super::event::{SessionEvent, SessionEventKind};
2use super::manager::Session;
3use super::metadata::{
4    metadata_path_for_session, read_session_metadata_for_append, report_session_diagnostic,
5    update_session_metadata_after_append_batch,
6};
7use super::read::validate_session_id;
8use crate::{
9    output::{is_credential_like_key, redact_sensitive_text},
10    persistence::CrossProcessFileLock,
11};
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub enum SessionAppendOutcome {
15    Durable,
16    DurableJsonlMetadataUpdateFailed,
17}
18
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub enum SessionAppendError {
21    NotWritten(String),
22    RolledBack(String),
23    Uncertain(String),
24}
25
26impl std::fmt::Display for SessionAppendError {
27    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
28        match self {
29            Self::NotWritten(message) => write!(formatter, "session append not written: {message}"),
30            Self::RolledBack(message) => write!(formatter, "session append rolled back: {message}"),
31            Self::Uncertain(message) => {
32                write!(
33                    formatter,
34                    "session append outcome uncertain: on-disk state uncertain: {message}"
35                )
36            }
37        }
38    }
39}
40
41impl std::error::Error for SessionAppendError {}
42use anyhow::Context;
43use serde_json::Value;
44use sha2::{Digest, Sha256};
45use std::{
46    collections::HashMap,
47    fs,
48    io::{Read, Seek, Write},
49    path::{Path, PathBuf},
50    sync::{Arc, Mutex, OnceLock, Weak},
51};
52
53pub const SESSION_TITLE_EVENT: &str = SessionEventKind::SessionTitle.as_str();
54
55/// Append a session event for paths where persistence is part of user-visible behavior.
56///
57/// Use this for CLI-managed user input, diagnostics, explicit session creation/resume flows,
58/// and other user-visible append paths where failures should be returned to the caller.
59pub fn record_session_event(
60    session: Option<&Session>,
61    cwd: &Path,
62    event_type: impl AsRef<str>,
63    payload: Value,
64) -> anyhow::Result<SessionAppendOutcome> {
65    let Some(session) = session else {
66        return Ok(SessionAppendOutcome::Durable);
67    };
68    session.append_with_outcome(&SessionEvent::new(
69        event_type.as_ref(),
70        session.id().to_string(),
71        cwd.to_path_buf(),
72        payload,
73    ))
74}
75
76/// Best-effort session event recording for runtime telemetry.
77pub fn try_record_session_event(
78    session: Option<&Session>,
79    cwd: &Path,
80    event_type: impl AsRef<str>,
81    payload: Value,
82) -> anyhow::Result<SessionAppendOutcome> {
83    record_session_event(session, cwd, event_type, payload)
84}
85
86pub fn sanitize_session_title(raw: &str) -> Option<String> {
87    let trimmed = raw
88        .trim()
89        .trim_matches(|ch| matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’'));
90    let collapsed = trimmed
91        .chars()
92        .map(|ch| {
93            if matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’') {
94                return ' ';
95            }
96            if ch.is_control() || ch.is_whitespace() {
97                ' '
98            } else {
99                ch
100            }
101        })
102        .collect::<String>();
103    let collapsed = collapsed.split_whitespace().collect::<Vec<_>>().join(" ");
104    let title = collapsed
105        .trim()
106        .trim_matches(|ch| matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’'))
107        .trim();
108    if title.is_empty() {
109        return None;
110    }
111    Some(
112        title
113            .chars()
114            .take(crate::sessions::titles::SESSION_TITLE_MAX_CHARS)
115            .collect(),
116    )
117}
118
119pub fn record_session_title(
120    session: &Session,
121    cwd: &Path,
122    title: &str,
123    provider: &str,
124    model: &str,
125) -> anyhow::Result<()> {
126    let Some(title) = sanitize_session_title(title) else {
127        anyhow::bail!("session title is empty after sanitization");
128    };
129    session.append(&SessionEvent::new_kind(
130        SessionEventKind::SessionTitle,
131        session.id().to_string(),
132        cwd.to_path_buf(),
133        serde_json::json!({"title": title, "provider": provider, "model": model}),
134    ))
135}
136
137pub(crate) const COMPACTION_SCHEMA_VERSION: u64 = 1;
138
139pub(crate) fn sanitize_compaction_summary(raw: &str) -> Option<String> {
140    let cleaned = raw
141        .trim()
142        .chars()
143        .map(|ch| {
144            if ch.is_control() && !matches!(ch, '\n' | '\t') {
145                ' '
146            } else {
147                ch
148            }
149        })
150        .collect::<String>();
151    let redacted = redact_sensitive_text(cleaned.trim());
152    if redacted.trim().is_empty() {
153        None
154    } else {
155        Some(redacted.trim().to_string())
156    }
157}
158
159#[cfg(test)]
160pub(crate) fn record_session_compaction(
161    session: &Session,
162    cwd: &Path,
163    summary: &str,
164    provider: &str,
165    model: &str,
166    cutoff_event_count: usize,
167) -> anyhow::Result<()> {
168    let bytes = session
169        .path
170        .exists()
171        .then(|| fs::read(&session.path))
172        .transpose()?
173        .unwrap_or_default();
174    let offset = bytes
175        .split_inclusive(|byte| *byte == b'\n')
176        .take(cutoff_event_count)
177        .map(|line| line.len())
178        .sum::<usize>();
179    record_session_compaction_at_byte_offset_with_snapshot(
180        session,
181        cwd,
182        summary,
183        provider,
184        model,
185        offset as u64,
186        capture_session_snapshot(&session.path)?,
187    )
188}
189
190pub(crate) fn record_session_compaction_at_byte_offset_with_snapshot(
191    session: &Session,
192    cwd: &Path,
193    summary: &str,
194    provider: &str,
195    model: &str,
196    cutoff_byte_offset: u64,
197    snapshot: Option<SessionSnapshot>,
198) -> anyhow::Result<()> {
199    let Some(summary) = sanitize_compaction_summary(summary) else {
200        anyhow::bail!("compaction summary is empty after sanitization");
201    };
202    let provider = provider.trim();
203    let model = model.trim();
204    if provider.is_empty() || model.is_empty() {
205        anyhow::bail!("compaction provider and model must be non-empty");
206    }
207    rotate_session_history(
208        session,
209        cwd,
210        &summary,
211        provider,
212        model,
213        cutoff_byte_offset,
214        snapshot,
215    )
216}
217
218fn rotate_session_history(
219    session: &Session,
220    cwd: &Path,
221    summary: &str,
222    provider: &str,
223    model: &str,
224    cutoff_byte_offset: u64,
225    snapshot: Option<SessionSnapshot>,
226) -> anyhow::Result<()> {
227    // Snapshot captured before provider work; validation under both locks rejects rotation/prune replacement.
228    validate_session_id(session.id.clone())?;
229    let parent = session
230        .path
231        .parent()
232        .ok_or_else(|| anyhow::anyhow!("session file has no parent directory"))?;
233    ensure_directory(parent)?;
234    let append_lock = session_append_lock(&session.path)?;
235    let _append_guard = append_lock
236        .lock()
237        .map_err(|_| anyhow::anyhow!("session append lock was poisoned"))?;
238    let _file_guard = CrossProcessFileLock::acquire(&session.path)?;
239    validate_session_snapshot(&session.path, snapshot.as_ref())?;
240    let previous_metadata = if session.path.exists() {
241        match super::metadata::read_complete_session_metadata(session)? {
242            Some(record) => Some(record),
243            None => Some(super::metadata::rebuild_session_metadata_from_jsonl(
244                session,
245            )?),
246        }
247    } else {
248        None
249    };
250    let old_len = open_secure_session_path(&session.path)?
251        .as_ref()
252        .map(|file| file.metadata().map(|metadata| metadata.len()))
253        .transpose()?
254        .unwrap_or(0);
255    let suffix_start = cutoff_byte_offset.min(old_len);
256    let checkpoint = SessionEvent::new_kind(
257        SessionEventKind::Compaction,
258        session.id.to_string(),
259        cwd.to_path_buf(),
260        serde_json::json!({"schema_version": COMPACTION_SCHEMA_VERSION, "summary": summary, "provider": provider, "model": model, "cutoff_event_count": 0, "aggregate": previous_metadata.as_ref().map(|record| record.checkpoint_aggregate()).unwrap_or_else(|| serde_json::json!({}))}),
261    );
262    let history_parent = parent.join(".history");
263    ensure_directory_or_create(&history_parent)?;
264    let history_root = history_parent.join(&session.id);
265    ensure_directory_or_create(&history_root)?;
266    remove_stale_rotation_temps(parent, &session.id)?;
267    let generation = fs::read_dir(&history_root)?
268        .filter_map(Result::ok)
269        .filter_map(|entry| {
270            entry
271                .file_name()
272                .to_str()
273                .and_then(|name| name.strip_suffix(".jsonl"))
274                .and_then(|name| name.parse::<u64>().ok())
275        })
276        .max()
277        .unwrap_or(0)
278        .saturating_add(1);
279    remove_stale_rotation_archive_temps(&history_root, generation)?;
280    let archive = history_root.join(format!("{generation}.jsonl"));
281    let archive_temp = history_root.join(format!(".{generation}.jsonl.compact-tmp"));
282    let archive_result = (|| -> anyhow::Result<()> {
283        let mut archive_file = fs::OpenOptions::new();
284        archive_file.write(true).create_new(true);
285        #[cfg(unix)]
286        {
287            use std::os::unix::fs::OpenOptionsExt;
288            archive_file.mode(super::store::SESSION_FILE_MODE);
289        }
290        let mut archive_file = archive_file.open(&archive_temp)?;
291        if let Some(mut source) = open_secure_session_path(&session.path)? {
292            std::io::copy(&mut source, &mut archive_file)?;
293        }
294        archive_file.sync_all()?;
295        fs::rename(&archive_temp, &archive)?;
296        sync_session_parent_dir(&history_root)?;
297        Ok(())
298    })();
299    if let Err(error) = archive_result {
300        remove_validated_rotation_temp(&archive_temp);
301        return Err(error);
302    }
303    let checkpoint_line = serialize_session_event_line(&checkpoint)?;
304    let temporary = parent.join(format!(
305        ".{}.jsonl.compact-{}-{}",
306        session.id,
307        std::process::id(),
308        generation
309    ));
310    let mut active_file = fs::OpenOptions::new();
311    active_file.write(true).create_new(true);
312    #[cfg(unix)]
313    {
314        use std::os::unix::fs::OpenOptionsExt;
315        active_file.mode(super::store::SESSION_FILE_MODE);
316    }
317    let mut active_file = active_file.open(&temporary)?;
318    active_file.write_all(&checkpoint_line)?;
319    if let Some(mut source) = open_secure_session_path(&session.path)? {
320        source.seek(std::io::SeekFrom::Start(suffix_start))?;
321        std::io::copy(&mut source, &mut active_file)?;
322    }
323    active_file.sync_all()?;
324    if let Err(error) = fs::rename(&temporary, &session.path) {
325        remove_validated_rotation_temp(&temporary);
326        return Err(error.into());
327    }
328    sync_session_parent_dir(parent).context("active session replacement committed but directory sync failed; active JSONL remains recoverable")?;
329    if let Some(previous) = previous_metadata
330        && let Err(error) =
331            super::metadata::write_rotated_session_metadata(session, previous, &checkpoint)
332    {
333        report_session_diagnostic(
334            super::metadata::SessionDiagnosticOperation::Metadata,
335            &metadata_path_for_session(session),
336            error,
337        );
338    }
339    Ok(())
340}
341
342fn open_secure_session_path(path: &Path) -> anyhow::Result<Option<fs::File>> {
343    let root = path
344        .parent()
345        .ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
346    let name = path
347        .file_name()
348        .and_then(|name| name.to_str())
349        .ok_or_else(|| anyhow::anyhow!("session file has no filename"))?;
350    let id = name
351        .strip_suffix(".jsonl")
352        .ok_or_else(|| anyhow::anyhow!("session file has unexpected name"))?;
353    super::store::open_existing_primary(root, id)
354}
355
356#[derive(Debug)]
357pub(crate) struct SessionSnapshot {
358    len: u64,
359    prefix: Vec<u8>,
360    digest: [u8; 32],
361    identity: Option<(u64, u64)>,
362}
363
364pub(crate) fn capture_session_snapshot(path: &Path) -> anyhow::Result<Option<SessionSnapshot>> {
365    let Some(mut file) = open_secure_session_path(path)? else {
366        return Ok(None);
367    };
368    let metadata = file.metadata()?;
369    let mut prefix = vec![0; usize::try_from(metadata.len().min(4096)).unwrap_or(4096)];
370    file.read_exact(&mut prefix)?;
371    let mut digest = Sha256::new();
372    digest.update(&prefix);
373    let mut remaining = metadata.len().saturating_sub(prefix.len() as u64);
374    let mut buffer = [0u8; 8192];
375    while remaining > 0 {
376        let read_len = usize::try_from(remaining.min(buffer.len() as u64)).unwrap_or(buffer.len());
377        file.read_exact(&mut buffer[..read_len])?;
378        digest.update(&buffer[..read_len]);
379        remaining -= read_len as u64;
380    }
381    Ok(Some(SessionSnapshot {
382        len: metadata.len(),
383        prefix,
384        digest: digest.finalize().into(),
385        identity: file_identity(&metadata),
386    }))
387}
388
389fn validate_session_snapshot(
390    path: &Path,
391    snapshot: Option<&SessionSnapshot>,
392) -> anyhow::Result<()> {
393    let current = capture_session_snapshot(path)?;
394    match (snapshot, current.as_ref()) {
395        (None, None) => Ok(()),
396        (Some(expected), Some(actual))
397            if actual.len >= expected.len
398                && actual.prefix.starts_with(&expected.prefix)
399                && (expected.identity == actual.identity
400                    || (expected.identity.is_none()
401                        && actual.len == expected.len
402                        && actual.digest == expected.digest)) =>
403        {
404            Ok(())
405        }
406        _ => anyhow::bail!(
407            "session changed during compaction; rotation aborted safely (snapshot/current generation mismatch)"
408        ),
409    }
410}
411
412fn file_identity(metadata: &fs::Metadata) -> Option<(u64, u64)> {
413    #[cfg(unix)]
414    {
415        use std::os::unix::fs::MetadataExt;
416        Some((metadata.dev(), metadata.ino()))
417    }
418    #[cfg(not(unix))]
419    {
420        let _ = metadata;
421        None
422    }
423}
424
425fn remove_validated_rotation_temp(path: &Path) {
426    if let Ok(metadata) = fs::symlink_metadata(path)
427        && metadata.file_type().is_file()
428        && !metadata.file_type().is_symlink()
429    {
430        let _ = fs::remove_file(path);
431    }
432}
433
434fn ensure_directory(path: &Path) -> anyhow::Result<()> {
435    let metadata = fs::symlink_metadata(path)?;
436    if !metadata.file_type().is_dir() || metadata.file_type().is_symlink() {
437        anyhow::bail!("session history path is not a directory");
438    }
439    Ok(())
440}
441
442fn ensure_directory_or_create(path: &Path) -> anyhow::Result<()> {
443    if !path.exists() {
444        fs::create_dir(path)?;
445        if let Some(parent) = path.parent() {
446            sync_session_parent_dir(parent)?;
447        }
448    }
449    ensure_directory(path)
450}
451
452fn remove_stale_rotation_temps(parent: &Path, id: &str) -> anyhow::Result<()> {
453    let prefix = format!(".{id}.jsonl.compact-");
454    for entry in fs::read_dir(parent)? {
455        let entry = entry?;
456        if !entry.file_name().to_string_lossy().starts_with(&prefix) {
457            continue;
458        }
459        let metadata = fs::symlink_metadata(entry.path())?;
460        if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
461            continue;
462        }
463        fs::remove_file(entry.path())?;
464    }
465    Ok(())
466}
467fn remove_stale_rotation_archive_temps(history_root: &Path, generation: u64) -> anyhow::Result<()> {
468    let prefix = format!(".{generation}.jsonl.compact-tmp");
469    for entry in fs::read_dir(history_root)? {
470        let entry = entry?;
471        if !entry.file_name().to_string_lossy().starts_with(&prefix) {
472            continue;
473        }
474        let metadata = fs::symlink_metadata(entry.path())?;
475        if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
476            continue;
477        }
478        fs::remove_file(entry.path())?;
479    }
480    Ok(())
481}
482const MAX_REDACTION_DEPTH: usize = 128;
483fn redact_json_value(value: Value) -> Value {
484    redact_json_value_at_depth(value, 0)
485}
486
487fn redact_json_value_at_depth(value: Value, depth: usize) -> Value {
488    match value {
489        Value::String(text) => Value::String(redact_sensitive_text(&text)),
490        Value::Array(items) => {
491            if depth >= MAX_REDACTION_DEPTH {
492                return Value::String("<redacted: max depth>".to_string());
493            }
494            Value::Array(
495                items
496                    .into_iter()
497                    .map(|value| redact_json_value_at_depth(value, depth + 1))
498                    .collect(),
499            )
500        }
501        Value::Object(map) => {
502            if depth >= MAX_REDACTION_DEPTH {
503                return Value::String("<redacted: max depth>".to_string());
504            }
505            Value::Object(
506                map.into_iter()
507                    .map(|(key, value)| {
508                        let redacted = if is_credential_like_key(&key) {
509                            match value {
510                                Value::String(_) => Value::String("<redacted>".to_string()),
511                                other => redact_json_value_at_depth(other, depth + 1),
512                            }
513                        } else {
514                            redact_json_value_at_depth(value, depth + 1)
515                        };
516                        (key, redacted)
517                    })
518                    .collect(),
519            )
520        }
521        other => other,
522    }
523}
524trait DurableSessionWriter {
525    fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()>;
526    fn flush_durable(&mut self) -> std::io::Result<()>;
527    fn sync_data_durable(&mut self) -> std::io::Result<()>;
528    fn set_len_durable(&mut self, len: u64) -> std::io::Result<()>;
529}
530
531impl DurableSessionWriter for fs::File {
532    fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
533        self.write_all(bytes)
534    }
535
536    fn flush_durable(&mut self) -> std::io::Result<()> {
537        self.flush()
538    }
539
540    fn sync_data_durable(&mut self) -> std::io::Result<()> {
541        self.sync_data()
542    }
543
544    fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
545        self.set_len(len)
546    }
547}
548
549fn serialize_session_event_line(event: &SessionEvent) -> anyhow::Result<Vec<u8>> {
550    let mut line = serde_json::to_vec(event).context("failed to serialize session event")?;
551    line.push(b'\n');
552    Ok(line)
553}
554
555fn rollback_session_append(writer: &mut dyn DurableSessionWriter, pre_append_len: u64) -> bool {
556    writer.set_len_durable(pre_append_len).is_ok()
557}
558
559fn append_event_line_durably(
560    writer: &mut dyn DurableSessionWriter,
561    line: &[u8],
562    pre_append_len: u64,
563) -> Result<(), SessionAppendError> {
564    let append = |writer: &mut dyn DurableSessionWriter, error: std::io::Error, stage: &str| {
565        let message = format!("failed to {stage} session JSONL event: {error}");
566        if rollback_session_append(writer, pre_append_len) {
567            SessionAppendError::RolledBack(message)
568        } else {
569            SessionAppendError::Uncertain(message)
570        }
571    };
572    if let Err(error) = writer.write_all_durable(line) {
573        return Err(append(writer, error, "write"));
574    }
575    if let Err(error) = writer.flush_durable() {
576        return Err(append(writer, error, "flush"));
577    }
578    if let Err(error) = writer.sync_data_durable() {
579        return Err(append(writer, error, "sync"));
580    }
581    Ok(())
582}
583
584fn sync_session_parent_dir(parent: &Path) -> std::io::Result<()> {
585    #[cfg(unix)]
586    {
587        fs::File::open(parent)?.sync_all()?;
588    }
589    #[cfg(not(unix))]
590    {
591        let _ = parent;
592    }
593    Ok(())
594}
595
596static SESSION_APPEND_LOCKS: OnceLock<Mutex<HashMap<PathBuf, Weak<Mutex<()>>>>> = OnceLock::new();
597
598fn session_append_lock(path: &Path) -> anyhow::Result<Arc<Mutex<()>>> {
599    let key = normalize_session_append_lock_path(path);
600    let registry = SESSION_APPEND_LOCKS.get_or_init(|| Mutex::new(HashMap::new()));
601    let mut locks = registry
602        .lock()
603        .map_err(|_| anyhow::anyhow!("session append lock registry was poisoned"))?;
604    if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) {
605        return Ok(lock);
606    }
607    locks.retain(|_, lock| lock.strong_count() > 0);
608    let lock = Arc::new(Mutex::new(()));
609    locks.insert(key, Arc::downgrade(&lock));
610    Ok(lock)
611}
612
613fn normalize_session_append_lock_path(path: &Path) -> PathBuf {
614    if let Ok(canonical) = path.canonicalize() {
615        return canonical;
616    }
617    if let (Some(parent), Some(file_name)) = (path.parent(), path.file_name())
618        && let Ok(parent) = parent.canonicalize()
619    {
620        return parent.join(file_name);
621    }
622    path.to_path_buf()
623}
624
625impl Session {
626    pub fn append(&self, event: &SessionEvent) -> anyhow::Result<()> {
627        self.append_with_outcome(event).map(|_| ())
628    }
629
630    pub fn append_with_outcome(
631        &self,
632        event: &SessionEvent,
633    ) -> anyhow::Result<SessionAppendOutcome> {
634        self.append_owned_batch(vec![event.clone()])
635            .map_err(anyhow::Error::new)
636    }
637
638    fn append_owned_batch(
639        &self,
640        mut events: Vec<SessionEvent>,
641    ) -> Result<SessionAppendOutcome, SessionAppendError> {
642        let not_written =
643            |error: &dyn std::fmt::Display| SessionAppendError::NotWritten(error.to_string());
644        validate_session_id(self.id.clone()).map_err(|error| not_written(&error))?;
645        let parent = self.path.parent().ok_or_else(|| {
646            SessionAppendError::NotWritten("session file has no parent".to_string())
647        })?;
648        super::store::prepare_session_root(parent).map_err(|error| not_written(&error))?;
649        let mut line = Vec::new();
650        for event in &mut events {
651            validate_session_id(event.session_id.clone()).map_err(|error| not_written(&error))?;
652            if event.session_id != self.id {
653                return Err(SessionAppendError::NotWritten(
654                    "event session id does not match session".to_string(),
655                ));
656            }
657            event.payload = redact_json_value(std::mem::take(&mut event.payload));
658            line.extend(serialize_session_event_line(event).map_err(|error| not_written(&error))?);
659        }
660        let append_lock = session_append_lock(&self.path).map_err(|error| not_written(&error))?;
661        let _append_guard = append_lock.lock().map_err(|_| {
662            SessionAppendError::NotWritten("session append lock was poisoned".to_string())
663        })?;
664        let _file_guard =
665            CrossProcessFileLock::acquire(&self.path).map_err(|error| not_written(&error))?;
666        let previous_metadata = read_session_metadata_for_append(self).ok().flatten();
667        let (mut file, is_new_session_file) =
668            super::store::open_primary(parent, &self.id).map_err(|error| not_written(&error))?;
669        let pre_append_len = file.metadata().map_err(|error| not_written(&error))?.len();
670        append_event_line_durably(&mut file, &line, pre_append_len)?;
671        if is_new_session_file {
672            sync_session_parent_dir(parent)
673                .map_err(|error| SessionAppendError::Uncertain(error.to_string()))?;
674        }
675        if let Err(error) =
676            update_session_metadata_after_append_batch(self, previous_metadata, &events)
677        {
678            report_session_diagnostic(
679                super::metadata::SessionDiagnosticOperation::Metadata,
680                &metadata_path_for_session(self),
681                &error,
682            );
683            return Ok(SessionAppendOutcome::DurableJsonlMetadataUpdateFailed);
684        }
685        Ok(SessionAppendOutcome::Durable)
686    }
687}
688pub(crate) fn try_record_session_event_batch(
689    session: Option<&Session>,
690    cwd: &Path,
691    event_type: impl AsRef<str>,
692    payloads: impl IntoIterator<Item = Value>,
693) -> anyhow::Result<SessionAppendOutcome> {
694    let Some(session) = session else {
695        return Ok(SessionAppendOutcome::Durable);
696    };
697    let event_type = event_type.as_ref();
698    let events = payloads
699        .into_iter()
700        .map(|payload| {
701            SessionEvent::new(
702                event_type,
703                session.id().to_string(),
704                cwd.to_path_buf(),
705                payload,
706            )
707        })
708        .collect();
709    session
710        .append_owned_batch(events)
711        .map_err(anyhow::Error::new)
712}
713
714#[cfg(test)]
715mod tests {
716    use super::super::manager::SessionManager;
717    use super::*;
718    use serde_json::json;
719    use tempfile::TempDir;
720
721    #[test]
722    fn session_jsonl_append_and_read_round_trip() {
723        let temp = TempDir::new().unwrap();
724        let manager = SessionManager::new(temp.path().join("sessions"));
725        let session = manager.create().unwrap();
726        let event = SessionEvent::new(
727            "user_input",
728            session.id().to_string(),
729            temp.path().to_path_buf(),
730            json!({"text":"hello"}),
731        );
732        session.append(&event).unwrap();
733
734        #[cfg(unix)]
735        {
736            use std::os::unix::fs::MetadataExt;
737            assert_eq!(
738                fs::symlink_metadata(session.path().parent().unwrap())
739                    .unwrap()
740                    .mode()
741                    & 0o777,
742                0o700
743            );
744            assert_eq!(
745                fs::symlink_metadata(session.path()).unwrap().mode() & 0o777,
746                0o600
747            );
748        }
749
750        let events = session.read_events().unwrap();
751        assert_eq!(events.len(), 1);
752        assert_eq!(events[0].event_type, "user_input");
753        assert_eq!(events[0].payload["text"], "hello");
754        assert_eq!(events[0].session_path.as_deref(), Some(temp.path()));
755    }
756
757    #[test]
758    fn compaction_rotates_history_and_keeps_active_suffix_bounded() {
759        let temp = TempDir::new().unwrap();
760        let manager = SessionManager::new(temp.path().join("sessions"));
761        let session = manager.create().unwrap();
762        session
763            .append(&SessionEvent::new(
764                "user_input",
765                session.id().to_string(),
766                temp.path().to_path_buf(),
767                json!({"text":"before"}),
768            ))
769            .unwrap();
770        let original = fs::read(session.path()).unwrap();
771        record_session_compaction(&session, temp.path(), "summary", "p", "m", 1).unwrap();
772
773        let history = temp
774            .path()
775            .join("sessions/.history")
776            .join(session.id())
777            .join("1.jsonl");
778        assert_eq!(fs::read(history).unwrap(), original);
779        let active = session.read_events().unwrap();
780        assert_eq!(active[0].event_type, "compaction");
781        assert_eq!(active[0].payload["cutoff_event_count"], 0);
782
783        assert_eq!(active.len(), 1);
784
785        session
786            .append(&SessionEvent::new(
787                "user_input",
788                session.id().to_string(),
789                temp.path().to_path_buf(),
790                json!({"text":"after"}),
791            ))
792            .unwrap();
793        record_session_compaction(&session, temp.path(), "summary 2", "p", "m", 2).unwrap();
794        assert!(
795            temp.path()
796                .join("sessions/.history")
797                .join(session.id())
798                .join("2.jsonl")
799                .exists()
800        );
801        assert!(
802            manager
803                .list()
804                .unwrap()
805                .iter()
806                .all(|item| item.id() == session.id())
807        );
808    }
809    #[test]
810    fn snapshot_rejects_replacement_with_shared_prefix_and_growth() {
811        let temp = TempDir::new().unwrap();
812        let path = temp.path().join("session.jsonl");
813        fs::write(&path, b"shared prefix\noriginal generation\n").unwrap();
814        crate::sessions::store::secure_test_session_root(temp.path());
815        let snapshot = capture_session_snapshot(&path).unwrap();
816        let backup = temp.path().join("session.backup");
817        fs::rename(&path, backup).unwrap();
818        fs::write(
819            &path,
820            b"shared prefix\nreplacement generation with growth\n",
821        )
822        .unwrap();
823        crate::sessions::store::secure_test_session_root(temp.path());
824
825        let error = validate_session_snapshot(&path, Some(snapshot.as_ref().unwrap())).unwrap_err();
826        assert!(error.to_string().contains("generation mismatch"));
827    }
828
829    #[test]
830    fn session_append_lock_is_shared_per_jsonl_path() {
831        let temp = TempDir::new().unwrap();
832        let path = temp.path().join("session.jsonl");
833
834        let first = session_append_lock(&path).unwrap();
835        let second = session_append_lock(&path).unwrap();
836        let other = session_append_lock(&temp.path().join("other.jsonl")).unwrap();
837
838        assert!(Arc::ptr_eq(&first, &second));
839        assert!(!Arc::ptr_eq(&first, &other));
840    }
841
842    #[test]
843    fn concurrent_session_appends_remain_parseable_jsonl() {
844        let temp = TempDir::new().unwrap();
845        let manager = SessionManager::new(temp.path().join("sessions"));
846        let session = manager.create().unwrap();
847        let mut handles = Vec::new();
848        for index in 0..24 {
849            let session = session.clone();
850            let cwd = temp.path().to_path_buf();
851            handles.push(std::thread::spawn(move || {
852                let event = SessionEvent::new(
853                    "diagnostic",
854                    session.id().to_string(),
855                    cwd,
856                    json!({"index": index}),
857                );
858                session.append(&event).unwrap();
859            }));
860        }
861        for handle in handles {
862            handle.join().unwrap();
863        }
864
865        let tolerant = session.read_events_tolerant().unwrap();
866        assert!(
867            tolerant.diagnostics.is_empty(),
868            "{:?}",
869            tolerant.diagnostics
870        );
871        assert_eq!(tolerant.events.len(), 24);
872    }
873
874    #[test]
875    fn session_append_rejects_mismatched_or_unsafe_event_ids() {
876        let temp = TempDir::new().unwrap();
877        let manager = SessionManager::new(temp.path().join("sessions"));
878        let session = manager.open("safe").unwrap();
879        let mismatched = SessionEvent::new(
880            "user_input",
881            "other".to_string(),
882            temp.path().to_path_buf(),
883            json!({"text":"hello"}),
884        );
885        assert!(session.append(&mismatched).is_err());
886
887        let unsafe_event = SessionEvent::new(
888            "user_input",
889            "../escape".to_string(),
890            temp.path().to_path_buf(),
891            json!({"text":"hello"}),
892        );
893        assert!(session.append(&unsafe_event).is_err());
894    }
895
896    #[test]
897    fn record_session_event_propagates_append_failure() {
898        let temp = TempDir::new().unwrap();
899        let invalid_session =
900            Session::unchecked_for_test("../unsafe".to_string(), temp.path().join("unsafe.jsonl"));
901
902        let error = record_session_event(
903            Some(&invalid_session),
904            temp.path(),
905            "diagnostic",
906            json!({"message":"persist me"}),
907        )
908        .unwrap_err();
909
910        assert!(error.to_string().contains("must not contain '..'"));
911        assert!(!invalid_session.path().exists());
912    }
913
914    #[test]
915    fn metadata_sidecar_write_failure_warns_but_append_succeeds() {
916        let temp = TempDir::new().unwrap();
917        let manager = SessionManager::new(temp.path().join("sessions"));
918        let session = manager.open("metadata-fails").unwrap();
919        fs::create_dir_all(metadata_path_for_session(&session)).unwrap();
920        crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
921
922        session
923            .append(&SessionEvent::new(
924                "user_input",
925                session.id().to_string(),
926                temp.path().to_path_buf(),
927                json!({"text":"durable"}),
928            ))
929            .unwrap();
930
931        let events = session.read_events().unwrap();
932        assert_eq!(events.len(), 1);
933        assert_eq!(events[0].payload["text"], "durable");
934        assert!(metadata_path_for_session(&session).is_dir());
935    }
936
937    #[test]
938    fn redacted_session_events_remain_readable_for_resume() {
939        let temp = TempDir::new().unwrap();
940        let manager = SessionManager::new(temp.path().join("sessions"));
941        let session = manager.create().unwrap();
942        let secret = "sk-testSecret123456";
943        session
944            .append(&SessionEvent::new(
945                "tool_result",
946                session.id().to_string(),
947                temp.path().to_path_buf(),
948                json!({"tool":"bash", "success": true, "output": secret}),
949            ))
950            .unwrap();
951        let events = session.read_events().unwrap();
952        assert_eq!(events[0].event_type, "tool_result");
953        assert_eq!(events[0].payload["tool"], "bash");
954        assert_eq!(events[0].payload["success"], true);
955        assert!(!events[0].payload.to_string().contains(secret));
956    }
957
958    #[test]
959    fn persisted_redaction_preserves_generic_fields_and_hides_credentials() {
960        let temp = TempDir::new().unwrap();
961        let manager = SessionManager::new(temp.path().join("sessions"));
962        let session = manager.create().unwrap();
963        session
964            .append(&SessionEvent::new(
965                "tool_result",
966                session.id().to_string(),
967                temp.path().to_path_buf(),
968                json!({
969                    "code": 200,
970                    "key": "value",
971                    "nested": {"code": "ok", "api_key": "api-secret"},
972                    "provider_api_key": "provider-secret"
973                }),
974            ))
975            .unwrap();
976
977        let payload = session.read_events().unwrap().remove(0).payload;
978        assert_eq!(payload["code"], 200);
979        assert_eq!(payload["key"], "value");
980        assert_eq!(payload["nested"]["code"], "ok");
981        assert_eq!(payload["nested"]["api_key"], "<redacted>");
982        assert_eq!(payload["provider_api_key"], "<redacted>");
983
984        let persisted = fs::read_to_string(session.path()).unwrap();
985        assert!(persisted.contains("\"code\":200"));
986        assert!(persisted.contains("\"key\":\"value\""));
987        assert!(!persisted.contains("api-secret"));
988        assert!(!persisted.contains("provider-secret"));
989    }
990
991    #[test]
992    fn redaction_preserves_numeric_token_usage_fields() {
993        let temp = TempDir::new().unwrap();
994        let manager = SessionManager::new(temp.path().join("sessions"));
995        let session = manager.create().unwrap();
996        session
997            .append(&SessionEvent::new(
998                "usage",
999                session.id().to_string(),
1000                temp.path().to_path_buf(),
1001                json!({
1002                    "input_tokens": 12,
1003                    "output_tokens": 34,
1004                    "cache_creation_input_tokens": 56,
1005                    "token_estimate": 78,
1006                    "access_token": "secret-token",
1007                }),
1008            ))
1009            .unwrap();
1010        let payload = session.read_events().unwrap().remove(0).payload;
1011        assert_eq!(payload["input_tokens"], 12);
1012        assert_eq!(payload["output_tokens"], 34);
1013        assert_eq!(payload["cache_creation_input_tokens"], 56);
1014        assert_eq!(payload["token_estimate"], 78);
1015        assert_eq!(payload["access_token"], "<redacted>");
1016    }
1017
1018    #[test]
1019    fn provider_stream_trace_session_event_is_redacted_before_write() {
1020        let temp = TempDir::new().unwrap();
1021        let manager = SessionManager::new(temp.path().join("sessions"));
1022        let session = manager.create().unwrap();
1023        session
1024            .append(&SessionEvent::new_kind(
1025                SessionEventKind::ProviderStreamTrace,
1026                session.id().to_string(),
1027                temp.path().to_path_buf(),
1028                json!({
1029                    "schema_version": 1,
1030                    "provider": "anthropic",
1031                    "api_key": "sk-secret123456",
1032                    "recent_events": [{"seq": 1, "usage": {"output_tokens": 7}}],
1033                    "pending_tools": [{
1034                        "index": 1,
1035                        "id": "toolu_1",
1036                        "name": "read",
1037                        "argument_bytes": 13,
1038                        "argument_sha256": "a".repeat(64)
1039                    }]
1040                }),
1041            ))
1042            .unwrap();
1043
1044        let payload = session.read_events().unwrap().remove(0).payload;
1045        assert_eq!(payload["api_key"], "<redacted>");
1046        assert_eq!(payload["schema_version"], 1);
1047        assert_eq!(payload["recent_events"][0]["seq"], 1);
1048        assert_eq!(payload["recent_events"][0]["usage"]["output_tokens"], 7);
1049        assert_eq!(payload["pending_tools"][0]["argument_bytes"], 13);
1050    }
1051
1052    #[test]
1053    fn redaction_caps_deep_json_nesting() {
1054        let mut value = json!({"api_key":"sk-testSecret123456"});
1055        for _ in 0..(MAX_REDACTION_DEPTH + 8) {
1056            value = json!([value]);
1057        }
1058
1059        let redacted = redact_json_value(value);
1060        let text = redacted.to_string();
1061
1062        assert!(text.contains("<redacted: max depth>"), "{text}");
1063        assert!(!text.contains("sk-testSecret"), "{text}");
1064    }
1065
1066    #[test]
1067    fn session_append_failure_rolls_back_partial_jsonl_tail() {
1068        struct PartialFailWriter {
1069            bytes: Vec<u8>,
1070        }
1071
1072        impl DurableSessionWriter for PartialFailWriter {
1073            fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
1074                self.bytes.extend_from_slice(&bytes[..bytes.len() / 2]);
1075                Err(std::io::Error::other("simulated partial write"))
1076            }
1077
1078            fn flush_durable(&mut self) -> std::io::Result<()> {
1079                Ok(())
1080            }
1081
1082            fn sync_data_durable(&mut self) -> std::io::Result<()> {
1083                Ok(())
1084            }
1085
1086            fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
1087                self.bytes.truncate(usize::try_from(len).unwrap());
1088                Ok(())
1089            }
1090        }
1091
1092        let original = b"{\"event_type\":\"existing\"}\n".to_vec();
1093        let mut writer = PartialFailWriter {
1094            bytes: original.clone(),
1095        };
1096        let error = append_event_line_durably(
1097            &mut writer,
1098            b"{\"event_type\":\"new\",\"payload\":{}}\n",
1099            original.len() as u64,
1100        )
1101        .unwrap_err();
1102
1103        assert!(matches!(error, SessionAppendError::RolledBack(_)));
1104        assert!(
1105            error
1106                .to_string()
1107                .contains("failed to write session JSONL event")
1108        );
1109        assert_eq!(writer.bytes, original);
1110    }
1111
1112    #[test]
1113    fn session_append_reports_uncertain_state_when_rollback_fails() {
1114        struct UnrecoverableWriter;
1115        impl DurableSessionWriter for UnrecoverableWriter {
1116            fn write_all_durable(&mut self, _: &[u8]) -> std::io::Result<()> {
1117                Err(std::io::Error::other("write"))
1118            }
1119            fn flush_durable(&mut self) -> std::io::Result<()> {
1120                Ok(())
1121            }
1122            fn sync_data_durable(&mut self) -> std::io::Result<()> {
1123                Ok(())
1124            }
1125            fn set_len_durable(&mut self, _: u64) -> std::io::Result<()> {
1126                Err(std::io::Error::other("rollback"))
1127            }
1128        }
1129        let mut writer = UnrecoverableWriter;
1130        let error = append_event_line_durably(&mut writer, b"line\n", 0).unwrap_err();
1131        assert!(matches!(error, SessionAppendError::Uncertain(_)));
1132        assert!(error.to_string().contains("on-disk state uncertain"));
1133    }
1134    #[test]
1135    fn session_append_flushes_and_syncs_event_line() {
1136        #[derive(Default)]
1137        struct RecordingWriter {
1138            bytes: Vec<u8>,
1139            flushed: bool,
1140            synced: bool,
1141        }
1142
1143        impl DurableSessionWriter for RecordingWriter {
1144            fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
1145                self.bytes.extend_from_slice(bytes);
1146                Ok(())
1147            }
1148
1149            fn flush_durable(&mut self) -> std::io::Result<()> {
1150                self.flushed = true;
1151                Ok(())
1152            }
1153
1154            fn sync_data_durable(&mut self) -> std::io::Result<()> {
1155                self.synced = true;
1156                Ok(())
1157            }
1158
1159            fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
1160                self.bytes.truncate(usize::try_from(len).unwrap());
1161                Ok(())
1162            }
1163        }
1164
1165        let temp = TempDir::new().unwrap();
1166        let event = SessionEvent::new(
1167            "user_input",
1168            "safe".to_string(),
1169            temp.path().to_path_buf(),
1170            json!({"text":"hello"}),
1171        );
1172        let mut writer = RecordingWriter::default();
1173
1174        let line = serialize_session_event_line(&event).unwrap();
1175        append_event_line_durably(&mut writer, &line, 0).unwrap();
1176
1177        assert!(writer.flushed);
1178        assert!(writer.synced);
1179        assert!(writer.bytes.ends_with(b"\n"));
1180        let parsed: SessionEvent = serde_json::from_slice(&writer.bytes[..writer.bytes.len() - 1])
1181            .expect("durable append writes one valid JSONL event");
1182        assert_eq!(parsed.event_type, "user_input");
1183    }
1184
1185    #[test]
1186    fn try_record_session_event_returns_append_failure() {
1187        let temp = TempDir::new().unwrap();
1188        let invalid_session =
1189            Session::unchecked_for_test("../unsafe".to_string(), temp.path().join("unsafe.jsonl"));
1190
1191        let error = try_record_session_event(
1192            Some(&invalid_session),
1193            temp.path(),
1194            "assistant_chunk",
1195            json!({"text":"best effort"}),
1196        )
1197        .unwrap_err();
1198
1199        assert!(error.to_string().contains("must not contain '..'"));
1200        assert!(matches!(
1201            error.downcast_ref::<SessionAppendError>(),
1202            Some(SessionAppendError::NotWritten(_))
1203        ));
1204        assert!(!invalid_session.path().exists());
1205    }
1206}