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