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 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 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}