1use super::event::{SessionEvent, SessionEventKind};
2use super::manager::{Session, SessionInternalDiagnostic, SessionListReport, SessionManager};
3use super::read::{MAX_METADATA_VISIT_BYTES, MAX_METADATA_VISIT_LINES, validate_session_id};
4use super::store::{open_existing_named, validate_existing_file, validate_path_file};
5use super::write::{
6 COMPACTION_SCHEMA_VERSION, sanitize_compaction_summary, sanitize_session_title,
7};
8use crate::{output::redact_sensitive_text, persistence::atomic_write_with_permissions};
9use chrono::{DateTime, Utc};
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12use std::{
13 fs,
14 io::Read,
15 path::{Path, PathBuf},
16 time::{SystemTime, UNIX_EPOCH},
17};
18
19const MAX_SESSION_METADATA_BYTES: u64 = 64 * 1024;
20const MAX_SESSION_INTERNAL_DIAGNOSTICS: usize = 64;
21pub(crate) fn push_session_diagnostic(
22 diagnostics: &mut Vec<SessionInternalDiagnostic>,
23 session_id: Option<String>,
24 message: String,
25) {
26 if diagnostics.len() < MAX_SESSION_INTERNAL_DIAGNOSTICS {
27 let redacted = redact_sensitive_text(&message);
28 let mut chars = redacted.chars();
29 let mut bounded = chars.by_ref().take(500).collect::<String>();
30 if chars.next().is_some() {
31 bounded.push('…');
32 }
33 diagnostics.push(SessionInternalDiagnostic {
34 session_id,
35 message: bounded,
36 });
37 }
38}
39
40#[derive(Debug, Clone, Copy)]
41pub(crate) enum SessionDiagnosticOperation {
42 Replay,
43 Metadata,
44 Title,
45}
46
47impl SessionDiagnosticOperation {
48 fn as_str(self) -> &'static str {
49 match self {
50 Self::Replay => "replay",
51 Self::Metadata => "metadata",
52 Self::Title => "title",
53 }
54 }
55}
56
57pub(crate) fn report_session_diagnostic(
58 operation: SessionDiagnosticOperation,
59 _path: &Path,
60 error: impl std::fmt::Display,
61) {
62 let error = redact_sensitive_text(&error.to_string())
63 .chars()
64 .take(500)
65 .collect::<String>();
66 crate::output::emit_terminal_warning(format!(
67 "warning: session operation={} category=session_jsonl failed: {error}",
68 operation.as_str()
69 ));
70}
71const SESSION_METADATA_SCHEMA_VERSION: u64 = 3;
72const SESSION_METADATA_EXTENSION: &str = "metadata.json";
73
74#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
75pub(in crate::sessions) struct SessionMetadataRecord {
76 schema_version: u64,
77 session_id: String,
78 latest_activity_timestamp: Option<DateTime<Utc>>,
79 latest_title: Option<String>,
80 session_path: Option<PathBuf>,
81 #[serde(default)]
82 first_user_input_text: Option<String>,
83 user_input_count: usize,
84 jsonl_len: u64,
85 jsonl_modified_ns: Option<u128>,
86 complete: bool,
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
91pub(crate) struct SessionMetadataSummary {
92 pub(crate) session: Session,
93 pub(crate) latest_activity_time: SystemTime,
94 pub(crate) latest_title: Option<String>,
95 pub(crate) session_path: Option<PathBuf>,
96}
97
98impl SessionMetadataRecord {
99 fn empty(session: &Session, marker: JsonlMarker, complete: bool) -> Self {
100 Self {
101 schema_version: SESSION_METADATA_SCHEMA_VERSION,
102 session_id: session.id.clone(),
103 latest_activity_timestamp: None,
104 latest_title: None,
105 session_path: None,
106 first_user_input_text: None,
107 user_input_count: 0,
108 jsonl_len: marker.len,
109 jsonl_modified_ns: marker.modified_ns,
110 complete,
111 }
112 }
113
114 fn apply_event(&mut self, event: &SessionEvent) {
115 if self.session_path.is_none() {
116 self.session_path = event.session_path.clone();
117 }
118 self.latest_activity_timestamp = Some(
119 self.latest_activity_timestamp
120 .map_or(event.timestamp, |current| current.max(event.timestamp)),
121 );
122 match event.kind() {
123 Some(SessionEventKind::SessionTitle) => {
124 if let Some(title) = event.payload.get("title").and_then(Value::as_str)
125 && let Some(title) = sanitize_session_title(title)
126 {
127 self.latest_title = Some(title);
128 }
129 self.first_user_input_text = None;
130 }
131 Some(SessionEventKind::UserInput) => {
132 self.user_input_count = self.user_input_count.saturating_add(1);
133 if self.first_user_input_text.is_none()
134 && let Some(text) = event.payload.get("text").and_then(Value::as_str)
135 && !text.trim().is_empty()
136 {
137 self.first_user_input_text = Some(text.to_string());
138 }
139 }
140 Some(SessionEventKind::Compaction) => {
141 if let Some(aggregate) = event.payload.get("aggregate") {
142 if let Some(title) = aggregate.get("latest_title").and_then(Value::as_str) {
143 self.latest_title = sanitize_session_title(title);
144 }
145 if let Some(count) = aggregate.get("user_input_count").and_then(Value::as_u64)
146 && let Ok(count) = usize::try_from(count)
147 {
148 self.user_input_count = count;
149 }
150 self.first_user_input_text = aggregate
151 .get("first_user_input_text")
152 .and_then(Value::as_str)
153 .map(|text| redact_sensitive_text(text).chars().take(4_000).collect());
154 if let Some(path) = aggregate.get("session_path").and_then(Value::as_str) {
155 self.session_path = Some(PathBuf::from(path));
156 }
157 if let Some(timestamp) = aggregate.get("latest_activity_timestamp")
158 && let Ok(timestamp) =
159 serde_json::from_value::<chrono::DateTime<Utc>>(timestamp.clone())
160 {
161 self.latest_activity_timestamp = Some(
162 self.latest_activity_timestamp
163 .map_or(timestamp, |current| current.max(timestamp)),
164 );
165 }
166 }
167 }
168 _ => {}
169 }
170 }
171
172 pub(crate) fn checkpoint_aggregate(&self) -> Value {
173 serde_json::json!({
174 "latest_title": self.latest_title,
175 "user_input_count": self.user_input_count,
176 "first_user_input_text": self.first_user_input_text.as_deref().map(|text| redact_sensitive_text(text).chars().take(4_000).collect::<String>()),
177 "session_path": self.session_path.as_ref().map(|path| path.to_string_lossy().to_string()),
178 "latest_activity_timestamp": self.latest_activity_timestamp,
179 })
180 }
181
182 pub(crate) fn activity_time(&self, session: &Session) -> SystemTime {
183 self.latest_activity_timestamp
184 .map(SystemTime::from)
185 .unwrap_or_else(|| jsonl_modified_time(session).unwrap_or(SystemTime::UNIX_EPOCH))
186 }
187}
188
189#[derive(Debug, Clone, Copy, PartialEq, Eq)]
190struct JsonlMarker {
191 len: u64,
192 modified_ns: Option<u128>,
193}
194
195impl JsonlMarker {
196 fn from_metadata(metadata: &fs::Metadata) -> Self {
197 Self {
198 len: metadata.len(),
199 modified_ns: metadata.modified().ok().map(system_time_ns),
200 }
201 }
202}
203impl SessionManager {
204 pub(crate) fn list_metadata_summaries(&self) -> anyhow::Result<Vec<SessionMetadataSummary>> {
205 Ok(self.list_metadata_report()?.summaries)
206 }
207
208 pub(crate) fn list_metadata_report(&self) -> anyhow::Result<SessionListReport> {
209 if !self.root.exists() {
210 return Ok(SessionListReport {
211 summaries: Vec::new(),
212 diagnostics: Vec::new(),
213 });
214 }
215 super::store::validate_session_root(&self.root)?;
216 let mut summaries = Vec::new();
217 let mut diagnostics = Vec::new();
218 for entry in fs::read_dir(&self.root)? {
219 let entry = match entry {
220 Ok(entry) => entry,
221 Err(error) => {
222 push_session_diagnostic(
223 &mut diagnostics,
224 None,
225 format!("failed to read session directory entry: {error}"),
226 );
227 continue;
228 }
229 };
230 let path = entry.path();
231 let Some((id, path)) = session_from_entry(path) else {
232 continue;
233 };
234 let id = match validate_session_id(id) {
235 Ok(id) => id,
236 Err(error) => {
237 push_session_diagnostic(
238 &mut diagnostics,
239 None,
240 format!("ignored invalid session file {}: {error}", path.display()),
241 );
242 continue;
243 }
244 };
245 if validate_path_file(&self.root, &id, &path).is_err() {
246 push_session_diagnostic(
247 &mut diagnostics,
248 Some(id),
249 "ignored unsafe session file".to_string(),
250 );
251 continue;
252 }
253 let session = Session::new(id, path);
254 match session_metadata_summary(session) {
255 Ok(summary) => summaries.push(summary),
256 Err(error) => push_session_diagnostic(
257 &mut diagnostics,
258 None,
259 format!("failed to load session metadata: {error}"),
260 ),
261 }
262 }
263 summaries.sort_by(|left, right| {
264 left.latest_activity_time
265 .cmp(&right.latest_activity_time)
266 .then_with(|| left.session.id.cmp(&right.session.id))
267 });
268 Ok(SessionListReport {
269 summaries,
270 diagnostics,
271 })
272 }
273}
274
275pub(crate) fn session_from_entry(path: PathBuf) -> Option<(String, PathBuf)> {
276 if path
277 .extension()
278 .is_none_or(|extension| extension != "jsonl")
279 {
280 return None;
281 }
282 let id = path.file_stem()?.to_string_lossy().to_string();
283 Some((id, path))
284}
285
286pub(crate) fn metadata_path_for_session(session: &Session) -> PathBuf {
287 session.path.with_extension(SESSION_METADATA_EXTENSION)
288}
289
290fn jsonl_marker(session: &Session) -> anyhow::Result<JsonlMarker> {
291 let root = session
292 .path
293 .parent()
294 .ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
295 let Some(file) = super::store::open_existing_primary(root, &session.id)? else {
296 anyhow::bail!("session JSONL is missing");
297 };
298 Ok(JsonlMarker::from_metadata(&file.metadata()?))
299}
300
301fn jsonl_modified_time(session: &Session) -> Option<SystemTime> {
302 let root = session.path.parent()?;
303 super::store::open_existing_primary(root, &session.id)
304 .ok()
305 .flatten()
306 .and_then(|file| file.metadata().ok())
307 .and_then(|metadata| metadata.modified().ok())
308}
309
310fn system_time_ns(time: SystemTime) -> u128 {
311 time.duration_since(UNIX_EPOCH)
312 .unwrap_or_default()
313 .as_nanos()
314}
315
316fn read_session_metadata_with_policy(
317 session: &Session,
318 require_complete: bool,
319) -> anyhow::Result<Option<SessionMetadataRecord>> {
320 let path = metadata_path_for_session(session);
321 let root = session
322 .path
323 .parent()
324 .ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
325 let name = path
326 .file_name()
327 .ok_or_else(|| anyhow::anyhow!("metadata path has no filename"))?
328 .to_string_lossy();
329 if !session.path.exists() {
330 return Ok(None);
331 }
332 let Some(file) = open_existing_named(root, &name)? else {
333 return Ok(None);
334 };
335 let mut bytes = Vec::new();
336 file.take(MAX_SESSION_METADATA_BYTES.saturating_add(1))
337 .read_to_end(&mut bytes)?;
338 if bytes.len() as u64 > MAX_SESSION_METADATA_BYTES {
339 return Ok(None);
340 }
341 let Ok(record) = serde_json::from_slice::<SessionMetadataRecord>(&bytes) else {
342 return Ok(None);
343 };
344 if record.schema_version != SESSION_METADATA_SCHEMA_VERSION
345 || record.session_id != session.id
346 || (require_complete && !record.complete)
347 {
348 return Ok(None);
349 }
350 let marker = jsonl_marker(session)?;
351 if record.jsonl_len != marker.len || record.jsonl_modified_ns != marker.modified_ns {
352 return Ok(None);
353 }
354 Ok(Some(record))
355}
356
357pub(crate) fn read_complete_session_metadata(
358 session: &Session,
359) -> anyhow::Result<Option<SessionMetadataRecord>> {
360 read_session_metadata(session)
361}
362
363fn read_session_metadata(session: &Session) -> anyhow::Result<Option<SessionMetadataRecord>> {
364 read_session_metadata_with_policy(session, true)
365}
366
367pub(crate) fn metadata_for_append(
368 session: &Session,
369 current_metadata: &fs::Metadata,
370) -> anyhow::Result<SessionMetadataRecord> {
371 let marker = JsonlMarker::from_metadata(current_metadata);
372 if let Some(record) = session
373 .metadata_cache()
374 .lock()
375 .map_err(|_| anyhow::anyhow!("session metadata cache was poisoned"))?
376 .clone()
377 && record.complete
378 && record.jsonl_len == marker.len
379 && record.jsonl_modified_ns == marker.modified_ns
380 {
381 return Ok(record);
382 }
383 let record = rebuild_session_metadata_record_from_jsonl(session)?;
384 cache_session_metadata(session, Some(record.clone()))?;
385 Ok(record)
386}
387fn write_session_metadata(session: &Session, record: &SessionMetadataRecord) -> anyhow::Result<()> {
388 let path = metadata_path_for_session(session);
389 validate_existing_file(&path)?;
390 let bytes = serde_json::to_vec_pretty(record)?;
391 atomic_write_with_permissions(&path, &bytes, Some(super::store::SESSION_FILE_MODE))
392}
393
394fn rebuild_session_metadata_record_from_jsonl(
395 session: &Session,
396) -> anyhow::Result<SessionMetadataRecord> {
397 validate_session_id(session.id.clone())?;
398 let initial_marker = jsonl_marker(session)?;
399 let mut record = SessionMetadataRecord::empty(session, initial_marker, true);
400 let (_, _) = session.visit_events_tolerant_bounded(
401 MAX_METADATA_VISIT_LINES,
402 MAX_METADATA_VISIT_BYTES,
403 |event| {
404 if event.session_id == session.id {
405 record.apply_event(&event);
406 }
407 },
408 )?;
409 let final_marker = jsonl_marker(session)?;
410 record.jsonl_len = final_marker.len;
411 record.jsonl_modified_ns = final_marker.modified_ns;
412 record.complete = initial_marker == final_marker;
413 Ok(record)
414}
415
416fn cache_session_metadata(
417 session: &Session,
418 record: Option<SessionMetadataRecord>,
419) -> anyhow::Result<()> {
420 *session
421 .metadata_cache()
422 .lock()
423 .map_err(|_| anyhow::anyhow!("session metadata cache was poisoned"))? = record;
424 Ok(())
425}
426
427pub(crate) fn invalidate_session_metadata_cache(session: &Session) -> anyhow::Result<()> {
428 cache_session_metadata(session, None)
429}
430
431pub(crate) fn rebuild_session_metadata_from_jsonl(
432 session: &Session,
433) -> anyhow::Result<SessionMetadataRecord> {
434 let record = rebuild_session_metadata_record_from_jsonl(session)?;
435 cache_session_metadata(session, Some(record.clone()))?;
436 if let Err(error) = write_session_metadata(session, &record)
437 && !is_session_file_not_found(&error)
438 {
439 report_session_diagnostic(
440 SessionDiagnosticOperation::Metadata,
441 &metadata_path_for_session(session),
442 &error,
443 );
444 }
445 Ok(record)
446}
447
448fn is_session_file_not_found(error: &anyhow::Error) -> bool {
451 error.chain().any(|cause| {
452 cause
453 .downcast_ref::<std::io::Error>()
454 .is_some_and(|io_error| io_error.kind() == std::io::ErrorKind::NotFound)
455 })
456}
457
458pub(crate) fn write_rotated_session_metadata(
459 session: &Session,
460 mut record: SessionMetadataRecord,
461 checkpoint: &SessionEvent,
462) -> anyhow::Result<()> {
463 record.apply_event(checkpoint);
464 let final_marker = jsonl_marker(session)?;
465 record.jsonl_len = final_marker.len;
466 record.jsonl_modified_ns = final_marker.modified_ns;
467 record.complete = true;
468 cache_session_metadata(session, Some(record.clone()))?;
469 write_session_metadata(session, &record)
470}
471
472pub(crate) fn session_metadata_for_listing(
473 session: &Session,
474) -> anyhow::Result<SessionMetadataRecord> {
475 if let Some(record) = read_session_metadata(session)? {
476 cache_session_metadata(session, Some(record.clone()))?;
477 return Ok(record);
478 }
479 rebuild_session_metadata_from_jsonl(session)
480}
481
482fn session_metadata_summary(session: Session) -> anyhow::Result<SessionMetadataSummary> {
483 let record = session_metadata_for_listing(&session)?;
484 let latest_activity_time = record.activity_time(&session);
485 Ok(SessionMetadataSummary {
486 session,
487 latest_activity_time,
488 latest_title: record.latest_title,
489 session_path: record.session_path,
490 })
491}
492
493#[derive(Debug, Clone, PartialEq, Eq)]
494pub(crate) struct CompactionCheckpoint {
495 pub(crate) summary: String,
496 pub(crate) provider: String,
497 pub(crate) model: String,
498 pub(crate) cutoff_event_count: usize,
499 pub(crate) event_index: usize,
500}
501pub(crate) fn update_session_metadata_after_append_batch(
502 session: &Session,
503 previous: Option<SessionMetadataRecord>,
504 events: &[SessionEvent],
505 sidecar_exists: bool,
506) -> anyhow::Result<()> {
507 let mut record = previous.unwrap_or_else(|| {
508 SessionMetadataRecord::empty(
509 session,
510 JsonlMarker {
511 len: 0,
512 modified_ns: None,
513 },
514 false,
515 )
516 });
517 let was_complete = record.complete;
518 for event in events {
519 record.apply_event(event);
520 }
521 let marker = jsonl_marker(session)?;
522 record.jsonl_len = marker.len;
523 record.jsonl_modified_ns = marker.modified_ns;
524 record.complete = true;
525 let requires_write = !sidecar_exists
526 || events.iter().any(|event| {
527 matches!(
528 event.kind(),
529 Some(
530 SessionEventKind::SessionTitle
531 | SessionEventKind::UserInput
532 | SessionEventKind::AssistantOutput
533 | SessionEventKind::TurnStatus
534 )
535 )
536 })
537 || (!was_complete && record.complete);
538 cache_session_metadata(session, Some(record.clone()))?;
540 if requires_write {
541 write_session_metadata(session, &record)?;
542 }
543 Ok(())
544}
545
546pub(crate) fn latest_valid_compaction_checkpoint_for_replay(
547 session_id: &str,
548 events: &[SessionEvent],
549) -> (Option<CompactionCheckpoint>, Vec<String>) {
550 latest_valid_compaction_checkpoint_with_policy(session_id, events, true)
551}
552
553fn latest_valid_compaction_checkpoint_with_policy(
554 session_id: &str,
555 events: &[SessionEvent],
556 skip_foreign_session_events: bool,
557) -> (Option<CompactionCheckpoint>, Vec<String>) {
558 let mut diagnostics = Vec::new();
559 for (event_index, event) in events.iter().enumerate().rev() {
560 if event.kind() != Some(SessionEventKind::Compaction) {
561 continue;
562 }
563 if skip_foreign_session_events && event.session_id != session_id {
564 continue;
565 }
566 match validate_compaction_checkpoint(session_id, event, event_index, events.len()) {
567 Ok(checkpoint) => return (Some(checkpoint), diagnostics),
568 Err(message) => diagnostics.push(message),
569 }
570 }
571 (None, diagnostics)
572}
573
574pub(super) fn validate_compaction_checkpoint(
575 session_id: &str,
576 event: &SessionEvent,
577 event_index: usize,
578 total_event_count: usize,
579) -> Result<CompactionCheckpoint, String> {
580 if event.session_id != session_id {
581 return Err(format!(
582 "ignored malformed compaction checkpoint at event_index={event_index}: session_id_mismatch"
583 ));
584 }
585 let schema_version = event
586 .payload
587 .get("schema_version")
588 .and_then(Value::as_u64)
589 .ok_or_else(|| {
590 format!(
591 "ignored malformed compaction checkpoint at event_index={event_index}: missing_schema_version"
592 )
593 })?;
594 if schema_version != COMPACTION_SCHEMA_VERSION {
595 return Err(format!(
596 "ignored malformed compaction checkpoint at event_index={event_index}: unsupported_schema_version"
597 ));
598 }
599 let summary = event
600 .payload
601 .get("summary")
602 .and_then(Value::as_str)
603 .and_then(sanitize_compaction_summary)
604 .ok_or_else(|| {
605 format!(
606 "ignored malformed compaction checkpoint at event_index={event_index}: empty_summary"
607 )
608 })?;
609 let provider = required_non_blank_payload_string(event, "provider", event_index)?;
610 let model = required_non_blank_payload_string(event, "model", event_index)?;
611 let cutoff_event_count = event
612 .payload
613 .get("cutoff_event_count")
614 .and_then(Value::as_u64)
615 .and_then(|value| usize::try_from(value).ok())
616 .ok_or_else(|| {
617 format!(
618 "ignored malformed compaction checkpoint at event_index={event_index}: missing_cutoff_event_count"
619 )
620 })?;
621 if cutoff_event_count > event_index || cutoff_event_count > total_event_count {
622 return Err(format!(
623 "ignored malformed compaction checkpoint at event_index={event_index}: cutoff_event_count_out_of_bounds"
624 ));
625 }
626 Ok(CompactionCheckpoint {
627 summary,
628 provider,
629 model,
630 cutoff_event_count,
631 event_index,
632 })
633}
634
635fn required_non_blank_payload_string(
636 event: &SessionEvent,
637 key: &str,
638 event_index: usize,
639) -> Result<String, String> {
640 event.payload.get(key).and_then(Value::as_str).map(str::trim)
641 .filter(|value| !value.is_empty())
642 .map(ToString::to_string)
643 .ok_or_else(|| format!("ignored malformed compaction checkpoint at event_index={event_index}: missing_{key}"))
644}
645
646#[derive(Debug, Clone, Default, PartialEq, Eq)]
647pub(crate) struct SessionTitleMetadata {
648 pub(crate) latest_title: Option<String>,
649 pub(crate) user_input_count: usize,
650 pub(crate) first_user_input_text: Option<String>,
651}
652
653pub(crate) fn session_title_metadata(session: &Session) -> anyhow::Result<SessionTitleMetadata> {
654 if let Some(record) = read_session_metadata(session)? {
655 return Ok(SessionTitleMetadata {
656 latest_title: record.latest_title,
657 user_input_count: record.user_input_count,
658 first_user_input_text: record.first_user_input_text,
659 });
660 }
661 let mut metadata = SessionTitleMetadata::default();
662 session
663 .visit_events_tolerant_bounded(
664 MAX_METADATA_VISIT_LINES,
665 MAX_METADATA_VISIT_BYTES,
666 |event| match event.kind() {
667 Some(SessionEventKind::SessionTitle) => {
668 if let Some(title) = event
669 .payload
670 .get("title")
671 .and_then(Value::as_str)
672 .and_then(sanitize_session_title)
673 {
674 metadata.latest_title = Some(title);
675 }
676 }
677 Some(SessionEventKind::UserInput) => {
678 metadata.user_input_count = metadata.user_input_count.saturating_add(1);
679 if metadata.first_user_input_text.is_none()
680 && let Some(text) = event.payload.get("text").and_then(Value::as_str)
681 && !text.trim().is_empty()
682 {
683 metadata.first_user_input_text = Some(text.to_string());
684 }
685 }
686 _ => {}
687 },
688 )
689 .inspect_err(|error| {
690 report_session_diagnostic(SessionDiagnosticOperation::Title, session.path(), error);
691 })?;
692 Ok(metadata)
693}
694
695pub fn latest_session_title(session: &Session) -> anyhow::Result<Option<String>> {
696 Ok(session_title_metadata(session)?.latest_title)
697}
698
699pub fn session_user_input_count(session: &Session) -> anyhow::Result<usize> {
700 Ok(session_title_metadata(session)?.user_input_count)
701}