supercode-interchange 0.5.46

Canonical, provider-neutral session interchange primitives for Volter Harness
Documentation
//! Read-only projection using the canonical graph and record decoders. Payloads
//! are materialized one assistant group at a time, never for the whole history.
use super::*;
use std::fs::File;
use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};

const REFERENCE: &str = "supercode.internal.read_index";

#[derive(Clone)]
struct Record {
    offset: u64,
    length: usize,
    assistant: bool,
    error: bool,
    message_id: Option<String>,
}

/// Internal read-only index. Skeleton messages are never exposed as transcript
/// content; all returned payloads come back through the canonical record decoder.
#[doc(hidden)]
pub struct ClaudeReadIndex {
    file: File,
    stamp: (u64, Option<std::time::SystemTime>),
    records: Vec<Record>,
    groups: Vec<Vec<usize>>,
    skeleton: Vec<ChatMessage>,
    meta: SessionMeta,
    parse_errors: usize,
    residue: Vec<String>,
}

impl ClaudeReadIndex {
    /// Whether the first decodable source record identifies this native format.
    pub fn supports(path: &Path) -> Result<bool> {
        if !path.is_file() || looks_like_sqlite(path) {
            return Ok(false);
        }
        let mut reader = BufReader::new(File::open(path)?);
        let mut line = String::new();
        loop {
            line.clear();
            if reader.read_line(&mut line)? == 0 {
                return Ok(false);
            }
            if let Some(source) = detect_source(&line) {
                return Ok(source == SessionSource::ClaudeCode);
            }
        }
    }

    /// Index the exact active branch without retaining raw transcript payloads.
    pub fn open(path: &Path, fidelity: Fidelity) -> Result<Self> {
        let file = File::open(path)?;
        let metadata = file.metadata()?;
        let stamp = (metadata.len(), metadata.modified().ok());
        let mut reader = BufReader::new(file.try_clone()?.take(stamp.0));
        let mut meta = SessionMeta::new(SessionSource::ClaudeCode);
        let mut graph = ClaudeReplayIndex::default();
        let mut records = Vec::new();
        let mut residue_records = Vec::new();
        let mut parse_errors = 0;
        let mut offset = 0;
        let mut line = String::new();
        loop {
            line.clear();
            let length = reader.read_line(&mut line)?;
            if length == 0 {
                break;
            }
            let index = records.len();
            let mut record = Record {
                offset,
                length,
                assistant: false,
                error: false,
                message_id: None,
            };
            offset += length as u64;
            if !line.trim().is_empty() {
                match serde_json::from_str::<Value>(&line) {
                    Ok(value) => {
                        let raw = line.strip_suffix('\n').unwrap_or(&line);
                        capture_claude_meta(&value, &mut meta, raw)?;
                        graph.observe(index, &value)?;
                        record.assistant =
                            value.get("type").and_then(Value::as_str) == Some("assistant");
                        record.error =
                            value.get("isApiErrorMessage").and_then(Value::as_bool) == Some(true);
                        record.message_id = claude_assistant_message_id(&value).map(str::to_owned);
                        if claude_residue_kind(&value).is_some() {
                            residue_records.push(index);
                        }
                    }
                    Err(_) => parse_errors += 1,
                }
            }
            records.push(record);
        }
        let selection = graph.select_lines(fidelity)?;
        let mut result = Self {
            file,
            stamp,
            records,
            groups: Vec::new(),
            skeleton: Vec::new(),
            meta,
            parse_errors,
            residue: selection.residue,
        };
        for index in residue_records {
            let raw = result.raw(index)?;
            let raw = raw.strip_suffix('\n').unwrap_or(&raw);
            let value = serde_json::from_str(raw)?;
            capture_claude_residue(&mut result.meta, index, raw, &value);
        }
        let mut pending = Vec::new();
        let mut pending_id: Option<String> = None;
        for index in selection.lines {
            let record = &result.records[index];
            if record.assistant && !record.error {
                if !pending.is_empty() && !(pending_id.is_some() && pending_id == record.message_id)
                {
                    result.add_group(std::mem::take(&mut pending))?;
                }
                pending_id = result.records[index].message_id.clone();
                pending.push(index);
            } else {
                result.add_group(std::mem::take(&mut pending))?;
                pending_id = None;
                if !result.records[index].error || !result.records[index].assistant {
                    result.add_group(vec![index])?;
                }
            }
        }
        result.add_group(pending)?;
        reorder_tool_results_after_calls(&mut result.skeleton);
        ensure_tool_results_paired(&mut result.skeleton);
        result.validate()?;
        Ok(result)
    }

    fn raw(&mut self, index: usize) -> Result<String> {
        let record = &self.records[index];
        self.file.seek(SeekFrom::Start(record.offset))?;
        let mut bytes = vec![0; record.length];
        self.file.read_exact(&mut bytes)?;
        String::from_utf8(bytes).map_err(|error| crate::Error::Other(error.to_string()))
    }

    fn decode_group(&mut self, group: &[usize]) -> Result<Vec<ChatMessage>> {
        let mut out = Vec::new();
        let mut pending = None;
        for &index in group {
            let value: Value = serde_json::from_str(&self.raw(index)?)?;
            match value.get("type").and_then(Value::as_str) {
                Some("assistant") => {
                    if let Some(previous) = pending.as_mut() {
                        merge_claude_assistant_chunk(previous, &value);
                    } else {
                        pending = Some(value);
                    }
                }
                _ => {
                    let before = out.len();
                    match value.get("type").and_then(Value::as_str) {
                        Some("user") => push_claude_user(&value, &mut out),
                        Some("attachment") => push_claude_attachment(&value, &mut out),
                        Some("system") => push_claude_system(&value, &mut out),
                        _ => {}
                    }
                    capture_claude_record_provenance(&value, &mut out[before..]);
                    restore_single_grok_message(&value, &mut out[before..]);
                }
            }
        }
        flush_claude_assistant(&mut pending, &mut out);
        Ok(out)
    }

    fn add_group(&mut self, group: Vec<usize>) -> Result<()> {
        if group.is_empty() {
            return Ok(());
        }
        let group_id = self.groups.len();
        let decoded = self.decode_group(&group)?;
        self.groups.push(group);
        for (ordinal, mut message) in decoded.into_iter().enumerate() {
            // Keep precisely what ordering, pairing and item counting inspect.
            // Actual payloads and provenance are restored by read_messages.
            message.content = message
                .content
                .as_ref()
                .map(|text| if text.trim().is_empty() { "" } else { "x" }.into());
            message.content_parts = message.content_parts.as_ref().map(|parts| {
                if parts.is_empty() {
                    vec![]
                } else {
                    vec![Value::Null]
                }
            });
            for call in message.tool_calls.iter_mut().flatten() {
                call.function.arguments = String::new();
            }
            message.metadata.clear();
            message
                .metadata
                .insert(REFERENCE.into(), format!("{group_id}:{ordinal}"));
            self.skeleton.push(message);
        }
        Ok(())
    }

    /// Number of exactly normalized messages, after tool pairing and ordering.
    pub fn len(&self) -> usize {
        self.skeleton.len()
    }
    /// Whether the active branch has no normalized messages.
    pub fn is_empty(&self) -> bool {
        self.skeleton.is_empty()
    }
    /// Raw record count, including blank records, without retaining their bytes.
    pub fn raw_record_count(&self) -> usize {
        self.records.len()
    }

    /// Count the conversation/tool items used by the existing window receipt.
    pub fn item_count(&self, range: std::ops::Range<usize>) -> usize {
        self.skeleton[range]
            .iter()
            .map(|message| {
                usize::from(
                    matches!(message.role, Role::User | Role::Assistant | Role::Tool)
                        && has_content(message),
                ) + message.tool_calls().len()
            })
            .sum()
    }

    /// Materialize just these normalized positions through the original decoders.
    pub fn read_messages(&mut self, indices: impl IntoIterator<Item = usize>) -> Result<Session> {
        let mut messages = Vec::new();
        let mut cache: Option<(usize, Vec<ChatMessage>)> = None;
        for index in indices {
            let message = &self.skeleton[index];
            if let Some(reference) = message.metadata.get(REFERENCE) {
                let (group, ordinal) = reference.split_once(':').expect("owned reference");
                let group: usize = group.parse().expect("owned group");
                let ordinal: usize = ordinal.parse().expect("owned ordinal");
                if cache.as_ref().map(|entry| entry.0) != Some(group) {
                    cache = Some((group, self.decode_group(&self.groups[group].clone())?));
                }
                messages.push(cache.as_ref().unwrap().1[ordinal].clone());
            } else {
                // Canonical interrupted-tool placeholder: no source payload exists.
                messages.push(message.clone());
            }
        }
        self.validate()?;
        Ok(Session {
            meta: self.meta.clone(),
            messages,
            subagents: Vec::new(),
            raw: Vec::new(),
            raw_trailing_newline: false,
            raw_is_verbatim: false,
            imported_message_count: Some(self.len()),
            parse_error_lines: self.parse_errors,
            load_residue: self.residue.clone(),
        })
    }

    /// Only the messages the existing first/latest/end-of-turn summary examines.
    pub fn read_summary(&mut self) -> Result<Session> {
        let mut indices = std::collections::BTreeSet::new();
        for assistant_only in [false, true] {
            let matches = self.skeleton.iter().enumerate().filter(|(_, message)| {
                (message.role == Role::Assistant || (!assistant_only && message.role == Role::User))
                    && has_content(message)
            });
            if let Some((i, _)) = matches.clone().next() {
                indices.insert(i);
            }
            if let Some((i, _)) = matches.last() {
                indices.insert(i);
            }
        }
        if let Some(i) = self.skeleton.iter().rposition(|m| m.role != Role::System) {
            indices.insert(i);
        }
        self.read_messages(indices)
    }

    fn validate(&self) -> Result<()> {
        let metadata = self.file.metadata()?;
        if (metadata.len(), metadata.modified().ok()) != self.stamp {
            return Err(crate::Error::Other(
                "Claude transcript changed during window read; retry the request".into(),
            ));
        }
        Ok(())
    }
}

fn has_content(message: &ChatMessage) -> bool {
    message
        .content
        .as_deref()
        .is_some_and(|text| !text.trim().is_empty())
        || message
            .content_parts
            .as_ref()
            .is_some_and(|parts| !parts.is_empty())
}