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