Skip to main content

codex_session_restore/
lib.rs

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;
24/// How many of the most recent tool calls are kept for the "Recent Tool
25/// Operations" section — a compact, ordered trace of what actually ran,
26/// separate from the full-session name histogram in `tools.counts`.
27const MAX_RECENT_TOOL_OPS: usize = 15;
28/// Cap on a single tool operation's key-argument excerpt.
29const 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    /// A tool call's output (`function_call_output` / `custom_tool_call_output`).
123    Tool,
124    /// A session-level failure (`event_msg.task_complete.error.message`).
125    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/// A single tool invocation's key argument, captured verbatim (bounded and
136/// redacted) so the "last N events" digest shows what actually ran instead
137/// of only a name tally. See `ToolInventory::counts` for the full-session
138/// histogram, which stays as a secondary, coarser view.
139#[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        // Titles come from session_index.jsonl when present. When the
253        // index has none, pay for one bounded head-read of the rollout
254        // itself and fall back to the first real human prompt, instead of
255        // leaving every untitled session looking identical in `list`.
256        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                // Forward-compat: current Codex builds wrap turns as
403                // `item_completed { item: { type: "user_message"/"agent_message", ... } }`
404                // instead of the flat `user_message`/`agent_message` variants above.
405                // Every session inspected so far also carries the same content via
406                // `response_item`, so this arm is additive, not yet load-bearing.
407                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                // A bare `agent_message` at the response_item level (e.g. a
464                // subagent's own reply) carries no `role` field and uses a
465                // `content` array instead of `message`; surface it as
466                // assistant text, same as the event_msg fallback above.
467                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        // Codex writes the same human/agent turn twice under two encodings
552        // (`response_item.message`, whose parts this parser joins with an
553        // extra "\n", and the flat `event_msg.*` string) that differ only by
554        // incidental whitespace. Compare on collapsed whitespace so that
555        // pair merges into one entry, while two messages with different
556        // words never collapse into each other.
557        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    // Provider title first; when the store has none, fall back to the
587    // earliest real human prompt (bounded head scan, independent of the
588    // tail-bounded record read above so it still works on sessions whose
589    // opening turn falls outside that window).
590    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
880/// Scan the first `MAX_HEAD_BYTES` of a rollout file for the earliest
881/// genuine human prompt, used as a topic fallback when a session has no
882/// title in `session_index.jsonl`. This is independent of `load`'s
883/// tail-bounded record read, so it still finds the opening prompt on a
884/// session whose first turn falls outside that tail window.
885fn 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
1032/// Extract plain text from a JSON value that may be a bare string, an
1033/// object carrying `content`/`text`, or an array of such parts. Codex uses
1034/// all three shapes across `function_call_output.output`,
1035/// `custom_tool_call_output.output`, and bare `response_item.agent_message`
1036/// / `item_completed` item payloads.
1037fn 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
1056/// Text of a `user_message`/`agent_message` item under either shape Codex
1057/// has used: a flat `message` string (what `event_msg`'s direct variants
1058/// carry) or a `content` array / `text` field (what a bare
1059/// `response_item.agent_message` and the newer `item_completed` wrapper
1060/// carry instead).
1061fn 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
1069/// The key argument of a tool call — the shell command, URL, query, or path
1070/// a reviewer actually needs to see, not just the tool's name.
1071/// `custom_tool_call.input` and `function_call.arguments` are read as
1072/// either a JSON object (pull a well-known field) or, when that fails
1073/// (`custom_tool_call.input` is often raw script text, not JSON), the raw
1074/// string itself, bounded and redacted the same as any other excerpt.
1075fn 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
1121/// Push a tool operation onto `target`. Callers append in ordinal order
1122/// during a single forward pass over the record stream, so `target` is
1123/// already chronological without a separate sort step.
1124fn 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, |_: &regex::Captures<'_>| {
1252                *redactions += 1;
1253                "[REDACTED]"
1254            })
1255            .into_owned();
1256        clean = token_regex()
1257            .replace_all(&clean, |_: &regex::Captures<'_>| {
1258                *redactions += 1;
1259                "[REDACTED]"
1260            })
1261            .into_owned();
1262        clean = uri_userinfo_regex()
1263            .replace_all(&clean, |captures: &regex::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
1324/// Truncate a short scalar (title, model, workspace label, …) to at most
1325/// `max_chars`, keeping the HEAD. These are identifiers, not narrative
1326/// text — the meaningful part is at the front. Message bodies use
1327/// [`truncate_message_tail`] instead.
1328fn 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
1336/// Truncate a message BODY to at most `max_chars`, keeping the TAIL. This
1337/// crate's design keeps "the last N" throughout — the last bytes of a
1338/// growing file, the last lines, the last messages — and per-message
1339/// character truncation used to be the one place that kept the *first* N
1340/// characters instead, cutting off a long tool-approval or error message
1341/// right before its actionable ending. The dropped prefix is called out
1342/// with an explicit `[truncated N chars]` marker instead of disappearing
1343/// silently.
1344fn 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
1354/// Collapse whitespace runs (including newlines) before a dedup
1355/// comparison. Codex writes the same human/agent turn twice under two
1356/// encodings — `response_item.message` (whose parts this parser joins with
1357/// an extra `"\n"`) and `event_msg.*` (a single flat string) — that differ
1358/// only by incidental whitespace, not content. Comparing on collapsed
1359/// whitespace merges exactly that pair while leaving two messages with
1360/// different words distinct.
1361fn 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        // D1: tool call arguments now surface in their own section.
1485        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        // D2: tool outputs now surface as their own message role.
1489        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                // response_item joins parts with "\n"; the trailing "\n" on
1532                // the first part plus the join's own "\n" produces a
1533                // double newline, unlike the flat event_msg string below.
1534                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}