supercode-interchange 0.5.46

Canonical, provider-neutral session interchange primitives for Volter Harness
Documentation
//! Conservative incremental normalization. A streaming comparison verifies every old
//! byte: growing rewrites must not be mistaken for appends. Tool/branch changes
//! fall back to the complete canonical loader, never an approximate projection.
use super::*;
use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};

pub(crate) struct ClaudeAppendState {
    graph: ClaudeReplayIndex,
    replay: Vec<usize>,
    bytes: u64,
    last_assistant: Option<String>,
}

impl ClaudeAppendState {
    pub(crate) fn new(session: &Session, fidelity: Fidelity) -> Result<Option<Self>> {
        if session.meta.source != SessionSource::ClaudeCode
            || !session.raw_is_verbatim
            || !session.raw_trailing_newline
            || session.parse_error_lines != 0
            || !session.subagents.is_empty()
        {
            return Ok(None);
        }
        let mut graph = ClaudeReplayIndex::default();
        let mut bytes = 0;
        for (i, line) in session.raw.iter().enumerate() {
            bytes += line.len() as u64 + 1;
            if !line.trim().is_empty() {
                graph.observe(i, &serde_json::from_str(line)?)?;
            }
        }
        let selection = graph.clone().select_lines(fidelity)?;
        if selection.residue != session.load_residue {
            // E.g. a semantic loader omitted an invalid child. Reuse its full
            // owner rather than dropping non-graph diagnostics on an append.
            return Ok(None);
        }
        let replay = selection.lines;
        let last_assistant = replay.last().and_then(|i| {
            let value: Value = serde_json::from_str(&session.raw[*i]).ok()?;
            (value.get("type").and_then(Value::as_str) == Some("assistant"))
                .then(|| claude_assistant_message_id(&value).map(str::to_owned))
                .flatten()
        });
        Ok(Some(Self {
            graph,
            replay,
            bytes,
            last_assistant,
        }))
    }

    pub(crate) fn append(
        &self,
        path: &Path,
        current: &Session,
        fidelity: Fidelity,
    ) -> Result<Option<(Session, Self)>> {
        let mut file = std::fs::File::open(path)?;
        let before = file.metadata()?;
        if before.len() <= self.bytes || before.len() - self.bytes > 8 * 1024 * 1024 {
            return Ok(None);
        }
        if !prefix_matches(&mut file, &current.raw)? {
            return Ok(None);
        }
        file.seek(SeekFrom::Start(self.bytes))?;
        let mut appended = String::new();
        (&mut file)
            .take(before.len() - self.bytes)
            .read_to_string(&mut appended)?;
        if !appended.ends_with('\n') {
            return Ok(None);
        }
        let lines = appended
            .strip_suffix('\n')
            .unwrap()
            .split('\n')
            .collect::<Vec<_>>();
        let mut values = Vec::with_capacity(lines.len());
        let mut graph = self.graph.clone();
        let mut meta = current.meta.clone();
        for (i, line) in lines.iter().enumerate() {
            if line.trim().is_empty() {
                values.push(Value::Null);
                continue;
            }
            let Ok(value) = serde_json::from_str::<Value>(line) else {
                return Ok(None);
            };
            // These shapes cannot reorder old tool results, restore a foreign
            // envelope or alter context. Everything else uses the full codec.
            if !simple_message(&value) {
                return Ok(None);
            }
            graph.observe(current.raw.len() + i, &value)?;
            capture_claude_meta(&value, &mut meta, line)?;
            values.push(value);
        }
        let selection = graph.clone().select_lines(fidelity)?;
        if !selection.lines.starts_with(&self.replay) {
            return Ok(None);
        }
        let added = &selection.lines[self.replay.len()..];
        if added.iter().any(|i| *i < current.raw.len()) {
            // A newly selected old record is not an appended payload.
            return Ok(None);
        }
        if let Some(&first) = added.first() {
            let value = &values[first - current.raw.len()];
            if self
                .last_assistant
                .as_deref()
                .is_some_and(|id| claude_assistant_message_id(value) == Some(id))
            {
                // A new chunk can change an already-emitted assistant message.
                return Ok(None);
            }
        }
        let mut messages = Vec::new();
        let mut pending = None;
        for &i in added {
            let value = &values[i - current.raw.len()];
            if value.get("type").and_then(Value::as_str) == Some("assistant") {
                if value.get("isApiErrorMessage").and_then(Value::as_bool) == Some(true) {
                    flush_claude_assistant(&mut pending, &mut messages);
                    continue;
                }
                if let Some(previous) = pending.as_mut() {
                    if claude_assistant_message_id(previous)
                        .is_some_and(|id| claude_assistant_message_id(value) == Some(id))
                    {
                        merge_claude_assistant_chunk(previous, value);
                        continue;
                    }
                    flush_claude_assistant(&mut pending, &mut messages);
                }
                pending = Some(value.clone());
            } else {
                flush_claude_assistant(&mut pending, &mut messages);
                let before = messages.len();
                push_claude_user(value, &mut messages);
                capture_claude_record_provenance(value, &mut messages[before..]);
                restore_single_grok_message(value, &mut messages[before..]);
            }
        }
        flush_claude_assistant(&mut pending, &mut messages);
        let last_assistant = added
            .last()
            .map(|i| &values[*i - current.raw.len()])
            .map(|value| {
                if value.get("type").and_then(Value::as_str) == Some("assistant") {
                    claude_assistant_message_id(value).map(str::to_owned)
                } else {
                    None
                }
            })
            .unwrap_or_else(|| self.last_assistant.clone());
        let after = file.metadata()?;
        if after.len() != before.len() || after.modified().ok() != before.modified().ok() {
            return Ok(None);
        }
        let count = current
            .imported_message_count
            .unwrap_or(current.messages.len())
            + messages.len();
        let mut session = current.clone();
        session.meta = meta;
        session.messages.extend(messages);
        session
            .raw
            .extend(lines.iter().map(|line| (*line).to_owned()));
        session.imported_message_count = Some(count);
        session.load_residue = selection.residue;
        Ok(Some((
            session,
            Self {
                graph,
                replay: selection.lines,
                bytes: before.len(),
                last_assistant,
            },
        )))
    }
}

fn prefix_matches(file: &mut std::fs::File, raw: &[String]) -> Result<bool> {
    // The full-fidelity follower already owns these verbatim records. Compare
    // directly, without another whole-file allocation or a probabilistic hash.
    // Seeking back to `bytes` after this helper handles BufReader's lookahead.
    let mut reader = BufReader::with_capacity(64 * 1024, file);
    for line in raw {
        for mut expected in [line.as_bytes(), b"\n".as_slice()] {
            while !expected.is_empty() {
                let available = reader.fill_buf()?;
                let length = available.len().min(expected.len());
                if length == 0 || available[..length] != expected[..length] {
                    return Ok(false);
                }
                reader.consume(length);
                expected = &expected[length..];
            }
        }
    }
    Ok(true)
}

fn simple_message(value: &Value) -> bool {
    if value
        .as_object()
        .is_some_and(|record| record.keys().any(|key| key.starts_with("_supercode_")))
    {
        return false;
    }
    let kind = value.get("type").and_then(Value::as_str);
    if !matches!(kind, Some("user" | "assistant")) {
        return false;
    }
    match value.get("message").and_then(|m| m.get("content")) {
        Some(Value::String(_)) => true,
        Some(Value::Array(blocks)) => blocks.iter().all(|block| {
            matches!(
                block.get("type").and_then(Value::as_str),
                Some("text" | "image" | "thinking" | "redacted_thinking" | "fallback")
            )
        }),
        _ => false,
    }
}