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