Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/app/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,7 @@ fn start_local(store: SessionStore) -> std::io::Result<(SpawnedHost, HostMux)> {
let mut host_services = HostServices {
background_startup_probes: true,
ai_title_generation: true,
drop_superseded_diffs: true,
..HostServices::default()
};
match mcp_host::Host::bind() {
Expand Down
45 changes: 28 additions & 17 deletions crates/core/src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,10 @@ use serde::{Deserialize, Serialize};

use crate::git::merge_file_changes_by_path;

mod superseded;

pub use superseded::{TurnSnapshots, drop_turn_diffs};

/// Claude Code's own prompt for resuming work after a usage-window reset.
pub const RESUME_PROMPT: &str = "Continue from where you left off.";

Expand Down Expand Up @@ -169,7 +173,7 @@ impl From<AgentEvent> for StoredEvent {
}

/// One renderable row in the chat timeline.
#[derive(Debug, Clone)]
#[derive(Debug, Clone, PartialEq)]
pub struct TimelineEntry {
/// Provider item id (or a synthetic id for errors).
pub id: String,
Expand All @@ -181,7 +185,7 @@ pub struct TimelineEntry {
}

/// Per-turn ("Work Log" section) metadata folded from turn lifecycle events.
#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone, Default, PartialEq)]
pub struct TurnMeta {
/// Provider-native id for this turn, used to attach replacement turn-diff
/// snapshots without relying on ambient "current turn" state during replay.
Expand Down Expand Up @@ -243,7 +247,7 @@ impl TurnTiming {
/// turn, intersected with the turn's own `TurnStarted`..`TurnCompleted` bounds.
/// Only the current turn can have items in flight, so one accumulator is
/// enough; it resets whenever a turn opens.
#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone, Default, PartialEq)]
struct ToolClock {
/// Item ids currently known to be in progress.
open: HashSet<String>,
Expand Down Expand Up @@ -408,13 +412,13 @@ fn tool_is_active(lifecycle: ToolLifecycle, state: ToolState) -> bool {
}
}

#[derive(Debug, Clone)]
#[derive(Debug, Clone, PartialEq)]
pub struct TurnChangeSet {
pub changes: Vec<FileChange>,
pub completeness: ChangeCompleteness,
}

#[derive(Debug, Clone)]
#[derive(Debug, Clone, PartialEq)]
pub enum EntryContent {
Item(ItemContent),
/// A user message injected into an already-open turn. Provider-originated
Expand Down Expand Up @@ -472,7 +476,7 @@ pub enum SteeringStatus {
/// A proposed plan captured this session (Codex plan item / Claude
/// `ExitPlanMode`). Streaming deltas accumulate into `markdown`; a `ProposedPlan`
/// event replaces it with the final text.
#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone, Default, PartialEq)]
pub struct ProposedPlan {
pub item_id: String,
pub markdown: String,
Expand Down Expand Up @@ -502,8 +506,10 @@ pub struct RunningTurn {
pub started_at: Option<u64>,
}

/// Folded view of a session's event history.
#[derive(Debug, Clone, Default)]
/// Folded view of a session's event history. Two timelines are equal only
/// when every later event folds the same onto both, so equality covers the
/// private state the fold continues from, not just what renders.
#[derive(Debug, Clone, Default, PartialEq)]
pub struct Timeline {
/// Top-level entries are shared so virtualized UI snapshots can retain a
/// turn without cloning its potentially large message, command-output, and
Expand Down Expand Up @@ -763,10 +769,7 @@ impl Timeline {
completeness,
} => {
let turn = self
.turns
.iter()
.position(|turn| turn.provider_turn_id.as_deref() == Some(turn_id.as_str()))
.or(self.current_turn)
.provider_turn(turn_id)
.unwrap_or_else(|| self.push_turn(ts));
if !turn_id.is_empty() {
self.turns[turn].provider_turn_id = Some(turn_id.clone());
Expand All @@ -781,10 +784,7 @@ impl Timeline {
checkpoint_id,
} => {
let turn = self
.turns
.iter()
.position(|turn| turn.provider_turn_id.as_deref() == Some(turn_id.as_str()))
.or(self.current_turn)
.provider_turn(turn_id)
.unwrap_or_else(|| self.push_turn(ts));
self.turns[turn].provider_turn_id = Some(turn_id.clone());
self.turns[turn].provider_checkpoint_id = Some(checkpoint_id.clone());
Expand Down Expand Up @@ -924,7 +924,9 @@ impl Timeline {
}
self.usage = Some(usage);
}
AgentEvent::Warning { message } => log::warn!("provider warning: {message}"),
// Logged once where the live provider event arrives; a fold also
// replays stored records, which would log every old warning again.
AgentEvent::Warning { .. } => {}
AgentEvent::ProviderStartFailed { error } => {
let turn = self.ensure_turn(ts);
let id = self.synthetic_id("error", ts);
Expand Down Expand Up @@ -1078,6 +1080,15 @@ impl Timeline {
}
}

/// The turn a provider-addressed record lands on: the first turn carrying
/// `turn_id`, else the current one. `None` means the record opens a turn.
fn provider_turn(&self, turn_id: &str) -> Option<usize> {
self.turns
.iter()
.position(|turn| turn.provider_turn_id.as_deref() == Some(turn_id))
.or(self.current_turn)
}

/// Whether the current turn is still accumulating. A turn is finished once
/// a `TurnCompleted` has been folded, which records a status even when the
/// event carried no timestamp to store as `end_ts`; both must be checked or
Expand Down
251 changes: 251 additions & 0 deletions crates/core/src/session/superseded.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,251 @@
use agent::AgentEvent;

use super::Timeline;

/// Follows a fold to name the turn-changes snapshot each new one supersedes.
///
/// A snapshot replaces its turn's changes wholesale, and nothing the fold does
/// between two snapshots reads a turn's diffs, so once a later snapshot lands
/// on the same turn the earlier one's diffs never reach a timeline again. The
/// turn is the one the fold itself resolves the snapshot to, in the same
/// lifetime: a rewind ends a turn's lifetime even when its index is reused.
#[derive(Debug, Clone)]
pub struct TurnSnapshots<K> {
/// The latest snapshot on each turn the fold holds.
latest: Vec<Option<K>>,
}

impl<K> Default for TurnSnapshots<K> {
fn default() -> Self {
Self { latest: Vec::new() }
}
}

impl<K> TurnSnapshots<K> {
/// Fold the event recorded at `ts`, which the caller knows as `key`, into
/// `fold`. When it is a snapshot that supersedes an earlier one, the
/// earlier one's key is returned.
pub fn apply_at(
&mut self,
fold: &mut Timeline,
ts: Option<u64>,
event: &AgentEvent,
key: K,
) -> Option<K> {
let target = match event {
AgentEvent::TurnChangesUpdated { turn_id, .. } => Some(fold.provider_turn(turn_id)),
_ => None,
};
fold.apply_at(ts, event);
self.latest.truncate(fold.turns.len());
self.latest.resize_with(fold.turns.len(), || None);
let turn = target?.unwrap_or(fold.turns.len() - 1);
self.latest[turn].replace(key)
}
}

/// Drop the diffs of a turn-changes snapshot, keeping its paths, kinds and
/// completeness. Returns whether there was any diff to drop.
pub fn drop_turn_diffs(event: &mut AgentEvent) -> bool {
let AgentEvent::TurnChangesUpdated { changes, .. } = event else {
return false;
};
let mut dropped = false;
for change in changes {
dropped |= change.diff.take().is_some();
}
dropped
}

#[cfg(test)]
mod tests {
use agent::{
ChangeCompleteness, FileChange, FileChangeKind, ItemContent, RewindMode, ThreadItem,
TurnStatus,
};

use super::*;
use crate::session::StoredEvent;

fn at(ts: u64, event: AgentEvent) -> StoredEvent {
StoredEvent {
ts: Some(ts),
event,
elided: None,
}
}

fn started(turn: &str) -> AgentEvent {
AgentEvent::TurnStarted {
turn_id: turn.into(),
}
}

fn completed(turn: &str) -> AgentEvent {
AgentEvent::TurnCompleted {
turn_id: turn.into(),
status: TurnStatus::Completed,
usage: None,
}
}

fn checkpoint(turn: &str, checkpoint: &str) -> AgentEvent {
AgentEvent::TurnCheckpoint {
turn_id: turn.into(),
checkpoint_id: checkpoint.into(),
}
}

fn rewind(checkpoint: &str) -> AgentEvent {
AgentEvent::RewindCompleted {
checkpoint_id: checkpoint.into(),
mode: RewindMode::Conversation,
prefill: None,
}
}

/// A snapshot of the turn's net changes: one file whose diff is `diff`.
fn changes(turn: &str, diff: &str) -> AgentEvent {
AgentEvent::TurnChangesUpdated {
turn_id: turn.into(),
changes: vec![FileChange {
path: "f".into(),
kind: FileChangeKind::Modify,
diff: Some(diff.into()),
}],
completeness: ChangeCompleteness::Exact,
}
}

fn message(id: &str, text: &str) -> AgentEvent {
AgentEvent::ItemCompleted(ThreadItem {
id: id.into(),
parent_item_id: None,
content: ItemContent::AssistantMessage { text: text.into() },
})
}

fn user(id: &str) -> AgentEvent {
AgentEvent::ItemCompleted(ThreadItem {
id: id.into(),
parent_item_id: None,
content: ItemContent::UserMessage {
text: id.into(),
context_len: None,
attachments: Vec::new(),
},
})
}

/// The indices of the snapshots a later one supersedes, ascending.
fn superseded(log: &[StoredEvent]) -> Vec<usize> {
let mut fold = Timeline::default();
let mut snapshots = TurnSnapshots::default();
let mut superseded: Vec<usize> = log
.iter()
.enumerate()
.filter_map(|(index, record)| {
snapshots.apply_at(&mut fold, record.ts, &record.event, index)
})
.collect();
superseded.sort_unstable();
superseded
}

#[test]
fn superseded_snapshots_follow_the_turn_the_fold_lands_them_on() {
let reused_id = vec![
at(1, started("x")),
at(2, changes("x", "a")),
at(3, message("m1", "one")),
at(4, completed("x")),
at(5, started("x")),
// Lands on the first turn named x, superseding record 1 there.
at(6, changes("x", "b")),
// Names no turn, so lands on the current one; superseded by record 8.
at(7, changes("", "c")),
at(8, message("m2", "two")),
at(9, changes("", "d")),
at(10, message("m3", "three")),
at(11, completed("x")),
];
let opens_turn = vec![
at(1, changes("x", "a")),
at(2, message("m1", "one")),
at(3, changes("x", "b")),
at(4, message("m2", "two")),
];
let names_turn = vec![
at(1, started("t")),
// No turn is named x yet: lands on the current turn and names it.
at(2, changes("x", "a")),
at(3, message("m1", "one")),
at(4, completed("t")),
at(5, started("u")),
// Finds the first turn by the name record 1 gave it.
at(6, changes("x", "b")),
at(7, message("m2", "two")),
at(8, completed("u")),
];
let rewound = vec![
at(1, user("u1")),
at(2, started("t1")),
at(3, checkpoint("t1", "c1")),
at(4, completed("t1")),
at(5, user("u2")),
at(6, started("t2")),
at(7, checkpoint("t2", "c2")),
at(8, changes("t2", "a")),
at(9, message("m1", "one")),
at(10, completed("t2")),
// Drops the second turn; the next user message reopens index 1.
at(11, rewind("c2")),
at(12, user("u3")),
at(13, started("t3")),
at(14, changes("t3", "b")),
at(15, message("m2", "two")),
at(16, changes("t3", "c")),
at(17, message("m3", "three")),
at(18, completed("t3")),
];
let late_changes = vec![
at(1, started("t")),
at(2, changes("t", "a")),
at(3, message("m", "done")),
at(4, completed("t")),
at(5, changes("t", "b")),
];
let cases = [
(
"a reused turn id lands on its first turn",
reused_id,
vec![1, 6],
),
("a snapshot that opened its turn", opens_turn, vec![0]),
("a snapshot that named its turn", names_turn, vec![1]),
(
"a rewind ends a turn's lifetime though its index is reused",
rewound,
vec![13],
),
("a snapshot after its turn completed", late_changes, vec![1]),
];
for (label, log, expected) in cases {
assert_eq!(superseded(&log), expected, "{label}");
let mut stubbed = log.clone();
for &index in &expected {
assert!(drop_turn_diffs(&mut stubbed[index].event), "{label}");
}
assert_eq!(
Timeline::fold_stored(&stubbed),
Timeline::fold_stored(&log),
"{label}: the fold is unchanged"
);
assert_eq!(
superseded(&stubbed),
expected,
"{label}: a stubbed log names the same snapshots"
);
}
}
}
1 change: 1 addition & 0 deletions crates/headless/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ fn serve_command(args: &[String]) -> Result<(), String> {
let mut services = HostServices {
background_startup_probes: true,
ai_title_generation: true,
drop_superseded_diffs: true,
..HostServices::default()
};
if let Ok(mut mcp_host) = mcp_host::Host::bind() {
Expand Down
Loading
Loading