fix: preserve live errors across segment rotation

This commit is contained in:
2026-08-19 08:29:36 +09:00
parent 9aeaa52bdb
commit 73902b03b6
3 changed files with 65 additions and 26 deletions
+40 -14
View File
@@ -271,6 +271,9 @@ pub struct App {
pub quit_confirm: Option<std::time::Instant>, pub quit_confirm: Option<std::time::Instant>,
/// Full display history in render order. /// Full display history in render order.
pub blocks: Vec<Block>, pub blocks: Vec<Block>,
/// Turn/protocol errors retained when a real `SegmentStart` replaces the
/// replayable conversation rows during segment rotation.
run_error_messages: Vec<String>,
pub scroll: Scroll, pub scroll: Scroll,
pub mode: Mode, pub mode: Mode,
pub cache: FileCache, pub cache: FileCache,
@@ -347,6 +350,7 @@ impl App {
quit: false, quit: false,
quit_confirm: None, quit_confirm: None,
blocks: Vec::new(), blocks: Vec::new(),
run_error_messages: Vec::new(),
scroll: Scroll::default(), scroll: Scroll::default(),
mode: Mode::Normal, mode: Mode::Normal,
cache: FileCache::new(), cache: FileCache::new(),
@@ -812,6 +816,15 @@ impl App {
}); });
} }
fn push_run_error(&mut self, message: String) {
self.run_error_messages.push(message.clone());
self.blocks.push(Block::Alert {
level: AlertLevel::Error,
source: AlertSource::Worker,
message,
});
}
fn handle_error(&mut self, code: ErrorCode, message: String) { fn handle_error(&mut self, code: ErrorCode, message: String) {
let text = format!("[{code:?}] {message}"); let text = format!("[{code:?}] {message}");
let was_applying = if let Some(picker) = self.rewind_picker.as_mut() { let was_applying = if let Some(picker) = self.rewind_picker.as_mut() {
@@ -829,7 +842,7 @@ impl App {
Duration::from_secs(6), Duration::from_secs(6),
); );
} }
self.push_error(text); self.push_run_error(text);
} }
fn rewind_submit_pending(&self) -> bool { fn rewind_submit_pending(&self) -> bool {
@@ -981,8 +994,16 @@ impl App {
self.assistant_streaming = false; self.assistant_streaming = false;
} }
Event::SegmentRotated { entry } => { Event::SegmentRotated { entry } => {
let retained_run_errors = self.run_error_messages.clone();
self.reset_for_rotation(); self.reset_for_rotation();
self.apply_log_entry_raw(&entry); self.apply_log_entry_raw(&entry);
for message in retained_run_errors {
self.blocks.push(Block::Alert {
level: AlertLevel::Error,
source: AlertSource::Worker,
message,
});
}
self.assistant_streaming = false; self.assistant_streaming = false;
} }
Event::SystemItem { item } => { Event::SystemItem { item } => {
@@ -2005,6 +2026,7 @@ impl App {
entries: &[serde_json::Value], entries: &[serde_json::Value],
greeting: Option<protocol::Greeting>, greeting: Option<protocol::Greeting>,
) { ) {
self.run_error_messages.clear();
self.turn_index = 0; self.turn_index = 0;
self.blocks.clear(); self.blocks.clear();
self.cache = FileCache::new(); self.cache = FileCache::new();
@@ -2082,11 +2104,7 @@ impl App {
self.apply_compaction_extension(&payload); self.apply_compaction_extension(&payload);
} }
session_store::LogEntry::RunErrored { message, .. } => { session_store::LogEntry::RunErrored { message, .. } => {
self.blocks.push(Block::Alert { self.push_run_error(message);
level: AlertLevel::Error,
source: AlertSource::Worker,
message,
});
} }
// Non-history-bearing variants don't affect the block view. // Non-history-bearing variants don't affect the block view.
_ => {} _ => {}
@@ -3269,7 +3287,7 @@ mod completion_flow_tests {
} }
#[test] #[test]
fn snapshot_and_segment_rotation_retain_one_durable_run_error_block() { fn snapshot_replaces_live_error_with_one_durable_run_error_block() {
let mut app = App::new("test".into()); let mut app = App::new("test".into());
app.handle_worker_event(Event::Error { app.handle_worker_event(Event::Error {
code: ErrorCode::ProviderError, code: ErrorCode::ProviderError,
@@ -3318,21 +3336,29 @@ mod completion_flow_tests {
}) })
.collect::<Vec<_>>(); .collect::<Vec<_>>();
assert_eq!(errors, ["provider unavailable"]); assert_eq!(errors, ["provider unavailable"]);
}
#[test]
fn segment_rotation_retains_live_error_across_real_segment_start() {
let mut app = App::new("test".into());
app.handle_worker_event(Event::Error { app.handle_worker_event(Event::Error {
code: ErrorCode::ProviderError, code: ErrorCode::ProviderError,
message: "retry unavailable".into(), message: "provider unavailable".into(),
}); });
let rotated_run_error = session_store::LogEntry::RunErrored { let segment_start = session_store::LogEntry::SegmentStart {
ts: 5, ts: 5,
interrupted: false, session_id: uuid::Uuid::nil(),
message: "retry unavailable".into(), system_prompt: None,
config: Default::default(),
history: Vec::new(),
forked_from: None,
compacted_from: None,
}; };
app.handle_worker_event(Event::SegmentRotated { app.handle_worker_event(Event::SegmentRotated {
entry: serde_json::to_value(rotated_run_error).unwrap(), entry: serde_json::to_value(segment_start).unwrap(),
}); });
let rotated_errors = app let errors = app
.blocks .blocks
.iter() .iter()
.filter_map(|block| match block { .filter_map(|block| match block {
@@ -3344,7 +3370,7 @@ mod completion_flow_tests {
_ => None, _ => None,
}) })
.collect::<Vec<_>>(); .collect::<Vec<_>>();
assert_eq!(rotated_errors, ["retry unavailable"]); assert_eq!(errors, ["[ProviderError] provider unavailable"]);
} }
#[test] #[test]
@@ -78,7 +78,7 @@ Deno.test("console routing projects live errors but not completion replies", ()
); );
}); });
Deno.test("snapshot and segment rotation retain one durable run_errored row", () => { Deno.test("snapshot replaces a live error with one durable run_errored row", () => {
const projector = createConsoleProjector(); const projector = createConsoleProjector();
let projection = projector.append([ let projection = projector.append([
{ {
@@ -115,13 +115,16 @@ Deno.test("snapshot and segment rotation retain one durable run_errored row", ()
assertEquals(errors[0].title, "Run error"); assertEquals(errors[0].title, "Run error");
assertEquals(errors[0].body, "provider unavailable"); assertEquals(errors[0].body, "provider unavailable");
assertEquals(errors[0].error, true); assertEquals(errors[0].error, true);
});
projection = projector.append([ Deno.test("segment rotation retains a live error beside the real SegmentStart history", () => {
const projector = createConsoleProjector();
const projection = projector.append([
{ {
eventId: "second-live-error", eventId: "live-error",
event: { event: {
event: "error", event: "error",
data: { code: "provider_error", message: "retry unavailable" }, data: { code: "provider_error", message: "provider unavailable" },
} satisfies Event, } satisfies Event,
}, },
{ {
@@ -130,21 +133,30 @@ Deno.test("snapshot and segment rotation retain one durable run_errored row", ()
event: "segment_rotated", event: "segment_rotated",
data: { data: {
entry: { entry: {
kind: "run_errored", kind: "segment_start",
ts: 5, ts: 5,
interrupted: false, session_id: "session-1",
message: "retry unavailable", system_prompt: null,
config: {},
history: [{
kind: "message",
role: "user",
content: [{ kind: "text", text: "retained conversation" }],
}],
}, },
}, },
} satisfies Event, } satisfies Event,
}, },
]); ]);
const rotatedErrors = projection.lines.filter((line) => const errors = projection.lines.filter((line) => line.kind === "error");
line.kind === "error" assertEquals(errors.length, 1);
assertEquals(errors[0].title, "error · provider_error");
assertEquals(errors[0].body, "provider unavailable");
assert(
projection.lines.some((line) => line.body === "retained conversation"),
"SegmentStart history should still seed the rotated projection",
); );
assertEquals(rotatedErrors.length, 1);
assertEquals(rotatedErrors[0].body, "retry unavailable");
}); });
Deno.test("workerConsoleHref encodes runtime and worker target authority", () => { Deno.test("workerConsoleHref encodes runtime and worker target authority", () => {
@@ -328,12 +328,13 @@ export function applyProtocolEvent(
next.status = event.data.status; next.status = event.data.status;
break; break;
case "segment_rotated": { case "segment_rotated": {
const retainedErrors = next.lines.filter((line) => line.kind === "error");
const segment = snapshotProjectionFromEntries( const segment = snapshotProjectionFromEntries(
envelope.eventId, envelope.eventId,
[event.data.entry], [event.data.entry],
next.cwd, next.cwd,
); );
next.lines = segment.lines; next.lines = [...segment.lines, ...retainedErrors];
next.tasks = segment.tasks; next.tasks = segment.tasks;
next.taskNextId = segment.taskNextId; next.taskNextId = segment.taskNextId;
break; break;