1use regex::Regex;
2use serde::Serialize;
3use serde_json::Value;
4use std::collections::{BTreeMap, BTreeSet};
5use std::ffi::OsStr;
6use std::fmt;
7use std::fs::{self, File, Metadata, OpenOptions};
8use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
9use std::path::{Component, Path, PathBuf};
10use std::sync::OnceLock;
11use std::time::{Duration, SystemTime, UNIX_EPOCH};
12
13pub const MAX_HEAD_BYTES: usize = 1024 * 1024;
14pub const DEFAULT_MAX_TAIL_BYTES: usize = 16 * 1024 * 1024;
15pub const DEFAULT_MAX_LINES: usize = 50_000;
16pub const DEFAULT_MAX_MESSAGES: usize = 48;
17pub const MAX_MESSAGE_CHARS: usize = 4096;
18pub const MAX_OUTPUT_BYTES: usize = 256 * 1024;
19const MAX_SESSION_FILES: usize = 50_000;
20const MAX_LINE_BYTES: usize = 1024 * 1024;
21const MAX_INDEX_BYTES: u64 = 8 * 1024 * 1024;
22const MAX_TOOL_NAMES: usize = 64;
23const MAX_FILE_HINTS: usize = 64;
24const MAX_RECENT_TOOL_OPS: usize = 15;
28const TOOL_OP_DETAIL_CHARS: usize = 200;
30
31#[derive(Clone, Copy, Debug)]
32pub struct RestoreLimits {
33 pub max_tail_bytes: usize,
34 pub max_lines: usize,
35 pub max_messages: usize,
36}
37
38impl Default for RestoreLimits {
39 fn default() -> Self {
40 Self {
41 max_tail_bytes: DEFAULT_MAX_TAIL_BYTES,
42 max_lines: DEFAULT_MAX_LINES,
43 max_messages: DEFAULT_MAX_MESSAGES,
44 }
45 }
46}
47
48impl RestoreLimits {
49 pub fn validate(self) -> Result<Self, RestoreError> {
50 if !(1024..=64 * 1024 * 1024).contains(&self.max_tail_bytes)
51 || !(1..=100_000).contains(&self.max_lines)
52 || !(1..=256).contains(&self.max_messages)
53 {
54 return Err(RestoreError::InvalidArgument(
55 "restore limits are outside the supported bounds".to_owned(),
56 ));
57 }
58 Ok(self)
59 }
60}
61
62#[derive(Debug)]
63pub enum RestoreError {
64 Io(std::io::Error),
65 Json(serde_json::Error),
66 HomeUnavailable,
67 InvalidArgument(String),
68 InvalidTarget,
69 UnsafeCandidate,
70 AmbiguousPrefix,
71 NotFound,
72 NoSessionMeta,
73 OutputLimit,
74}
75
76impl fmt::Display for RestoreError {
77 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
78 match self {
79 Self::Io(_) => f.write_str("session storage is unavailable"),
80 Self::Json(_) => f.write_str("session metadata is malformed"),
81 Self::HomeUnavailable => f.write_str("Codex home is unavailable"),
82 Self::InvalidArgument(message) => f.write_str(message),
83 Self::InvalidTarget => f.write_str("session selector is invalid"),
84 Self::UnsafeCandidate => f.write_str("session candidate is outside the trusted root"),
85 Self::AmbiguousPrefix => f.write_str("session ID prefix is ambiguous"),
86 Self::NotFound => f.write_str("session was not found"),
87 Self::NoSessionMeta => f.write_str("session metadata is missing"),
88 Self::OutputLimit => f.write_str("bounded report exceeds the output limit"),
89 }
90 }
91}
92
93impl std::error::Error for RestoreError {}
94
95impl From<std::io::Error> for RestoreError {
96 fn from(value: std::io::Error) -> Self {
97 Self::Io(value)
98 }
99}
100
101impl From<serde_json::Error> for RestoreError {
102 fn from(value: serde_json::Error) -> Self {
103 Self::Json(value)
104 }
105}
106
107#[derive(Clone, Debug, Serialize)]
108pub struct SessionCandidate {
109 #[serde(skip)]
110 pub path: PathBuf,
111 pub id: String,
112 pub title: Option<String>,
113 pub updated_unix_ms: u64,
114 pub size_bytes: u64,
115}
116
117#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
118#[serde(rename_all = "lowercase")]
119pub enum MessageRole {
120 User,
121 Assistant,
122 Tool,
124 Error,
126}
127
128#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
129pub struct Message {
130 pub role: MessageRole,
131 pub text: String,
132 pub timestamp: Option<String>,
133}
134
135#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
140pub struct ToolOperation {
141 pub name: String,
142 pub detail: String,
143 pub timestamp: Option<String>,
144}
145
146#[derive(Clone, Debug, Default, Serialize)]
147pub struct ToolInventory {
148 pub counts: BTreeMap<String, u64>,
149 pub changed_files: BTreeSet<String>,
150}
151
152#[derive(Clone, Debug, Default, Serialize)]
153pub struct GitHints {
154 pub recorded_branch: Option<String>,
155 pub recorded_commit: Option<String>,
156 pub repository_label: Option<String>,
157}
158
159#[derive(Clone, Debug, Serialize)]
160pub struct SessionMeta {
161 pub id: String,
162 pub title: Option<String>,
163 pub started_at: Option<String>,
164 pub updated_unix_ms: u64,
165 pub workspace_label: Option<String>,
166 pub model: Option<String>,
167 pub model_provider: Option<String>,
168}
169
170#[derive(Clone, Debug, Serialize)]
171pub struct SessionReport {
172 pub schema: &'static str,
173 pub meta: SessionMeta,
174 pub messages: Vec<Message>,
175 pub tool_ops: Vec<ToolOperation>,
176 pub tools: ToolInventory,
177 pub git: GitHints,
178 pub truncated: bool,
179 pub malformed_records: u64,
180 pub redactions: u64,
181}
182
183#[derive(Clone, Debug)]
184pub struct SessionSource {
185 pub path: PathBuf,
186 pub metadata: Metadata,
187 pub id: String,
188 pub title: Option<String>,
189}
190
191#[derive(Clone, Debug)]
192struct ParsedMeta {
193 id: String,
194 timestamp: Option<String>,
195 cwd: Option<PathBuf>,
196 model_provider: Option<String>,
197 git: GitHints,
198}
199
200#[derive(Clone, Debug)]
201struct TimedMessage {
202 ordinal: usize,
203 timestamp: Option<String>,
204 text: String,
205}
206
207pub fn default_codex_home() -> Result<PathBuf, RestoreError> {
208 if let Some(home) = std::env::var_os("CODEX_HOME") {
209 if !home.is_empty() {
210 return Ok(PathBuf::from(home));
211 }
212 }
213 std::env::var_os("USERPROFILE")
214 .filter(|value| !value.is_empty())
215 .map(PathBuf::from)
216 .map(|home| home.join(".codex"))
217 .ok_or(RestoreError::HomeUnavailable)
218}
219
220pub fn list_sessions(
221 home: &Path,
222 max_age_hours: Option<u64>,
223 limit: usize,
224) -> Result<Vec<SessionCandidate>, RestoreError> {
225 if !(1..=100).contains(&limit) {
226 return Err(RestoreError::InvalidArgument(
227 "list limit must be between 1 and 100".to_owned(),
228 ));
229 }
230 let root = trusted_sessions_root(home)?;
231 let titles = read_session_index(home)?;
232 let now = SystemTime::now();
233 let cutoff = max_age_hours
234 .map(|hours| Duration::from_secs(hours.saturating_mul(3600)))
235 .and_then(|age| now.checked_sub(age));
236 let mut candidates = Vec::new();
237 for path in discover_session_paths(&root)? {
238 let metadata = safe_candidate_metadata(&root, &path)?;
239 let modified = metadata.modified().unwrap_or(UNIX_EPOCH);
240 if cutoff.is_some_and(|value| modified < value) {
241 continue;
242 }
243 let Some(meta) = read_session_meta(&path, &metadata)? else {
244 continue;
245 };
246 let Some(file_id) = rollout_filename_id(&path) else {
247 continue;
248 };
249 if meta.id != file_id {
250 continue;
251 }
252 let mut title = titles.get(&meta.id).cloned();
257 if title.is_none() {
258 let mut ignored_redactions = 0;
259 title = first_human_prompt(&path)
260 .ok()
261 .flatten()
262 .and_then(|prompt| safe_scalar(&prompt, 256, &mut ignored_redactions));
263 }
264 candidates.push(SessionCandidate {
265 path,
266 id: meta.id.clone(),
267 title,
268 updated_unix_ms: system_time_millis(modified),
269 size_bytes: metadata.len(),
270 });
271 }
272 candidates.sort_by(|left, right| {
273 right
274 .updated_unix_ms
275 .cmp(&left.updated_unix_ms)
276 .then_with(|| left.id.cmp(&right.id))
277 });
278 candidates.truncate(limit);
279 Ok(candidates)
280}
281
282pub fn resolve_target(home: &Path, target: &OsStr) -> Result<SessionSource, RestoreError> {
283 let root = trusted_sessions_root(home)?;
284 let target_path = PathBuf::from(target);
285 let path = if target_path.components().count() > 1 || target_path.is_absolute() {
286 if !target_path.is_absolute() {
287 return Err(RestoreError::InvalidTarget);
288 }
289 safe_candidate_metadata(&root, &target_path)?;
290 let canonical = fs::canonicalize(&target_path).map_err(|_| RestoreError::NotFound)?;
291 ensure_descendant(&root, &canonical)?;
292 canonical
293 } else {
294 let selector = target.to_str().ok_or(RestoreError::InvalidTarget)?;
295 if selector.len() < 16 || !selector.chars().all(is_id_selector_char) {
296 return Err(RestoreError::InvalidTarget);
297 }
298 let mut matches = Vec::new();
299 for candidate_path in discover_session_paths(&root)? {
300 let Some(candidate_id) = rollout_filename_id(&candidate_path) else {
301 continue;
302 };
303 if candidate_id == selector || candidate_id.starts_with(selector) {
304 let metadata = safe_candidate_metadata(&root, &candidate_path)?;
305 let Some(meta) = read_session_meta(&candidate_path, &metadata)? else {
306 continue;
307 };
308 if meta.id == candidate_id {
309 matches.push(candidate_path);
310 }
311 }
312 }
313 if matches.len() > 1 {
314 return Err(RestoreError::AmbiguousPrefix);
315 }
316 matches.pop().ok_or(RestoreError::NotFound)?
317 };
318 let metadata = safe_candidate_metadata(&root, &path)?;
319 let parsed = read_session_meta(&path, &metadata)?.ok_or(RestoreError::NoSessionMeta)?;
320 let file_id = rollout_filename_id(&path).ok_or(RestoreError::InvalidTarget)?;
321 if parsed.id != file_id {
322 return Err(RestoreError::UnsafeCandidate);
323 }
324 let titles = read_session_index(home)?;
325 Ok(SessionSource {
326 path,
327 metadata,
328 id: parsed.id.clone(),
329 title: titles.get(&parsed.id).cloned(),
330 })
331}
332
333pub fn load_session(
334 source: &SessionSource,
335 limits: RestoreLimits,
336) -> Result<SessionReport, RestoreError> {
337 let limits = limits.validate()?;
338 let (meta_line, records, mut truncated, malformed) = read_bounded_records(source, limits)?;
339 let parsed_meta = parse_session_meta(&meta_line)?.ok_or(RestoreError::NoSessionMeta)?;
340 if parsed_meta.id != source.id {
341 return Err(RestoreError::UnsafeCandidate);
342 }
343 let mut response_user = Vec::new();
344 let mut response_assistant = Vec::new();
345 let mut response_tool_output = Vec::new();
346 let mut event_user = Vec::new();
347 let mut event_assistant = Vec::new();
348 let mut event_error = Vec::new();
349 let mut tool_ops: Vec<ToolOperation> = Vec::new();
350 let mut tool_call_names: BTreeMap<String, String> = BTreeMap::new();
351 let mut tools = ToolInventory::default();
352 let mut cwd = parsed_meta.cwd.clone();
353 let mut model = None;
354 let mut redactions = 0_u64;
355 let mut malformed_records = malformed;
356
357 for (ordinal, line) in records.into_iter().enumerate() {
358 if line.len() > MAX_LINE_BYTES {
359 malformed_records += 1;
360 truncated = true;
361 continue;
362 }
363 let value: Value = match serde_json::from_str(&line) {
364 Ok(value) => value,
365 Err(_) => {
366 malformed_records += 1;
367 continue;
368 }
369 };
370 let record_type = value.get("type").and_then(Value::as_str).unwrap_or_default();
371 let payload = value.get("payload").unwrap_or(&Value::Null);
372 let timestamp = value
373 .get("timestamp")
374 .and_then(Value::as_str)
375 .map(bounded_scalar);
376 match record_type {
377 "turn_context" => {
378 if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
379 cwd = Some(PathBuf::from(value));
380 }
381 if let Some(value) = payload.get("model").and_then(Value::as_str) {
382 model = safe_scalar(value, 128, &mut redactions);
383 }
384 }
385 "event_msg" => match payload.get("type").and_then(Value::as_str) {
386 Some("user_message") => {
387 if let Some(text) = payload.get("message").and_then(Value::as_str) {
388 push_message(&mut event_user, ordinal, timestamp, text, &mut redactions);
389 }
390 }
391 Some("agent_message") => {
392 if let Some(text) = payload.get("message").and_then(Value::as_str) {
393 push_message(
394 &mut event_assistant,
395 ordinal,
396 timestamp,
397 text,
398 &mut redactions,
399 );
400 }
401 }
402 Some("item_completed") => {
408 if let Some(item) = payload.get("item") {
409 match item.get("type").and_then(Value::as_str) {
410 Some("user_message") => {
411 if let Some(text) = item_text(item) {
412 push_message(
413 &mut event_user,
414 ordinal,
415 timestamp,
416 &text,
417 &mut redactions,
418 );
419 }
420 }
421 Some("agent_message") => {
422 if let Some(text) = item_text(item) {
423 push_message(
424 &mut event_assistant,
425 ordinal,
426 timestamp,
427 &text,
428 &mut redactions,
429 );
430 }
431 }
432 _ => {}
433 }
434 }
435 }
436 Some("task_complete") => {
437 if let Some(text) = payload
438 .get("error")
439 .and_then(|error| error.get("message"))
440 .and_then(Value::as_str)
441 {
442 push_message(&mut event_error, ordinal, timestamp, text, &mut redactions);
443 }
444 }
445 Some("patch_apply_end") => record_changed_files(&mut tools, payload, cwd.as_deref()),
446 _ => {}
447 },
448 "response_item" => match payload.get("type").and_then(Value::as_str) {
449 Some("message") => {
450 let role = payload.get("role").and_then(Value::as_str);
451 if matches!(role, Some("user") | Some("assistant")) {
452 let text = message_content(payload, role.unwrap());
453 if !text.is_empty() && !looks_injected(&text) {
454 let target = if role == Some("user") {
455 &mut response_user
456 } else {
457 &mut response_assistant
458 };
459 push_message(target, ordinal, timestamp, &text, &mut redactions);
460 }
461 }
462 }
463 Some("agent_message") => {
468 if let Some(text) = item_text(payload) {
469 if !looks_injected(&text) {
470 push_message(
471 &mut response_assistant,
472 ordinal,
473 timestamp,
474 &text,
475 &mut redactions,
476 );
477 }
478 }
479 }
480 Some("function_call") | Some("custom_tool_call") => {
481 if let Some(name) = payload.get("name").and_then(Value::as_str) {
482 record_tool_name(&mut tools, name);
483 if let Some(call_id) = payload.get("call_id").and_then(Value::as_str) {
484 tool_call_names.insert(call_id.to_owned(), name.to_owned());
485 }
486 if let Some(detail) = extract_tool_detail(payload) {
487 push_tool_op(&mut tool_ops, timestamp, name, &detail, &mut redactions);
488 }
489 }
490 }
491 Some("function_call_output") | Some("custom_tool_call_output") => {
492 if let Some(text) = payload.get("output").and_then(extract_text_value) {
493 let name = payload
494 .get("call_id")
495 .and_then(Value::as_str)
496 .and_then(|call_id| tool_call_names.get(call_id))
497 .cloned()
498 .unwrap_or_else(|| "tool".to_owned());
499 let combined = format!("{name} -> {text}");
500 push_message(
501 &mut response_tool_output,
502 ordinal,
503 timestamp,
504 &combined,
505 &mut redactions,
506 );
507 }
508 }
509 _ => {}
510 },
511 "patch_apply_end" => record_changed_files(&mut tools, payload, cwd.as_deref()),
512 _ => {}
513 }
514 }
515
516 if tool_ops.len() > MAX_RECENT_TOOL_OPS {
517 let drop_count = tool_ops.len() - MAX_RECENT_TOOL_OPS;
518 tool_ops.drain(..drop_count);
519 truncated = true;
520 }
521
522 let mut messages = Vec::with_capacity(
523 event_user.len()
524 + response_user.len()
525 + event_assistant.len()
526 + response_assistant.len()
527 + response_tool_output.len()
528 + event_error.len(),
529 );
530 messages.extend(response_user.into_iter().map(|message| (MessageRole::User, message)));
531 messages.extend(event_user.into_iter().map(|message| (MessageRole::User, message)));
532 messages.extend(
533 response_assistant
534 .into_iter()
535 .map(|message| (MessageRole::Assistant, message)),
536 );
537 messages.extend(
538 event_assistant
539 .into_iter()
540 .map(|message| (MessageRole::Assistant, message)),
541 );
542 messages.extend(
543 response_tool_output
544 .into_iter()
545 .map(|message| (MessageRole::Tool, message)),
546 );
547 messages.extend(event_error.into_iter().map(|message| (MessageRole::Error, message)));
548 messages.sort_by_key(|(_, message)| message.ordinal);
549 let mut deduplicated = Vec::with_capacity(messages.len());
550 for message in messages {
551 let duplicate = deduplicated.last().is_some_and(
558 |(last_role, last_message): &(MessageRole, TimedMessage)| {
559 *last_role == message.0
560 && normalize_for_dedup(&last_message.text) == normalize_for_dedup(&message.1.text)
561 },
562 );
563 if !duplicate {
564 deduplicated.push(message);
565 }
566 }
567 let mut messages = deduplicated;
568 if messages.len() > limits.max_messages {
569 let drop_count = messages.len() - limits.max_messages;
570 messages.drain(..drop_count);
571 truncated = true;
572 }
573 let messages = messages
574 .into_iter()
575 .map(|(role, message)| Message {
576 role,
577 text: message.text,
578 timestamp: message.timestamp,
579 })
580 .collect();
581 let workspace_label = cwd
582 .as_deref()
583 .and_then(Path::file_name)
584 .and_then(OsStr::to_str)
585 .and_then(|value| safe_scalar(value, 128, &mut redactions));
586 let title = source
591 .title
592 .as_deref()
593 .and_then(|value| safe_scalar(value, 256, &mut redactions))
594 .or_else(|| {
595 first_human_prompt(&source.path)
596 .ok()
597 .flatten()
598 .and_then(|prompt| safe_scalar(&prompt, 256, &mut redactions))
599 });
600 let mut report = SessionReport {
601 schema: "codex-session-restore-v1",
602 meta: SessionMeta {
603 id: source.id.clone(),
604 title,
605 started_at: parsed_meta.timestamp,
606 updated_unix_ms: system_time_millis(
607 source.metadata.modified().unwrap_or(UNIX_EPOCH),
608 ),
609 workspace_label,
610 model,
611 model_provider: parsed_meta
612 .model_provider
613 .as_deref()
614 .and_then(|value| safe_scalar(value, 128, &mut redactions)),
615 },
616 messages,
617 tool_ops,
618 tools,
619 git: parsed_meta.git,
620 truncated,
621 malformed_records,
622 redactions,
623 };
624 enforce_output_bound(&mut report)?;
625 Ok(report)
626}
627
628pub fn encode_json<T: Serialize>(value: &T) -> Result<String, RestoreError> {
629 let encoded = serde_json::to_string_pretty(value)?;
630 if encoded.len() > MAX_OUTPUT_BYTES {
631 return Err(RestoreError::OutputLimit);
632 }
633 Ok(encoded)
634}
635
636pub fn render_report(report: &SessionReport) -> String {
637 let mut output = String::new();
638 output.push_str("Codex session restore report\n");
639 output.push_str(&format!("session_id: {}\n", report.meta.id));
640 if let Some(title) = &report.meta.title {
641 output.push_str(&format!("title: {title}\n"));
642 }
643 if let Some(started_at) = &report.meta.started_at {
644 output.push_str(&format!("started_at: {started_at}\n"));
645 }
646 output.push_str(&format!(
647 "updated_unix_ms: {}\n",
648 report.meta.updated_unix_ms
649 ));
650 if let Some(workspace) = &report.meta.workspace_label {
651 output.push_str(&format!("workspace_label: {workspace}\n"));
652 }
653 if let Some(model) = &report.meta.model {
654 output.push_str(&format!("model: {model}\n"));
655 }
656 output.push_str(&format!(
657 "truncated: {}\nmalformed_records: {}\nredactions: {}\n",
658 report.truncated, report.malformed_records, report.redactions
659 ));
660 output.push_str("messages:\n");
661 for message in &report.messages {
662 let role = match message.role {
663 MessageRole::User => "user",
664 MessageRole::Assistant => "assistant",
665 MessageRole::Tool => "tool",
666 MessageRole::Error => "error",
667 };
668 output.push_str(&format!("- {role}: {}\n", message.text.replace('\n', " ")));
669 }
670 if !report.tool_ops.is_empty() {
671 output.push_str("recent_tool_operations:\n");
672 for op in &report.tool_ops {
673 output.push_str(&format!("- {}: {}\n", op.name, op.detail.replace('\n', " ")));
674 }
675 }
676 if !report.tools.counts.is_empty() {
677 output.push_str("tool_counts:\n");
678 for (name, count) in &report.tools.counts {
679 output.push_str(&format!("- {name}: {count}\n"));
680 }
681 }
682 if !report.tools.changed_files.is_empty() {
683 output.push_str("changed_files:\n");
684 for path in &report.tools.changed_files {
685 output.push_str(&format!("- {path}\n"));
686 }
687 }
688 if let Some(branch) = &report.git.recorded_branch {
689 output.push_str(&format!("recorded_branch: {branch}\n"));
690 }
691 if let Some(commit) = &report.git.recorded_commit {
692 output.push_str(&format!("recorded_commit: {commit}\n"));
693 }
694 output
695}
696
697fn trusted_sessions_root(home: &Path) -> Result<PathBuf, RestoreError> {
698 let home_meta = fs::symlink_metadata(home).map_err(|_| RestoreError::HomeUnavailable)?;
699 if !home_meta.is_dir() || metadata_is_reparse(&home_meta) {
700 return Err(RestoreError::HomeUnavailable);
701 }
702 let root = home.join("sessions");
703 let metadata = fs::symlink_metadata(&root).map_err(|_| RestoreError::HomeUnavailable)?;
704 if !metadata.is_dir() || metadata_is_reparse(&metadata) {
705 return Err(RestoreError::HomeUnavailable);
706 }
707 fs::canonicalize(root).map_err(RestoreError::Io)
708}
709
710fn discover_session_paths(root: &Path) -> Result<Vec<PathBuf>, RestoreError> {
711 let mut paths = Vec::new();
712 for year in safe_read_directories(root)? {
713 for month in safe_read_directories(&year)? {
714 for day in safe_read_directories(&month)? {
715 for entry in fs::read_dir(&day)? {
716 let entry = entry?;
717 if paths.len() >= MAX_SESSION_FILES {
718 return Err(RestoreError::InvalidArgument(
719 "session file inventory exceeds the supported bound".to_owned(),
720 ));
721 }
722 let path = entry.path();
723 if path.extension() == Some(OsStr::new("jsonl"))
724 && rollout_filename_id(&path).is_some()
725 {
726 paths.push(path);
727 }
728 }
729 }
730 }
731 }
732 Ok(paths)
733}
734
735fn safe_read_directories(root: &Path) -> Result<Vec<PathBuf>, RestoreError> {
736 let mut result = Vec::new();
737 for entry in fs::read_dir(root)? {
738 let entry = entry?;
739 let file_type = entry.file_type()?;
740 if file_type.is_dir() && !file_type.is_symlink() {
741 let metadata = fs::symlink_metadata(entry.path())?;
742 if !metadata_is_reparse(&metadata) {
743 result.push(entry.path());
744 }
745 }
746 }
747 Ok(result)
748}
749
750fn safe_candidate_metadata(root: &Path, path: &Path) -> Result<Metadata, RestoreError> {
751 let link_meta = fs::symlink_metadata(path).map_err(|_| RestoreError::NotFound)?;
752 if !link_meta.is_file() || metadata_is_reparse(&link_meta) {
753 return Err(RestoreError::UnsafeCandidate);
754 }
755 let canonical = fs::canonicalize(path).map_err(|_| RestoreError::UnsafeCandidate)?;
756 ensure_descendant(root, &canonical)?;
757 let metadata = fs::metadata(&canonical)?;
758 if !metadata.is_file() {
759 return Err(RestoreError::UnsafeCandidate);
760 }
761 Ok(metadata)
762}
763
764fn ensure_descendant(root: &Path, path: &Path) -> Result<(), RestoreError> {
765 if path == root || !path.starts_with(root) {
766 return Err(RestoreError::UnsafeCandidate);
767 }
768 Ok(())
769}
770
771fn read_session_index(home: &Path) -> Result<BTreeMap<String, String>, RestoreError> {
772 let path = home.join("session_index.jsonl");
773 let metadata = match fs::symlink_metadata(&path) {
774 Ok(metadata) => metadata,
775 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
776 Err(error) => return Err(error.into()),
777 };
778 if !metadata.is_file() || metadata_is_reparse(&metadata) || metadata.len() > MAX_INDEX_BYTES {
779 return Ok(BTreeMap::new());
780 }
781 let mut file = open_shared_read(&path)?;
782 let snapshot_len = file.metadata()?.len();
783 if snapshot_len > MAX_INDEX_BYTES {
784 return Ok(BTreeMap::new());
785 }
786 let mut bytes = Vec::with_capacity(snapshot_len as usize);
787 (&mut file).take(snapshot_len).read_to_end(&mut bytes)?;
788 if bytes.len() as u64 != snapshot_len {
789 return Err(RestoreError::UnsafeCandidate);
790 }
791 let mut titles = BTreeMap::new();
792 for line in bytes.split(|byte| *byte == b'\n') {
793 if line.len() > MAX_LINE_BYTES {
794 continue;
795 }
796 let Ok(value) = serde_json::from_slice::<Value>(line) else {
797 continue;
798 };
799 let Some(id) = value.get("id").and_then(Value::as_str).filter(|id| valid_uuid(id)) else {
800 continue;
801 };
802 let Some(title) = value.get("thread_name").and_then(Value::as_str) else {
803 continue;
804 };
805 let mut redactions = 0;
806 if let Some(title) = safe_scalar(title, 256, &mut redactions) {
807 titles.insert(id.to_owned(), title);
808 }
809 }
810 Ok(titles)
811}
812
813fn read_session_meta(path: &Path, metadata: &Metadata) -> Result<Option<ParsedMeta>, RestoreError> {
814 let mut file = open_shared_read(path)?;
815 let mut reader = BufReader::new((&mut file).take(MAX_HEAD_BYTES as u64));
816 let mut line = Vec::with_capacity((metadata.len() as usize).min(8192));
817 if reader.read_until(b'\n', &mut line)? == 0 {
818 return Ok(None);
819 }
820 if line.last() == Some(&b'\n') {
821 line.pop();
822 if line.last() == Some(&b'\r') {
823 line.pop();
824 }
825 }
826 if line.len() > MAX_LINE_BYTES {
827 return Ok(None);
828 }
829 let line = std::str::from_utf8(&line).map_err(|_| RestoreError::NoSessionMeta)?;
830 parse_session_meta(line)
831}
832
833fn parse_session_meta(line: &str) -> Result<Option<ParsedMeta>, RestoreError> {
834 let value: Value = serde_json::from_str(line)?;
835 if value.get("type").and_then(Value::as_str) != Some("session_meta") {
836 return Ok(None);
837 }
838 let payload = value.get("payload").ok_or(RestoreError::NoSessionMeta)?;
839 let id = payload
840 .get("id")
841 .and_then(Value::as_str)
842 .filter(|value| valid_uuid(value))
843 .ok_or(RestoreError::NoSessionMeta)?
844 .to_owned();
845 let git_value = payload.get("git").unwrap_or(&Value::Null);
846 let mut ignored_redactions = 0;
847 let recorded_branch = git_value
848 .get("branch")
849 .and_then(Value::as_str)
850 .and_then(|value| safe_scalar(value, 256, &mut ignored_redactions));
851 let recorded_commit = git_value
852 .get("commit_hash")
853 .and_then(Value::as_str)
854 .filter(|value| (7..=64).contains(&value.len()) && value.chars().all(|ch| ch.is_ascii_hexdigit()))
855 .map(str::to_owned);
856 let repository_label = git_value
857 .get("repository_url")
858 .and_then(Value::as_str)
859 .and_then(repository_label)
860 .and_then(|value| safe_scalar(&value, 128, &mut ignored_redactions));
861 Ok(Some(ParsedMeta {
862 id,
863 timestamp: value
864 .get("timestamp")
865 .and_then(Value::as_str)
866 .map(bounded_scalar),
867 cwd: payload.get("cwd").and_then(Value::as_str).map(PathBuf::from),
868 model_provider: payload
869 .get("model_provider")
870 .and_then(Value::as_str)
871 .map(str::to_owned),
872 git: GitHints {
873 recorded_branch,
874 recorded_commit,
875 repository_label,
876 },
877 }))
878}
879
880fn first_human_prompt(path: &Path) -> Result<Option<String>, RestoreError> {
886 let mut file = open_shared_read(path)?;
887 let mut head = Vec::new();
888 (&mut file).take(MAX_HEAD_BYTES as u64).read_to_end(&mut head)?;
889 for raw in head.split(|byte| *byte == b'\n') {
890 if raw.is_empty() || raw.len() > MAX_LINE_BYTES {
891 continue;
892 }
893 let Ok(line) = std::str::from_utf8(raw) else {
894 continue;
895 };
896 let Ok(value) = serde_json::from_str::<Value>(line) else {
897 continue;
898 };
899 let record_type = value.get("type").and_then(Value::as_str).unwrap_or_default();
900 let payload = value.get("payload").unwrap_or(&Value::Null);
901 let text = match record_type {
902 "event_msg" if payload.get("type").and_then(Value::as_str) == Some("user_message") => {
903 payload.get("message").and_then(Value::as_str).map(str::to_owned)
904 }
905 "response_item"
906 if payload.get("type").and_then(Value::as_str) == Some("message")
907 && payload.get("role").and_then(Value::as_str) == Some("user") =>
908 {
909 let text = message_content(payload, "user");
910 (!text.is_empty()).then_some(text)
911 }
912 _ => None,
913 };
914 let Some(text) = text else {
915 continue;
916 };
917 if looks_injected(&text) {
918 continue;
919 }
920 let trimmed = text.trim();
921 if !trimmed.is_empty() {
922 return Ok(Some(trimmed.to_owned()));
923 }
924 }
925 Ok(None)
926}
927
928fn read_bounded_records(
929 source: &SessionSource,
930 limits: RestoreLimits,
931) -> Result<(String, Vec<String>, bool, u64), RestoreError> {
932 let mut file = open_shared_read(&source.path)?;
933 let snapshot_len = file.metadata()?.len();
934 let mut head = Vec::with_capacity((snapshot_len as usize).min(MAX_HEAD_BYTES));
935 (&mut file).take(MAX_HEAD_BYTES as u64).read_to_end(&mut head)?;
936 let meta_line = head
937 .split(|byte| *byte == b'\n')
938 .next()
939 .and_then(|line| std::str::from_utf8(line).ok())
940 .ok_or(RestoreError::NoSessionMeta)?
941 .to_owned();
942 let mut truncated = snapshot_len as usize > limits.max_tail_bytes;
943 let start = snapshot_len.saturating_sub(limits.max_tail_bytes as u64);
944 file.seek(SeekFrom::Start(start))?;
945 let snapshot_tail_len = snapshot_len - start;
946 let mut tail = Vec::with_capacity(snapshot_tail_len as usize);
947 (&mut file)
948 .take(snapshot_tail_len)
949 .read_to_end(&mut tail)?;
950 if tail.len() as u64 != snapshot_tail_len {
951 return Err(RestoreError::UnsafeCandidate);
952 }
953 if start > 0 {
954 if let Some(index) = tail.iter().position(|byte| *byte == b'\n') {
955 tail.drain(..=index);
956 } else {
957 tail.clear();
958 }
959 }
960 let mut malformed = 0_u64;
961 let mut lines = Vec::new();
962 for raw in tail.split(|byte| *byte == b'\n') {
963 if raw.is_empty() {
964 continue;
965 }
966 if raw.len() > MAX_LINE_BYTES {
967 malformed += 1;
968 truncated = true;
969 continue;
970 }
971 match std::str::from_utf8(raw) {
972 Ok(line) => lines.push(line.to_owned()),
973 Err(_) => malformed += 1,
974 }
975 }
976 if lines.len() > limits.max_lines {
977 let drop_count = lines.len() - limits.max_lines;
978 lines.drain(..drop_count);
979 truncated = true;
980 }
981 Ok((meta_line, lines, truncated, malformed))
982}
983
984#[cfg(windows)]
985fn open_shared_read(path: &Path) -> Result<File, RestoreError> {
986 use std::os::windows::fs::OpenOptionsExt;
987 const FILE_SHARE_READ: u32 = 0x00000001;
988 const FILE_SHARE_WRITE: u32 = 0x00000002;
989 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x00200000;
990 let file = OpenOptions::new()
991 .read(true)
992 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
993 .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
994 .open(path)?;
995 if metadata_is_reparse(&file.metadata()?) {
996 return Err(RestoreError::UnsafeCandidate);
997 }
998 Ok(file)
999}
1000
1001#[cfg(not(windows))]
1002fn open_shared_read(path: &Path) -> Result<File, RestoreError> {
1003 Ok(OpenOptions::new().read(true).open(path)?)
1004}
1005
1006#[cfg(windows)]
1007fn metadata_is_reparse(metadata: &Metadata) -> bool {
1008 use std::os::windows::fs::MetadataExt;
1009 const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x00000400;
1010 metadata.file_type().is_symlink()
1011 || metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0
1012}
1013
1014#[cfg(not(windows))]
1015fn metadata_is_reparse(metadata: &Metadata) -> bool {
1016 metadata.file_type().is_symlink()
1017}
1018
1019fn message_content(payload: &Value, role: &str) -> String {
1020 let expected = if role == "user" { "input_text" } else { "output_text" };
1021 payload
1022 .get("content")
1023 .and_then(Value::as_array)
1024 .into_iter()
1025 .flatten()
1026 .filter(|part| part.get("type").and_then(Value::as_str) == Some(expected))
1027 .filter_map(|part| part.get("text").and_then(Value::as_str))
1028 .collect::<Vec<_>>()
1029 .join("\n")
1030}
1031
1032fn extract_text_value(value: &Value) -> Option<String> {
1038 match value {
1039 Value::String(text) => Some(text.clone()),
1040 Value::Object(_) => value
1041 .get("content")
1042 .and_then(extract_text_value)
1043 .or_else(|| value.get("text").and_then(Value::as_str).map(str::to_owned)),
1044 Value::Array(items) => {
1045 let joined = items
1046 .iter()
1047 .filter_map(extract_text_value)
1048 .collect::<Vec<_>>()
1049 .join("\n");
1050 (!joined.is_empty()).then_some(joined)
1051 }
1052 _ => None,
1053 }
1054}
1055
1056fn item_text(item: &Value) -> Option<String> {
1062 item.get("message")
1063 .and_then(Value::as_str)
1064 .map(str::to_owned)
1065 .or_else(|| item.get("content").and_then(extract_text_value))
1066 .or_else(|| item.get("text").and_then(Value::as_str).map(str::to_owned))
1067}
1068
1069fn extract_tool_detail(payload: &Value) -> Option<String> {
1076 let raw = payload
1077 .get("arguments")
1078 .or_else(|| payload.get("input"))
1079 .and_then(Value::as_str)?;
1080 let detail = serde_json::from_str::<Value>(raw)
1081 .ok()
1082 .and_then(|parsed| salient_arg_field(&parsed))
1083 .unwrap_or_else(|| raw.to_owned());
1084 Some(detail)
1085}
1086
1087fn salient_arg_field(value: &Value) -> Option<String> {
1088 const KEYS: [&str; 8] = [
1089 "command", "cmd", "url", "query", "path", "file_path", "pattern", "script",
1090 ];
1091 let object = value.as_object()?;
1092 for key in KEYS {
1093 if let Some(text) = object.get(key).and_then(Value::as_str) {
1094 return Some(text.to_owned());
1095 }
1096 }
1097 None
1098}
1099
1100fn push_message(
1101 target: &mut Vec<TimedMessage>,
1102 ordinal: usize,
1103 timestamp: Option<String>,
1104 text: &str,
1105 redactions: &mut u64,
1106) {
1107 if looks_injected(text) {
1108 return;
1109 }
1110 let text = redact_text(text, redactions);
1111 let text = truncate_message_tail(text.trim(), MAX_MESSAGE_CHARS);
1112 if !text.is_empty() {
1113 target.push(TimedMessage {
1114 ordinal,
1115 timestamp,
1116 text,
1117 });
1118 }
1119}
1120
1121fn push_tool_op(
1125 target: &mut Vec<ToolOperation>,
1126 timestamp: Option<String>,
1127 name: &str,
1128 detail: &str,
1129 redactions: &mut u64,
1130) {
1131 if looks_injected(detail) {
1132 return;
1133 }
1134 let redacted = redact_text(detail, redactions);
1135 let first_line = redacted.lines().next().unwrap_or(&redacted);
1136 let detail = truncate_chars(first_line.trim(), TOOL_OP_DETAIL_CHARS);
1137 if !detail.is_empty() {
1138 target.push(ToolOperation {
1139 name: name.to_owned(),
1140 detail,
1141 timestamp,
1142 });
1143 }
1144}
1145
1146fn looks_injected(value: &str) -> bool {
1147 let lowered = value.to_ascii_lowercase();
1148 [
1149 "<environment_context>",
1150 "<permissions instructions>",
1151 "<collaboration_mode>",
1152 "<skills_instructions>",
1153 "<app-context>",
1154 "# agents.md instructions",
1155 "========= memory_summary begins =========",
1156 ]
1157 .iter()
1158 .any(|marker| lowered.contains(marker))
1159}
1160
1161fn record_tool_name(tools: &mut ToolInventory, name: &str) {
1162 if tools.counts.len() >= MAX_TOOL_NAMES && !tools.counts.contains_key(name) {
1163 return;
1164 }
1165 if !name.is_empty()
1166 && name.len() <= 64
1167 && name
1168 .chars()
1169 .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-' | '.' | ':'))
1170 {
1171 *tools.counts.entry(name.to_owned()).or_default() += 1;
1172 }
1173}
1174
1175fn record_changed_files(tools: &mut ToolInventory, payload: &Value, cwd: Option<&Path>) {
1176 let Some(cwd) = cwd else {
1177 return;
1178 };
1179 let Some(changes) = payload.get("changes").and_then(Value::as_object) else {
1180 return;
1181 };
1182 for path in changes.keys() {
1183 if tools.changed_files.len() >= MAX_FILE_HINTS {
1184 break;
1185 }
1186 let path = Path::new(path);
1187 let Ok(relative) = path.strip_prefix(cwd) else {
1188 continue;
1189 };
1190 if !safe_relative_path(relative) {
1191 continue;
1192 }
1193 let text = relative.to_string_lossy().replace('\\', "/");
1194 if text.len() <= 512 && !contains_credential_marker(&text) {
1195 tools.changed_files.insert(text);
1196 }
1197 }
1198}
1199
1200fn safe_relative_path(path: &Path) -> bool {
1201 path.components().next().is_some()
1202 && path.components().all(|component| match component {
1203 Component::Normal(value) => value
1204 .to_str()
1205 .is_some_and(|text| !text.is_empty() && !text.chars().any(char::is_control)),
1206 _ => false,
1207 })
1208}
1209
1210fn enforce_output_bound(report: &mut SessionReport) -> Result<(), RestoreError> {
1211 loop {
1212 let encoded = serde_json::to_vec(report)?;
1213 if encoded.len() <= MAX_OUTPUT_BYTES {
1214 return Ok(());
1215 }
1216 if report.messages.is_empty() {
1217 return Err(RestoreError::OutputLimit);
1218 }
1219 report.messages.remove(0);
1220 report.truncated = true;
1221 }
1222}
1223
1224fn redact_text(value: &str, redactions: &mut u64) -> String {
1225 let normalized = value.replace("\r\n", "\n").replace('\r', "\n");
1226 let mut output = String::new();
1227 let mut in_pem = false;
1228 for line in normalized.lines() {
1229 let lowered = line.to_ascii_lowercase();
1230 if lowered.contains("-----begin ") && lowered.contains("private key-----") {
1231 in_pem = true;
1232 *redactions += 1;
1233 if !output.is_empty() {
1234 output.push('\n');
1235 }
1236 output.push_str("[REDACTED]");
1237 continue;
1238 }
1239 if in_pem {
1240 if lowered.contains("-----end ") && lowered.contains("private key-----") {
1241 in_pem = false;
1242 }
1243 continue;
1244 }
1245 let mut clean: String = line.chars().filter(|ch| !ch.is_control() || *ch == '\t').collect();
1246 if credential_regex().is_match(&clean) || auth_header_regex().is_match(&clean) {
1247 *redactions += 1;
1248 clean = "[REDACTED]".to_owned();
1249 }
1250 clean = jwt_regex()
1251 .replace_all(&clean, |_: ®ex::Captures<'_>| {
1252 *redactions += 1;
1253 "[REDACTED]"
1254 })
1255 .into_owned();
1256 clean = token_regex()
1257 .replace_all(&clean, |_: ®ex::Captures<'_>| {
1258 *redactions += 1;
1259 "[REDACTED]"
1260 })
1261 .into_owned();
1262 clean = uri_userinfo_regex()
1263 .replace_all(&clean, |captures: ®ex::Captures<'_>| {
1264 *redactions += 1;
1265 format!("{}[REDACTED]@", &captures[1])
1266 })
1267 .into_owned();
1268 if !output.is_empty() {
1269 output.push('\n');
1270 }
1271 output.push_str(&clean);
1272 }
1273 output
1274}
1275
1276fn credential_regex() -> &'static Regex {
1277 static VALUE: OnceLock<Regex> = OnceLock::new();
1278 VALUE.get_or_init(|| {
1279 Regex::new(r#"(?i)(?:^|[^A-Za-z0-9_])["']?(?:password|passwd|secret|token|api[_-]?key|authorization|[A-Za-z0-9_]+_(?:password|passwd|secret|token|api[_-]?key|authorization))["']?\s*[:=]"#)
1280 .unwrap()
1281 })
1282}
1283
1284fn auth_header_regex() -> &'static Regex {
1285 static VALUE: OnceLock<Regex> = OnceLock::new();
1286 VALUE.get_or_init(|| Regex::new(r#"(?i)\b(?:bearer|basic)\s+[^\s,;"']+"#).unwrap())
1287}
1288
1289fn jwt_regex() -> &'static Regex {
1290 static VALUE: OnceLock<Regex> = OnceLock::new();
1291 VALUE.get_or_init(|| Regex::new(r"\beyJ[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]{8,}\b").unwrap())
1292}
1293
1294fn token_regex() -> &'static Regex {
1295 static VALUE: OnceLock<Regex> = OnceLock::new();
1296 VALUE.get_or_init(|| {
1297 Regex::new(r"\b(?:sk-[A-Za-z0-9_-]{16,}|gh[pousr]_[A-Za-z0-9_]{20,}|github_pat_[A-Za-z0-9_]{20,}|AKIA[A-Z0-9]{16})\b").unwrap()
1298 })
1299}
1300
1301fn uri_userinfo_regex() -> &'static Regex {
1302 static VALUE: OnceLock<Regex> = OnceLock::new();
1303 VALUE.get_or_init(|| Regex::new(r"([A-Za-z][A-Za-z0-9+.-]*://)[^/@\s]+:[^/@\s]+@").unwrap())
1304}
1305
1306fn contains_credential_marker(value: &str) -> bool {
1307 credential_regex().is_match(value)
1308 || auth_header_regex().is_match(value)
1309 || jwt_regex().is_match(value)
1310 || token_regex().is_match(value)
1311 || uri_userinfo_regex().is_match(value)
1312}
1313
1314fn safe_scalar(value: &str, max_chars: usize, redactions: &mut u64) -> Option<String> {
1315 let value = redact_text(value, redactions);
1316 let value = truncate_chars(value.trim(), max_chars);
1317 (!value.is_empty()).then_some(value)
1318}
1319
1320fn bounded_scalar(value: &str) -> String {
1321 truncate_chars(value.trim(), 128)
1322}
1323
1324fn truncate_chars(value: &str, max_chars: usize) -> String {
1329 let mut result: String = value.chars().take(max_chars).collect();
1330 if value.chars().count() > max_chars {
1331 result.push('…');
1332 }
1333 result
1334}
1335
1336fn truncate_message_tail(text: &str, max_chars: usize) -> String {
1345 let total_chars = text.chars().count();
1346 if total_chars <= max_chars {
1347 return text.to_owned();
1348 }
1349 let cut = total_chars - max_chars;
1350 let tail: String = text.chars().skip(cut).collect();
1351 format!("[truncated {cut} chars]…{tail}")
1352}
1353
1354fn normalize_for_dedup(text: &str) -> String {
1362 text.split_whitespace().collect::<Vec<_>>().join(" ")
1363}
1364
1365fn repository_label(value: &str) -> Option<String> {
1366 let without_query = value.split(['?', '#']).next().unwrap_or_default();
1367 let tail = without_query
1368 .trim_end_matches(['/', '\\'])
1369 .rsplit(['/', '\\', ':'])
1370 .next()
1371 .unwrap_or_default()
1372 .trim_end_matches(".git");
1373 (!tail.is_empty()).then_some(tail.to_owned())
1374}
1375
1376fn rollout_filename_id(path: &Path) -> Option<String> {
1377 let stem = path.file_stem()?.to_str()?;
1378 let id = stem.rsplit('-').take(5).collect::<Vec<_>>();
1379 if id.len() != 5 {
1380 return None;
1381 }
1382 let candidate = format!("{}-{}-{}-{}-{}", id[4], id[3], id[2], id[1], id[0]);
1383 valid_uuid(&candidate).then_some(candidate)
1384}
1385
1386fn valid_uuid(value: &str) -> bool {
1387 if value.len() != 36 {
1388 return false;
1389 }
1390 value.chars().enumerate().all(|(index, ch)| {
1391 if matches!(index, 8 | 13 | 18 | 23) {
1392 ch == '-'
1393 } else {
1394 ch.is_ascii_hexdigit() && !ch.is_ascii_uppercase()
1395 }
1396 })
1397}
1398
1399fn is_id_selector_char(ch: char) -> bool {
1400 ch == '-' || (ch.is_ascii_hexdigit() && !ch.is_ascii_uppercase())
1401}
1402
1403fn system_time_millis(value: SystemTime) -> u64 {
1404 value
1405 .duration_since(UNIX_EPOCH)
1406 .unwrap_or_default()
1407 .as_millis()
1408 .min(u64::MAX as u128) as u64
1409}
1410
1411#[cfg(test)]
1412mod tests {
1413 use super::*;
1414 use std::io::Write;
1415 use std::sync::atomic::{AtomicBool, Ordering};
1416 use std::sync::{mpsc, Arc};
1417 use std::thread;
1418
1419 fn fixture_home() -> tempfile::TempDir {
1420 let temp = tempfile::tempdir().unwrap();
1421 fs::create_dir_all(temp.path().join("sessions/2026/08/10")).unwrap();
1422 temp
1423 }
1424
1425 fn write_session(home: &Path, id: &str, records: &[Value]) -> PathBuf {
1426 let path = home
1427 .join("sessions/2026/08/10")
1428 .join(format!("rollout-2026-08-10T10-00-00-{id}.jsonl"));
1429 let mut file = File::create(&path).unwrap();
1430 let meta = serde_json::json!({
1431 "timestamp": "2026-08-10T10:00:00Z",
1432 "type": "session_meta",
1433 "payload": {
1434 "id": id,
1435 "cwd": "C:\\work\\demo",
1436 "model_provider": "openai",
1437 "git": {"branch":"feature/restore","commit_hash":"0123456789abcdef","repository_url":"https://user:pass@example.invalid/acme/demo.git"}
1438 }
1439 });
1440 writeln!(file, "{}", serde_json::to_string(&meta).unwrap()).unwrap();
1441 for record in records {
1442 writeln!(file, "{}", serde_json::to_string(record).unwrap()).unwrap();
1443 }
1444 path
1445 }
1446
1447 fn source(home: &Path, id: &str) -> SessionSource {
1448 resolve_target(home, OsStr::new(id)).unwrap()
1449 }
1450
1451 #[test]
1452 fn parser_surfaces_only_user_and_assistant_and_excludes_hidden_records() {
1453 let temp = fixture_home();
1454 let id = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa";
1455 write_session(
1456 temp.path(),
1457 id,
1458 &[
1459 serde_json::json!({"timestamp":"1","type":"response_item","payload":{"type":"reasoning","summary":["HIDDEN_REASONING"]}}),
1460 serde_json::json!({"timestamp":"2","type":"response_item","payload":{"type":"message","role":"developer","content":[{"type":"input_text","text":"HIDDEN_DEVELOPER"}]}}),
1461 serde_json::json!({"timestamp":"3","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"fallback user"}]}}),
1462 serde_json::json!({"timestamp":"4","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"fallback assistant"}]}}),
1463 serde_json::json!({"timestamp":"5","type":"event_msg","payload":{"type":"user_message","message":"Choose checked_add"}}),
1464 serde_json::json!({"timestamp":"6","type":"event_msg","payload":{"type":"agent_message","message":"Preserve the public API"}}),
1465 serde_json::json!({"timestamp":"7","type":"response_item","payload":{"type":"function_call","name":"shell_command","call_id":"call_1","arguments":"VISIBLE_COMMAND_ARG"}}),
1466 serde_json::json!({"timestamp":"8","type":"response_item","payload":{"type":"function_call_output","call_id":"call_1","output":"VISIBLE_TOOL_OUTPUT"}}),
1467 serde_json::json!({"timestamp":"9","type":"compacted","payload":{"replacement_history":"HIDDEN_COMPACTION"}}),
1468 serde_json::json!({"timestamp":"10","type":"event_msg","payload":{"type":"patch_apply_end","changes":{"C:\\work\\demo\\src\\lib.rs":{"kind":"update"}},"stdout":"HIDDEN_PATCH_OUTPUT"}}),
1469 ],
1470 );
1471 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1472 assert_eq!(
1473 report
1474 .messages
1475 .iter()
1476 .filter(|message| matches!(message.role, MessageRole::User | MessageRole::Assistant))
1477 .count(),
1478 4
1479 );
1480 assert!(report.messages.iter().any(|message| message.text == "fallback user"));
1481 assert!(report.messages.iter().any(|message| message.text == "fallback assistant"));
1482 assert!(report.messages.iter().any(|message| message.text == "Choose checked_add"));
1483 assert!(report.messages.iter().any(|message| message.text == "Preserve the public API"));
1484 assert_eq!(report.tool_ops.len(), 1);
1486 assert_eq!(report.tool_ops[0].name, "shell_command");
1487 assert_eq!(report.tool_ops[0].detail, "VISIBLE_COMMAND_ARG");
1488 assert!(report
1490 .messages
1491 .iter()
1492 .any(|message| message.role == MessageRole::Tool && message.text.contains("VISIBLE_TOOL_OUTPUT")));
1493 let encoded = encode_json(&report).unwrap();
1494 for hidden in ["HIDDEN_REASONING", "HIDDEN_DEVELOPER", "HIDDEN_COMPACTION"] {
1495 assert!(!encoded.contains(hidden));
1496 }
1497 assert!(encoded.contains("VISIBLE_COMMAND_ARG"));
1498 assert!(encoded.contains("VISIBLE_TOOL_OUTPUT"));
1499 assert_eq!(report.tools.counts.get("shell_command"), Some(&1));
1500 assert!(report.tools.changed_files.contains("src/lib.rs"));
1501 assert!(!encoded.contains("HIDDEN_PATCH_OUTPUT"));
1502 }
1503
1504 #[test]
1505 fn event_msg_is_fallback_without_duplicates() {
1506 let temp = fixture_home();
1507 let id = "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb";
1508 write_session(
1509 temp.path(),
1510 id,
1511 &[
1512 serde_json::json!({"timestamp":"1","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"canonical user"}]}}),
1513 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"user_message","message":"canonical user"}}),
1514 serde_json::json!({"timestamp":"3","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"assistant fallback"}]}}),
1515 ],
1516 );
1517 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1518 assert_eq!(report.messages.iter().filter(|m| m.role == MessageRole::User).count(), 1);
1519 assert!(report.messages.iter().any(|m| m.text == "canonical user"));
1520 assert!(report.messages.iter().any(|m| m.text == "assistant fallback"));
1521 }
1522
1523 #[test]
1524 fn dedup_collapses_dual_encoding_whitespace_variant_but_keeps_distinct_messages() {
1525 let temp = fixture_home();
1526 let id = "66666666-6666-4666-8666-666666666666";
1527 write_session(
1528 temp.path(),
1529 id,
1530 &[
1531 serde_json::json!({"timestamp":"1","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"line one\n"},{"type":"input_text","text":"line two"}]}}),
1535 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"user_message","message":"line one\nline two"}}),
1536 serde_json::json!({"timestamp":"3","type":"event_msg","payload":{"type":"user_message","message":"a genuinely different follow-up"}}),
1537 ],
1538 );
1539 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1540 let user_messages: Vec<_> =
1541 report.messages.iter().filter(|m| m.role == MessageRole::User).collect();
1542 assert_eq!(
1543 user_messages.len(),
1544 2,
1545 "dual-encoding pair must collapse but the distinct follow-up must survive: {user_messages:?}"
1546 );
1547 assert!(user_messages[0].text.contains("line one"));
1548 assert!(user_messages[1].text.contains("a genuinely different follow-up"));
1549 }
1550
1551 #[test]
1552 fn tool_call_arguments_are_captured_as_recent_tool_operations() {
1553 let temp = fixture_home();
1554 let id = "22222222-2222-4222-8222-222222222222";
1555 write_session(
1556 temp.path(),
1557 id,
1558 &[
1559 serde_json::json!({"timestamp":"1","type":"response_item","payload":{"type":"function_call","name":"shell","call_id":"call_1","arguments":"{\"command\":\"rg -n TODO\"}"}}),
1560 serde_json::json!({"timestamp":"2","type":"response_item","payload":{"type":"custom_tool_call","name":"exec","call_id":"call_2","input":"tools.exec_command({cmd:\"ls -la\"})"}}),
1561 ],
1562 );
1563 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1564 assert_eq!(report.tool_ops.len(), 2);
1565 assert_eq!(report.tool_ops[0].name, "shell");
1566 assert_eq!(report.tool_ops[0].detail, "rg -n TODO");
1567 assert_eq!(report.tool_ops[1].name, "exec");
1568 assert!(report.tool_ops[1].detail.contains("tools.exec_command"));
1569 }
1570
1571 #[test]
1572 fn tool_output_and_task_complete_error_are_surfaced() {
1573 let temp = fixture_home();
1574 let id = "33333333-3333-4333-8333-333333333333";
1575 write_session(
1576 temp.path(),
1577 id,
1578 &[
1579 serde_json::json!({"timestamp":"1","type":"response_item","payload":{"type":"function_call","name":"shell","call_id":"call_9","arguments":"{\"command\":\"ls\"}"}}),
1580 serde_json::json!({"timestamp":"2","type":"response_item","payload":{"type":"function_call_output","call_id":"call_9","output":"total 0\nfile.txt"}}),
1581 serde_json::json!({"timestamp":"3","type":"response_item","payload":{"type":"agent_message","content":[{"type":"text","text":"Agent errored: usage limit reached"}]}}),
1582 serde_json::json!({"timestamp":"4","type":"event_msg","payload":{"type":"task_complete","error":{"message":"You have hit your usage limit"}}}),
1583 ],
1584 );
1585 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1586 let tool_message = report
1587 .messages
1588 .iter()
1589 .find(|m| m.role == MessageRole::Tool)
1590 .expect("tool output message present");
1591 assert!(tool_message.text.contains("shell -> total 0"));
1592 let error_message = report
1593 .messages
1594 .iter()
1595 .find(|m| m.role == MessageRole::Error)
1596 .expect("error message present");
1597 assert_eq!(error_message.text, "You have hit your usage limit");
1598 assert!(report
1599 .messages
1600 .iter()
1601 .any(|m| m.role == MessageRole::Assistant && m.text.contains("Agent errored")));
1602 }
1603
1604 #[test]
1605 fn item_completed_event_is_handled_for_forward_compatibility() {
1606 let temp = fixture_home();
1607 let id = "77777777-7777-4777-8777-777777777777";
1608 write_session(
1609 temp.path(),
1610 id,
1611 &[
1612 serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"item_completed","item":{"type":"user_message","message":"future schema user turn"}}}),
1613 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"item_completed","item":{"type":"agent_message","message":"future schema assistant turn"}}}),
1614 ],
1615 );
1616 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1617 assert!(report
1618 .messages
1619 .iter()
1620 .any(|m| m.role == MessageRole::User && m.text == "future schema user turn"));
1621 assert!(report
1622 .messages
1623 .iter()
1624 .any(|m| m.role == MessageRole::Assistant && m.text == "future schema assistant turn"));
1625 }
1626
1627 #[test]
1628 fn title_falls_back_to_first_human_prompt_when_untitled() {
1629 let temp = fixture_home();
1630 let id = "44444444-4444-4444-8444-444444444444";
1631 write_session(
1632 temp.path(),
1633 id,
1634 &[
1635 serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"user_message","message":"Investigate the failing build"}}),
1636 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"agent_message","message":"Looking into it"}}),
1637 ],
1638 );
1639 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1640 assert_eq!(report.meta.title.as_deref(), Some("Investigate the failing build"));
1641
1642 let candidates = list_sessions(temp.path(), None, 10).unwrap();
1643 let candidate = candidates.iter().find(|c| c.id == id).expect("candidate present");
1644 assert_eq!(candidate.title.as_deref(), Some("Investigate the failing build"));
1645 }
1646
1647 #[test]
1648 fn long_message_bodies_are_tail_truncated_with_marker() {
1649 let temp = fixture_home();
1650 let id = "55555555-5555-4555-8555-555555555555";
1651 let filler = "A".repeat(5000);
1652 let actionable_tail = "APPROVAL REQUEST: run rm -rf /tmp/example";
1653 let long_message = format!("{filler}{actionable_tail}");
1654 write_session(
1655 temp.path(),
1656 id,
1657 &[serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"user_message","message": long_message}})],
1658 );
1659 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1660 let message = &report.messages[0];
1661 assert!(message.text.starts_with("[truncated "));
1662 assert!(
1663 message.text.ends_with(actionable_tail),
1664 "tail must keep the actionable ending, got: {}",
1665 message.text
1666 );
1667 }
1668
1669 #[test]
1670 fn redactor_removes_credentials_from_all_report_fields() {
1671 let temp = fixture_home();
1672 let id = "cccccccc-cccc-4ccc-8ccc-cccccccccccc";
1673 write_session(
1674 temp.path(),
1675 id,
1676 &[
1677 serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"user_message","message":"\"api_key\": \"JSON_SECRET_MUST_NOT_ESCAPE\"\nOPENAI_API_KEY=ENV_SECRET_MUST_NOT_ESCAPE\nuse sk-ABCDEFGHIJKLMNOPQRSTUV"}}),
1678 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"agent_message","message":"Bearer dotted.secret.value"}}),
1679 ],
1680 );
1681 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1682 let encoded = encode_json(&report).unwrap();
1683 assert!(!encoded.contains("JSON_SECRET_MUST_NOT_ESCAPE"));
1684 assert!(!encoded.contains("ENV_SECRET_MUST_NOT_ESCAPE"));
1685 assert!(!encoded.contains("sk-ABCDEFGHIJKLMNOPQRSTUV"));
1686 assert!(!encoded.contains("dotted.secret.value"));
1687 assert!(report.redactions >= 4);
1688 assert_eq!(report.git.repository_label.as_deref(), Some("demo"));
1689 }
1690
1691 #[test]
1692 fn redactor_preserves_noncredential_policy_and_budget_fields() {
1693 let mut redactions = 0;
1694 let input = "token_budget=4096\ntoken_count: 20\nauthorization_policy=deny";
1695 let output = redact_text(input, &mut redactions);
1696 assert_eq!(output, input);
1697 assert_eq!(redactions, 0);
1698 }
1699
1700 #[test]
1701 fn resolve_unique_prefix_and_reject_ambiguous() {
1702 let temp = fixture_home();
1703 let first = "dddddddd-dddd-4ddd-8ddd-dddddddddddd";
1704 let second = "dddddddd-dddd-4ddd-8ddd-ddddddddddde";
1705 write_session(temp.path(), first, &[]);
1706 write_session(temp.path(), second, &[]);
1707 assert!(matches!(
1708 resolve_target(temp.path(), OsStr::new("dddddddd-dddd-4ddd")),
1709 Err(RestoreError::AmbiguousPrefix)
1710 ));
1711 let resolved = resolve_target(temp.path(), OsStr::new(first)).unwrap();
1712 assert_eq!(resolved.id, first);
1713 }
1714
1715 #[test]
1716 fn exact_id_resolution_is_not_limited_to_the_newest_hundred_sessions() {
1717 let temp = fixture_home();
1718 let target = "00000000-0000-4000-8000-000000000000";
1719 write_session(temp.path(), target, &[]);
1720 for index in 1..=120_u64 {
1721 let id = format!(
1722 "{index:08x}-0000-4000-8000-{index:012x}"
1723 );
1724 write_session(temp.path(), &id, &[]);
1725 }
1726 let resolved = resolve_target(temp.path(), OsStr::new(target)).unwrap();
1727 assert_eq!(resolved.id, target);
1728 }
1729
1730 #[test]
1731 fn bounded_reader_never_reads_more_than_tail_and_keeps_latest_messages() {
1732 let temp = fixture_home();
1733 let id = "eeeeeeee-eeee-4eee-8eee-eeeeeeeeeeee";
1734 let records = (0..200)
1735 .map(|index| serde_json::json!({"timestamp":index.to_string(),"type":"event_msg","payload":{"type":"user_message","message":format!("message-{index:03}")}}))
1736 .collect::<Vec<_>>();
1737 write_session(temp.path(), id, &records);
1738 let report = load_session(
1739 &source(temp.path(), id),
1740 RestoreLimits {
1741 max_tail_bytes: 4096,
1742 max_lines: 20,
1743 max_messages: 3,
1744 },
1745 )
1746 .unwrap();
1747 assert!(report.truncated);
1748 assert_eq!(report.messages.len(), 3);
1749 assert!(report.messages.last().unwrap().text.contains("199"));
1750 }
1751
1752 #[test]
1753 fn active_rollout_growth_cannot_extend_the_opened_snapshot_read() {
1754 let temp = fixture_home();
1755 let id = "11111111-1111-4111-8111-111111111111";
1756 let path = write_session(
1757 temp.path(),
1758 id,
1759 &[serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"user_message","message":"bounded active session"}})],
1760 );
1761 let source = source(temp.path(), id);
1762 let keep_writing = Arc::new(AtomicBool::new(true));
1763 let writer_flag = Arc::clone(&keep_writing);
1764 let writer = thread::spawn(move || {
1765 let mut file = OpenOptions::new().append(true).open(path).unwrap();
1766 let record = serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"agent_message","message":"active append"}}).to_string();
1767 while writer_flag.load(Ordering::Acquire) {
1768 writeln!(file, "{record}").unwrap();
1769 file.flush().unwrap();
1770 thread::sleep(Duration::from_millis(1));
1771 }
1772 });
1773 thread::sleep(Duration::from_millis(20));
1774 let (send, receive) = mpsc::channel();
1775 let reader = thread::spawn(move || {
1776 let result = load_session(
1777 &source,
1778 RestoreLimits {
1779 max_tail_bytes: 4096,
1780 max_lines: 64,
1781 max_messages: 8,
1782 },
1783 );
1784 let _ = send.send(result.map(|report| report.messages.len()));
1785 });
1786 let result = receive.recv_timeout(Duration::from_secs(2));
1787 keep_writing.store(false, Ordering::Release);
1788 writer.join().unwrap();
1789 reader.join().unwrap();
1790 assert!(result.unwrap().is_ok());
1791 }
1792
1793 #[test]
1794 fn injected_context_is_not_surfaced() {
1795 let temp = fixture_home();
1796 let id = "ffffffff-ffff-4fff-8fff-ffffffffffff";
1797 write_session(
1798 temp.path(),
1799 id,
1800 &[
1801 serde_json::json!({"timestamp":"1","type":"event_msg","payload":{"type":"user_message","message":"<environment_context>HIDDEN_ENV</environment_context>"}}),
1802 serde_json::json!({"timestamp":"2","type":"event_msg","payload":{"type":"user_message","message":"Real user decision"}}),
1803 ],
1804 );
1805 let report = load_session(&source(temp.path(), id), RestoreLimits::default()).unwrap();
1806 assert_eq!(report.messages.len(), 1);
1807 assert_eq!(report.messages[0].text, "Real user decision");
1808 }
1809}