From 1d06d4b08a36b43528bb77a4fa2823dfb7eb7dfa Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 16 Sep 2026 04:19:02 +0900 Subject: [PATCH] fix: keep compaction lifecycle runtime-only --- crates/protocol/src/lib.rs | 65 ++-- crates/protocol/src/typescript.rs | 26 +- crates/session-store/src/worker_metadata.rs | 18 + crates/tui/src/app.rs | 44 +-- crates/worker/src/in_flight.rs | 105 ++---- crates/worker/src/worker.rs | 309 ++++++------------ crates/worker/tests/compact_events_test.rs | 214 ++++++------ web/workspace/src/lib/generated/protocol.ts | 8 +- .../src/lib/workspace/console/model.test.ts | 11 +- .../src/lib/workspace/console/model.ts | 29 +- 10 files changed, 352 insertions(+), 477 deletions(-) diff --git a/crates/protocol/src/lib.rs b/crates/protocol/src/lib.rs index 8ca42409..f90490ce 100644 --- a/crates/protocol/src/lib.rs +++ b/crates/protocol/src/lib.rs @@ -1228,8 +1228,12 @@ pub enum Event { /// This is not part of LLM history or prompt context; clients may display it /// briefly as operational status. MemoryWorker(MemoryWorkerEvent), - /// Worker has started compacting the current session, or bound the run to its - /// observable Internal Worker. Revisions upsert one stable lifecycle item. + /// Runtime-only compaction progress. `None` clears the current status. + /// This never enters Session history and carries no operation or Segment identity. + CompactionProgress { + compaction: Option, + }, + /// Legacy compaction lifecycle event retained for wire read compatibility. CompactStart { lifecycle: CompactionLifecycle, }, @@ -1416,26 +1420,30 @@ pub enum CommandEvent { }, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(rename_all = "snake_case")] +pub enum CompactionPhase { + Preparing, + Summarizing, + Committing, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(rename_all = "snake_case")] +pub enum CompactionTrigger { + Manual, + PreRun, + RequestThreshold, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct InFlightCompaction { - pub schema_version: u32, - pub compaction_id: String, - pub revision: u64, - pub internal_worker: Option, + pub phase: CompactionPhase, pub started_at_ms: u64, -} - -impl InFlightCompaction { - pub fn from_running(lifecycle: &CompactionLifecycle) -> Option { - (lifecycle.state == CompactionLifecycleState::Running).then(|| Self { - schema_version: lifecycle.schema_version, - compaction_id: lifecycle.compaction_id.clone(), - revision: lifecycle.revision, - internal_worker: lifecycle.internal_worker.clone(), - started_at_ms: lifecycle.started_at_ms, - }) - } + pub trigger: CompactionTrigger, } /// Unfinished model output and active command state included in @@ -2309,11 +2317,9 @@ mod tests { exit_code: None, }], compaction: Some(InFlightCompaction { - schema_version: 3, - compaction_id: "compaction-1".into(), - revision: 1, - internal_worker: None, + phase: CompactionPhase::Summarizing, started_at_ms: 99, + trigger: CompactionTrigger::Manual, }), }, internal_workers: Vec::new(), @@ -2327,14 +2333,20 @@ mod tests { "streaming_args" ); assert_eq!( - parsed["data"]["in_flight"]["compaction"]["compaction_id"], - "compaction-1" + parsed["data"]["in_flight"]["compaction"]["phase"], + "summarizing" + ); + assert_eq!( + parsed["data"]["in_flight"]["compaction"]["trigger"], + "manual" ); assert!( parsed["data"]["in_flight"]["compaction"] .as_object() .is_some_and(|value| { !value.contains_key("state") + && !value.contains_key("compaction_id") + && !value.contains_key("internal_worker") && !value.contains_key("summary") && !value.contains_key("new_segment_id") }), @@ -2344,7 +2356,10 @@ mod tests { match serde_json::from_str::(&json).unwrap() { Event::Snapshot { in_flight, .. } => { assert_eq!(in_flight.blocks.len(), 3); - assert_eq!(in_flight.compaction.unwrap().compaction_id, "compaction-1"); + assert_eq!( + in_flight.compaction.unwrap().phase, + CompactionPhase::Summarizing + ); } other => panic!("expected Snapshot, got {other:?}"), } diff --git a/crates/protocol/src/typescript.rs b/crates/protocol/src/typescript.rs index b5fec423..22527d63 100644 --- a/crates/protocol/src/typescript.rs +++ b/crates/protocol/src/typescript.rs @@ -4,18 +4,18 @@ use ts_rs::{Config, TS}; use crate::{ Alert, AlertLevel, AlertSource, CommandEvent, CommandSnapshot, CommandStatus, CommandStream, - CommandStreamSlice, CompactionLifecycle, CompactionLifecycleState, CompletionEntry, - CompletionKind, ErrorCode, Event, Greeting, InFlightBlock, InFlightCompaction, - InFlightSnapshot, InFlightToolCallState, InternalWorkerKind, InternalWorkerRef, - InternalWorkerSnapshot, InvokeKind, MemoryWorkerEvent, Method, PasteArtifactAvailability, - PasteArtifactMediaType, PasteArtifactRef, PendingSubmissionSummary, PendingSubmissionsSnapshot, - Permission, RewindSummary, RewindTarget, RewindTargetId, RunResult, ScopeRule, Segment, - SessionContentPart, SessionEntryProvenance, SessionMessageRole, SessionSnapshot, - SessionSnapshotEntry, SessionSnapshotEntryData, SessionToolAttachment, SubmissionDisposition, - SymlinkPolicy, ToolResultDisposition, TurnResult, UploadedFileAvailability, UploadedFileRef, - WorkerBusyState, WorkerCommandAcknowledgement, WorkerCommandDisposition, WorkerCommandEnvelope, - WorkerCommandKind, WorkerEvent, WorkerMaintenanceState, WorkerRunState, WorkerState, - WorkerStateSnapshot, WorkerStatus, + CommandStreamSlice, CompactionLifecycle, CompactionLifecycleState, CompactionPhase, + CompactionTrigger, CompletionEntry, CompletionKind, ErrorCode, Event, Greeting, InFlightBlock, + InFlightCompaction, InFlightSnapshot, InFlightToolCallState, InternalWorkerKind, + InternalWorkerRef, InternalWorkerSnapshot, InvokeKind, MemoryWorkerEvent, Method, + PasteArtifactAvailability, PasteArtifactMediaType, PasteArtifactRef, PendingSubmissionSummary, + PendingSubmissionsSnapshot, Permission, RewindSummary, RewindTarget, RewindTargetId, RunResult, + ScopeRule, Segment, SessionContentPart, SessionEntryProvenance, SessionMessageRole, + SessionSnapshot, SessionSnapshotEntry, SessionSnapshotEntryData, SessionToolAttachment, + SubmissionDisposition, SymlinkPolicy, ToolResultDisposition, TurnResult, + UploadedFileAvailability, UploadedFileRef, WorkerBusyState, WorkerCommandAcknowledgement, + WorkerCommandDisposition, WorkerCommandEnvelope, WorkerCommandKind, WorkerEvent, + WorkerMaintenanceState, WorkerRunState, WorkerState, WorkerStateSnapshot, WorkerStatus, subscription::{ EventSubscriptionSelector, SubscriptionEvent, SubscriptionEventPayload, SubscriptionFrame, SubscriptionFramePayload, SubscriptionId, SubscriptionRejectionCode, SubscriptionRequest, @@ -73,6 +73,8 @@ pub fn generated_protocol_types() -> String { push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); diff --git a/crates/session-store/src/worker_metadata.rs b/crates/session-store/src/worker_metadata.rs index 387e5f98..3549dbb9 100644 --- a/crates/session-store/src/worker_metadata.rs +++ b/crates/session-store/src/worker_metadata.rs @@ -182,6 +182,24 @@ pub trait WorkerMetadataStore: Send + Sync { Ok(metadata) } + /// Compare and swap the active Segment pointer while preserving unrelated metadata. + /// Returns `false` without mutation when the durable pointer no longer matches. + fn compare_and_swap_active( + &self, + worker_name: &str, + expected: &WorkerActiveSegmentRef, + replacement: WorkerActiveSegmentRef, + ) -> Result { + let mut matched = false; + self.update_by_name(worker_name, |metadata| { + if metadata.active.as_ref() == Some(expected) { + metadata.active = Some(replacement); + matched = true; + } + })?; + Ok(matched) + } + /// Set the active pointer while preserving spawned children, workspace ownership, and manifest snapshot. fn set_active( &self, diff --git a/crates/tui/src/app.rs b/crates/tui/src/app.rs index b0483bf6..13bf8ee9 100644 --- a/crates/tui/src/app.rs +++ b/crates/tui/src/app.rs @@ -1400,6 +1400,19 @@ impl App { self.reset_run_state(); } } + Event::CompactionProgress { compaction } => { + if compaction.is_some() { + if self.last_streaming_compact_mut().is_none() { + self.blocks.push(Block::Compact(CompactEvent::Streaming { + started_at: Instant::now(), + })); + } + } else if let Some(Block::Compact(CompactEvent::Streaming { .. })) = + self.blocks.last() + { + self.blocks.pop(); + } + } Event::CompactStart { lifecycle } => { let should_apply = match &self.active_compaction { None => true, @@ -1688,9 +1701,7 @@ impl App { } } } - self.active_compaction = compaction - .as_ref() - .map(|lifecycle| (lifecycle.compaction_id.clone(), lifecycle.revision)); + self.active_compaction = None; if compaction.is_some() && self.last_streaming_compact_mut().is_none() { self.blocks.push(Block::Compact(CompactEvent::Streaming { started_at: Instant::now(), @@ -4301,28 +4312,21 @@ mod completion_flow_tests { } #[test] - fn snapshot_restores_running_compaction_and_fences_unrelated_terminal() { + fn snapshot_restores_and_runtime_clear_removes_compaction_progress() { let mut app = App::new("test".into()); - let lifecycle = test_compaction_lifecycle(protocol::CompactionLifecycleState::Running); app.apply_in_flight_snapshot(InFlightSnapshot { - compaction: Some(protocol::InFlightCompaction::from_running(&lifecycle).unwrap()), + compaction: Some(protocol::InFlightCompaction { + phase: protocol::CompactionPhase::Summarizing, + started_at_ms: 100, + trigger: protocol::CompactionTrigger::Manual, + }), ..InFlightSnapshot::default() }); - - let mut unrelated = lifecycle; - unrelated.compaction_id = "another-compaction".into(); - unrelated.revision = 2; - unrelated.state = protocol::CompactionLifecycleState::Failed; - unrelated.error = Some("must not replace".into()); - app.handle_worker_event(Event::CompactFailed { - lifecycle: unrelated, - }); - assert_eq!(compact_block_count(&app), 1); - assert!(matches!( - app.blocks.as_slice(), - [Block::Compact(CompactEvent::Streaming { .. })] - )); + + app.handle_worker_event(Event::CompactionProgress { compaction: None }); + + assert_eq!(compact_block_count(&app), 0); } #[test] diff --git a/crates/worker/src/in_flight.rs b/crates/worker/src/in_flight.rs index 539bead1..a11fcf3e 100644 --- a/crates/worker/src/in_flight.rs +++ b/crates/worker/src/in_flight.rs @@ -221,33 +221,13 @@ impl InFlightEvents { self.lock().commands = commands; } - /// Publish current compaction progress into reconnect snapshots. - /// Terminal lifecycle events clear the in-flight value; their durable - /// session record remains the historical authority. - pub(crate) fn update_compaction(&self, lifecycle: &protocol::CompactionLifecycle) { + /// Atomically update reconnect state and publish the matching live progress event. + pub(crate) fn set_compaction(&self, compaction: Option) { let mut inner = self.lock(); - match lifecycle.state { - protocol::CompactionLifecycleState::Running => match &inner.compaction { - None => inner.compaction = protocol::InFlightCompaction::from_running(lifecycle), - Some(current) - if current.compaction_id == lifecycle.compaction_id - && lifecycle.revision > current.revision => - { - inner.compaction = protocol::InFlightCompaction::from_running(lifecycle); - } - Some(_) => {} - }, - protocol::CompactionLifecycleState::Done - | protocol::CompactionLifecycleState::Failed - | protocol::CompactionLifecycleState::Interrupted => { - if inner.compaction.as_ref().is_some_and(|current| { - current.compaction_id == lifecycle.compaction_id - && lifecycle.revision > current.revision - }) { - inner.compaction = None; - } - } - } + inner.compaction = compaction.clone(); + let _ = self + .working_event_tx + .send(Event::CompactionProgress { compaction }); } pub(crate) fn clear(&self) { @@ -773,79 +753,42 @@ mod tests { } #[test] - fn compaction_progress_is_snapshot_only_while_running() { + fn compaction_progress_updates_snapshot_and_live_event_atomically() { let (working_event_tx, _) = broadcast::channel(16); + let mut rx = working_event_tx.subscribe(); let in_flight = InFlightEvents::new(working_event_tx); - let running = protocol::CompactionLifecycle { - schema_version: 3, - compaction_id: "compact-1".into(), - revision: 1, - internal_worker: None, - state: protocol::CompactionLifecycleState::Running, + let progress = protocol::InFlightCompaction { + phase: protocol::CompactionPhase::Preparing, started_at_ms: 100, - ended_at_ms: None, - summary: None, - error: None, - new_segment_id: None, + trigger: protocol::CompactionTrigger::Manual, }; - in_flight.update_compaction(&running); + in_flight.set_compaction(Some(progress.clone())); let guard = in_flight.snapshot_guard(); assert_eq!( - snapshot_from_guard(&guard) - .compaction - .as_ref() - .map(|item| item.compaction_id.as_str()), - Some("compact-1") + snapshot_from_guard(&guard).compaction, + Some(progress.clone()) ); assert!(!snapshot_from_guard(&guard).is_empty()); drop(guard); - - let mut stale = running.clone(); - stale.revision = 0; - in_flight.update_compaction(&stale); - let mut other = running.clone(); - other.compaction_id = "compact-2".into(); - other.revision = 2; - in_flight.update_compaction(&other); - let guard = in_flight.snapshot_guard(); - assert_eq!( - snapshot_from_guard(&guard) - .compaction - .as_ref() - .map(|item| (item.compaction_id.as_str(), item.revision)), - Some(("compact-1", 1)) - ); - drop(guard); + assert!(matches!( + rx.try_recv().unwrap(), + Event::CompactionProgress { compaction: Some(item) } if item == progress + )); in_flight.clear(); let guard = in_flight.snapshot_guard(); - assert_eq!( - snapshot_from_guard(&guard) - .compaction - .as_ref() - .map(|item| item.compaction_id.as_str()), - Some("compact-1") - ); - drop(guard); - - let mut stale_done = running.clone(); - stale_done.state = protocol::CompactionLifecycleState::Done; - in_flight.update_compaction(&stale_done); - let guard = in_flight.snapshot_guard(); assert!(snapshot_from_guard(&guard).compaction.is_some()); drop(guard); - let mut done = running; - done.revision = 2; - done.state = protocol::CompactionLifecycleState::Done; - done.ended_at_ms = Some(200); - done.summary = Some("private summary".into()); - done.new_segment_id = Some("staged-segment".into()); - in_flight.update_compaction(&done); - + in_flight.set_compaction(None); let guard = in_flight.snapshot_guard(); assert!(snapshot_from_guard(&guard).compaction.is_none()); + drop(guard); + assert!(matches!( + rx.try_recv().unwrap(), + Event::CompactionProgress { compaction: None } + )); } #[test] diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 97147ac8..b46f922d 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -64,7 +64,6 @@ use crate::internal_worker::{ prepare_internal_worker_from_spec, }; -const COMPACTION_EXTENSION_DOMAIN: &str = "yoi.compaction"; const LARGE_PASTE_INLINE_MAX_BYTES: usize = 32 * 1024; const WORKER_ORCHESTRATION_INSTRUCTION_ID: &str = "worker.orchestration"; const WORKER_ORCHESTRATION_PROMPT_REF: &str = "common.worker_orchestration"; @@ -337,8 +336,9 @@ use crate::skill::{SkillActivationResponse, SkillClientError}; #[cfg(test)] use async_trait::async_trait; use protocol::{ - AlertLevel, AlertSource, CompactionLifecycle, CompactionLifecycleState, ErrorCode, Event, - RewindSummary, RewindTarget, RewindTargetId, Segment, + AlertLevel, AlertSource, CompactionLifecycle, CompactionLifecycleState, CompactionPhase, + CompactionTrigger, ErrorCode, Event, InFlightCompaction, RewindSummary, RewindTarget, + RewindTargetId, Segment, }; use tokio::net::UnixStream; use tokio::sync::broadcast; @@ -977,6 +977,11 @@ pub struct SegmentLocation { type WorkerMetadataWriter = Arc Result<(), WorkerStoreError> + Send + Sync>; +type WorkerMetadataSegmentCas = Arc< + dyn Fn(&str, &WorkerActiveSegmentRef, WorkerActiveSegmentRef) -> Result + + Send + + Sync, +>; fn worker_metadata_writer_for_store(store: &St) -> WorkerMetadataWriter where @@ -996,6 +1001,16 @@ where }) } +fn worker_metadata_segment_cas_for_store(store: &St) -> WorkerMetadataSegmentCas +where + St: WorkerMetadataStore + Clone + Send + Sync + 'static, +{ + let store = store.clone(); + Arc::new(move |worker_name, expected, replacement| { + store.compare_and_swap_active(worker_name, expected, replacement) + }) +} + /// Lock-free shared session/segment pointer. /// /// Holds the current `(SessionId, SegmentId)` pair and the append tally @@ -2099,6 +2114,7 @@ pub struct Worker { /// constructors install this from the same FsStore that owns the session /// logs; low-level `Worker::new` tests leave it absent. worker_metadata_writer: Option, + worker_metadata_segment_cas: Option, /// Shared session pointer. Source of truth for the Worker's current /// `segment_id` and append tally. `self.segment_id()` is a thin /// wrapper over `segment_state.segment_id()`. @@ -2428,6 +2444,7 @@ impl Worker { last_run_interrupted: false, store, worker_metadata_writer: None, + worker_metadata_segment_cas: None, segment_state: SegmentState::new(session_id, segment_id, 0), filesystem_authority, workdir_session, @@ -3059,6 +3076,27 @@ impl Worker { Ok(()) } + fn compare_and_swap_worker_metadata_segment( + &self, + expected: SegmentLocation, + replacement: SegmentLocation, + ) -> Result<(), WorkerError> { + let Some(compare_and_swap) = &self.worker_metadata_segment_cas else { + return Ok(()); + }; + let matched = compare_and_swap( + &self.manifest.worker.name, + &WorkerActiveSegmentRef::active_segment(expected.session_id, expected.segment_id), + WorkerActiveSegmentRef::active_segment(replacement.session_id, replacement.segment_id), + )?; + if !matched { + return Err(WorkerError::InvalidState( + "active Segment changed before compaction commit".into(), + )); + } + Ok(()) + } + /// Enable name-keyed Worker metadata write-through for Workers built through /// the low-level constructor. High-level manifest constructors enable it /// automatically; this hook lets tests and custom embedders opt into the @@ -3068,6 +3106,7 @@ impl Worker { St: WorkerMetadataStore + Clone + Send + Sync + 'static, { self.worker_metadata_writer = Some(worker_metadata_writer_for_store(&self.store)); + self.worker_metadata_segment_cas = Some(worker_metadata_segment_cas_for_store(&self.store)); self.write_worker_metadata_pending() } @@ -3330,15 +3369,6 @@ impl Worker { } } - /// Broadcast a typed `Event` to connected clients. No-op when no - /// `working_event_tx` is attached (tests / direct `Worker::new` usage) or when - /// no clients are currently subscribed. - fn send_event(&self, event: Event) { - if let Some(tx) = self.working_event_tx.as_ref() { - let _ = tx.send(event); - } - } - /// Push a `Method::Notify` entry onto the pending buffer. /// /// The notification will be appended to `worker.history` as an @@ -4546,44 +4576,12 @@ impl Worker { } } - fn persist_compaction_lifecycle( - &mut self, - lifecycle: &CompactionLifecycle, - ) -> Result<(), WorkerError> { - Ok(self.commit_entry(LogEntry::Extension { - ts: segment_log::now_millis(), - domain: COMPACTION_EXTENSION_DOMAIN.into(), - payload: serde_json::to_value(lifecycle).map_err(|error| { - WorkerError::InvalidState(format!( - "serialize compaction lifecycle {}: {error}", - lifecycle.compaction_id - )) - })?, - })?) - } - - fn persist_and_send_compact_start( - &mut self, - lifecycle: CompactionLifecycle, - ) -> Result<(), WorkerError> { - self.persist_compaction_lifecycle(&lifecycle)?; + fn set_compaction_progress(&self, compaction: Option) { if let Some(in_flight) = &self.in_flight { - in_flight.update_compaction(&lifecycle); + in_flight.set_compaction(compaction); + } else if let Some(tx) = &self.working_event_tx { + let _ = tx.send(Event::CompactionProgress { compaction }); } - self.send_event(Event::CompactStart { lifecycle }); - Ok(()) - } - - fn persist_and_send_compact_failed( - &mut self, - lifecycle: CompactionLifecycle, - ) -> Result<(), WorkerError> { - self.persist_compaction_lifecycle(&lifecycle)?; - if let Some(in_flight) = &self.in_flight { - in_flight.update_compaction(&lifecycle); - } - self.send_event(Event::CompactFailed { lifecycle }); - Ok(()) } /// Perform compaction after a `compact_needed` abort and resume execution. @@ -4614,7 +4612,10 @@ impl Worker { .map(|s| s.retained_tokens()) .unwrap_or(manifest::defaults::COMPACT_RETAINED_TOKENS); - match self.compact(retained).await { + match self + .compact_with_cancel(retained, None, CompactionTrigger::RequestThreshold) + .await + { Ok(new_segment_id) => { info!( new_segment_id = %new_segment_id, @@ -4659,7 +4660,10 @@ impl Worker { } let retained = state.retained_tokens(); - match self.compact(retained).await { + match self + .compact_with_cancel(retained, None, CompactionTrigger::PreRun) + .await + { Ok(new_segment_id) => { info!( new_segment_id = %new_segment_id, @@ -4721,79 +4725,9 @@ impl Worker { Ok(rewrite_guard) } - /// Terminalize and clean up any compaction that was left active by the - /// previous controller generation. This runs before the restored - /// controller publishes its first Idle state. + /// Legacy compaction extensions are ignored on restore. Current compaction + /// progress is runtime-only and cannot survive a process restart. pub async fn recover_unfinished_compaction(&mut self) -> Result<(), WorkerError> { - let (entries, _) = self.sink.subscribe_with_snapshot(); - let latest_payload = entries.iter().rev().find_map(|entry| match entry { - LogEntry::Extension { - domain, payload, .. - } if domain == COMPACTION_EXTENSION_DOMAIN => Some(payload.clone()), - _ => None, - }); - let Some(payload) = latest_payload else { - return Ok(()); - }; - #[derive(serde::Deserialize)] - #[serde(deny_unknown_fields)] - struct CompactionLifecycleWire { - schema_version: u32, - compaction_id: String, - revision: u64, - #[serde(default)] - internal_worker: Option, - state: CompactionLifecycleState, - started_at_ms: u64, - #[serde(default)] - ended_at_ms: Option, - #[serde(default)] - summary: Option, - #[serde(default)] - error: Option, - #[serde(default)] - new_segment_id: Option, - } - let wire: CompactionLifecycleWire = serde_json::from_value(payload).map_err(|error| { - WorkerError::InvalidState(format!("decode compaction lifecycle: {error}")) - })?; - if !matches!(wire.schema_version, 2 | 3) { - return Err(WorkerError::InvalidState(format!( - "unsupported compaction lifecycle schema version {}", - wire.schema_version - ))); - } - let mut lifecycle = CompactionLifecycle { - schema_version: wire.schema_version, - compaction_id: wire.compaction_id, - revision: wire.revision, - internal_worker: wire.internal_worker, - state: wire.state, - started_at_ms: wire.started_at_ms, - ended_at_ms: wire.ended_at_ms, - summary: wire.summary, - error: wire.error, - new_segment_id: wire.new_segment_id, - }; - match lifecycle.state { - CompactionLifecycleState::Running => { - lifecycle.schema_version = 3; - lifecycle.revision = lifecycle.revision.saturating_add(1); - lifecycle.state = CompactionLifecycleState::Interrupted; - lifecycle.ended_at_ms = Some(segment_log::now_millis()); - lifecycle.error = - Some("worker execution restarted before compaction completed".into()); - self.persist_compaction_lifecycle(&lifecycle)?; - self.send_event(Event::CompactFailed { - lifecycle: lifecycle.clone(), - }); - self.release_compaction_service(&lifecycle).await; - } - CompactionLifecycleState::Interrupted => { - self.release_compaction_service(&lifecycle).await; - } - CompactionLifecycleState::Done | CompactionLifecycleState::Failed => {} - } Ok(()) } @@ -4851,7 +4785,10 @@ impl Worker { return Ok(ManualCompactResult::Skipped { message }); } - match self.compact_with_cancel(retained, cancel.take()).await { + match self + .compact_with_cancel(retained, cancel.take(), CompactionTrigger::Manual) + .await + { Ok(new_segment_id) => { info!(new_segment_id = %new_segment_id, "Manual compaction succeeded"); if let Some(ref state) = state { @@ -5024,13 +4961,15 @@ impl Worker { /// Runs one parent-owned observable compaction service and returns the new /// Segment ID. Lifecycle revisions are committed before they are broadcast. pub async fn compact(&mut self, retained_tokens: u64) -> Result { - self.compact_with_cancel(retained_tokens, None).await + self.compact_with_cancel(retained_tokens, None, CompactionTrigger::RequestThreshold) + .await } async fn compact_with_cancel( &mut self, retained_tokens: u64, mut cancel: Option>, + trigger: CompactionTrigger, ) -> Result { let _rewrite_guard = self .prepare_session_rewrite(SessionRewriteKind::Compact) @@ -5047,7 +4986,12 @@ impl Worker { error: None, new_segment_id: None, }; - self.persist_and_send_compact_start(lifecycle.clone())?; + let started_at_ms = lifecycle.started_at_ms; + self.set_compaction_progress(Some(InFlightCompaction { + phase: CompactionPhase::Preparing, + started_at_ms, + trigger, + })); let outcome = if let Some(cancel) = cancel.as_mut() { tokio::select! { biased; @@ -5055,20 +4999,15 @@ impl Worker { let _ = changed; Err(WorkerError::CompactCancelled) } - result = self.compact_impl(retained_tokens, &mut lifecycle) => result, + result = self.compact_impl(retained_tokens, &mut lifecycle, trigger) => result, } } else { - self.compact_impl(retained_tokens, &mut lifecycle).await + self.compact_impl(retained_tokens, &mut lifecycle, trigger) + .await }; match outcome { Ok((new_segment_id, _summary)) => { debug_assert_eq!(lifecycle.state, CompactionLifecycleState::Done); - if let Some(in_flight) = &self.in_flight { - in_flight.update_compaction(&lifecycle); - } - self.send_event(Event::CompactDone { - lifecycle: lifecycle.clone(), - }); self.release_compaction_service(&lifecycle).await; Ok(new_segment_id) } @@ -5080,10 +5019,8 @@ impl Worker { CompactionLifecycleState::Failed }; lifecycle.ended_at_ms = Some(segment_log::now_millis()); - lifecycle.error = Some(error.to_string().chars().take(2_000).collect()); - let terminal = self.persist_and_send_compact_failed(lifecycle.clone()); + self.set_compaction_progress(None); self.release_compaction_service(&lifecycle).await; - terminal?; Err(error) } } @@ -5114,6 +5051,7 @@ impl Worker { &mut self, retained_tokens: u64, lifecycle: &mut CompactionLifecycle, + trigger: CompactionTrigger, ) -> Result<(SegmentId, String), WorkerError> { use crate::compact::worker::{ CompactWorkerContext, CompactWorkerInterceptor, CompactionOutputFeature, @@ -5344,7 +5282,11 @@ impl Worker { .map_err(|error| WorkerError::InvalidState(error.to_string()))?; lifecycle.revision = lifecycle.revision.saturating_add(1); lifecycle.internal_worker = Some(internal_ref); - self.persist_and_send_compact_start(lifecycle.clone())?; + self.set_compaction_progress(Some(InFlightCompaction { + phase: CompactionPhase::Summarizing, + started_at_ms: lifecycle.started_at_ms, + trigger, + })); if let Err(error) = handle.send(summary_input.text).await { let _ = registry.remove_service(&handle.session_id_string()); @@ -5650,24 +5592,11 @@ impl Worker { })?, }); } - // Commit the terminal lifecycle in the same atomic replacement-segment - // creation as the rewritten history. Restore can therefore never see a - // replacement segment without the Done fact for the compaction that - // created it. - lifecycle.revision = lifecycle.revision.saturating_add(1); - lifecycle.state = CompactionLifecycleState::Done; - lifecycle.ended_at_ms = Some(segment_log::now_millis()); - lifecycle.summary = Some(summary_text.clone()); - lifecycle.new_segment_id = Some(new_segment_id.to_string()); - initial_entries.push(LogEntry::Extension { - ts: segment_log::now_millis(), - domain: COMPACTION_EXTENSION_DOMAIN.to_string(), - payload: serde_json::to_value(&*lifecycle).map_err(|error| { - WorkerError::InvalidState(format!( - "serialize terminal compaction lifecycle: {error}" - )) - })?, - }); + self.set_compaction_progress(Some(InFlightCompaction { + phase: CompactionPhase::Committing, + started_at_ms: lifecycle.started_at_ms, + trigger, + })); self.store .create_segment(old_loc.session_id, new_segment_id, &initial_entries)?; @@ -5677,10 +5606,11 @@ impl Worker { // replacement may be collected later. The metadata store publishes the // complete record with an atomic replace, so restore observes either the // old Segment or this fully-written replacement, never a partial switch. - self.write_worker_metadata_active(SegmentLocation { + let new_location = SegmentLocation { session_id: old_loc.session_id, segment_id: new_segment_id, - })?; + }; + self.compare_and_swap_worker_metadata_segment(old_loc, new_location)?; // All live mutations after the durable commit are infallible and happen // before the replacement SegmentStart is broadcast. This keeps the @@ -5730,6 +5660,10 @@ impl Worker { // their blocks from `SegmentStart.history`. No per-item // broadcast is required. let _ = &compact_introduced_system_messages; + lifecycle.revision = lifecycle.revision.saturating_add(1); + lifecycle.state = CompactionLifecycleState::Done; + lifecycle.ended_at_ms = Some(segment_log::now_millis()); + self.set_compaction_progress(None); Ok((new_segment_id, summary_text)) } @@ -5853,6 +5787,7 @@ where apply_worker_manifest(&mut worker, &manifest.engine); worker.set_cache_key(Some(segment_id.to_string())); let worker_metadata_writer = Some(worker_metadata_writer_for_store(&store)); + let worker_metadata_segment_cas = Some(worker_metadata_segment_cas_for_store(&store)); let scope = SharedScope::new(common.scope); let workdir_session = workdir_session_from_authority(&common.filesystem_authority, &scope); @@ -5863,6 +5798,7 @@ where last_run_interrupted: false, store, worker_metadata_writer, + worker_metadata_segment_cas, segment_state: SegmentState::new(session_id, segment_id, 0), filesystem_authority: common.filesystem_authority, workdir_session, @@ -5947,6 +5883,7 @@ where last_run_interrupted: false, store, worker_metadata_writer: None, + worker_metadata_segment_cas: None, segment_state: SegmentState::new(session_id, segment_id, 0), filesystem_authority: common.filesystem_authority, workdir_session, @@ -6055,6 +5992,7 @@ where apply_worker_manifest(&mut worker, &manifest.engine); worker.set_cache_key(Some(segment_id.to_string())); let worker_metadata_writer = Some(worker_metadata_writer_for_store(&store)); + let worker_metadata_segment_cas = Some(worker_metadata_segment_cas_for_store(&store)); let scope = SharedScope::new(common.scope); let workdir_session = workdir_session_from_authority(&common.filesystem_authority, &scope); @@ -6065,6 +6003,7 @@ where last_run_interrupted: false, store, worker_metadata_writer, + worker_metadata_segment_cas, segment_state: SegmentState::new(session_id, segment_id, 0), filesystem_authority: common.filesystem_authority, workdir_session, @@ -6426,6 +6365,7 @@ where let task_feature = TaskFeature::from_history(&state.history); let worker_metadata_writer = Some(worker_metadata_writer_for_store(&store)); + let worker_metadata_segment_cas = Some(worker_metadata_segment_cas_for_store(&store)); let scope = SharedScope::new(common.scope); let workdir_session = workdir_session_from_authority(&common.filesystem_authority, &scope); @@ -6436,6 +6376,7 @@ where last_run_interrupted: state.last_run_interrupted, store, worker_metadata_writer, + worker_metadata_segment_cas, segment_state: SegmentState::new(session_id, segment_id, state.entries_count), filesystem_authority: common.filesystem_authority, workdir_session, @@ -10232,53 +10173,21 @@ mod build_summary_prompt_tests { } #[tokio::test] - async fn restore_terminalizes_running_compaction_before_idle_publication() { + async fn restore_does_not_synthesize_compaction_lifecycle_history() { let (_dir, mut worker) = rewind_test_worker().await; - let lifecycle = CompactionLifecycle { - schema_version: 3, - compaction_id: "compact-before-restart".into(), - revision: 1, - internal_worker: None, - state: CompactionLifecycleState::Running, - started_at_ms: segment_log::now_millis(), - ended_at_ms: None, - summary: None, - error: None, - new_segment_id: None, - }; - worker.persist_compaction_lifecycle(&lifecycle).unwrap(); + let before = worker.sink.subscribe_with_snapshot().0; worker.recover_unfinished_compaction().await.unwrap(); - let (entries, _) = worker.sink.subscribe_with_snapshot(); - let restored = entries.iter().rev().find_map(|entry| match entry { - LogEntry::Extension { - domain, payload, .. - } if domain == COMPACTION_EXTENSION_DOMAIN => { - serde_json::from_value::(payload.clone()).ok() - } - _ => None, - }); - let restored = restored.expect("terminal compaction lifecycle"); - assert_eq!(restored.state, CompactionLifecycleState::Interrupted); - assert_eq!(restored.revision, 2); - assert!( - restored - .error - .as_deref() - .is_some_and(|error| error.contains("restarted")) - ); - - let mut future = lifecycle; - future.schema_version = 4; - future.compaction_id = "future-compaction".into(); - worker.persist_compaction_lifecycle(&future).unwrap(); - let error = worker.recover_unfinished_compaction().await.unwrap_err(); - assert!( - error - .to_string() - .contains("unsupported compaction lifecycle schema version 4") - ); + let after = worker.sink.subscribe_with_snapshot().0; + assert_eq!(format!("{after:?}"), format!("{before:?}")); + assert!(!after.iter().any(|entry| { + matches!( + entry, + LogEntry::Extension { domain, .. } + if domain == "yoi.compaction" + ) + })); } fn minimal_manifest() -> WorkerManifest { diff --git a/crates/worker/tests/compact_events_test.rs b/crates/worker/tests/compact_events_test.rs index 719adafd..1bce696e 100644 --- a/crates/worker/tests/compact_events_test.rs +++ b/crates/worker/tests/compact_events_test.rs @@ -406,6 +406,37 @@ fn system_texts_in_sink_session_start( Vec::new() } +#[tokio::test] +async fn active_segment_cas_rejects_stale_compaction_writer() { + let client = MockClient::new(vec![ + single_text_events("seed response"), + write_summary_tool_use_events("summary-1", "replacement summary"), + ]); + let (mut worker, metadata_store, _segment_store) = make_faulting_worker(client).await; + worker.run_text("seed input").await.unwrap(); + let old_segment_id = worker.segment_id(); + let competing_segment_id = uuid::Uuid::now_v7(); + metadata_store + .update_by_name("test-worker", |metadata| { + metadata.active.as_mut().unwrap().segment_id = Some(competing_segment_id); + }) + .unwrap(); + + let _error = worker.compact(0).await.unwrap_err(); + + assert_eq!(worker.segment_id(), old_segment_id); + assert_eq!( + metadata_store + .read_by_name("test-worker") + .unwrap() + .unwrap() + .active + .unwrap() + .segment_id, + Some(competing_segment_id) + ); +} + #[tokio::test] async fn failed_active_segment_commit_keeps_live_and_durable_history_on_old_segment() { let client = MockClient::new(vec![ @@ -614,7 +645,7 @@ async fn compact_emits_session_start_carrying_summary_and_task_snapshot() { } #[tokio::test] -async fn pre_run_compact_success_broadcasts_start_and_done() { +async fn pre_run_compact_publishes_runtime_progress_phases() { // Responses: (1) first run returns short text, (2) compact worker // emits write_summary then closes (two LLM calls inside the compact // worker: one for write_summary, one that the compact loop consumes @@ -640,86 +671,46 @@ async fn pre_run_compact_success_broadcasts_start_and_done() { assert_ne!(worker.segment_id(), segment_before); let events = drain(&mut rx); - let kinds: Vec<&str> = events - .iter() - .map(|e| match e { - Event::CompactStart { .. } => "start", - Event::CompactDone { .. } => "done", - Event::CompactFailed { .. } => "failed", - _ => "other", - }) - .collect(); - assert!( - kinds.contains(&"start") && kinds.contains(&"done"), - "expected CompactStart + CompactDone in {kinds:?}" - ); - assert!( - !kinds.contains(&"failed"), - "unexpected CompactFailed in {kinds:?}" - ); - let starts = events + let progress = events .iter() .filter_map(|event| match event { - Event::CompactStart { lifecycle } => Some(lifecycle), + Event::CompactionProgress { compaction } => { + Some(compaction.as_ref().map(|item| item.phase)) + } _ => None, }) .collect::>(); assert_eq!( - starts.len(), - 2, - "start and Internal Worker binding revisions" + progress, + vec![ + Some(protocol::CompactionPhase::Preparing), + Some(protocol::CompactionPhase::Summarizing), + Some(protocol::CompactionPhase::Committing), + None, + ] ); - assert_eq!(starts[0].compaction_id, starts[1].compaction_id); - assert_eq!(starts[0].revision, 1); - assert!(starts[0].internal_worker.is_none()); - assert_eq!(starts[1].revision, 2); - assert!(matches!( - starts[1].internal_worker.as_ref().map(|worker| &worker.kind), - Some(protocol::InternalWorkerKind::Service { kind }) if kind == "compaction" - )); + assert!(events.iter().all(|event| !matches!( + event, + Event::CompactStart { .. } | Event::CompactDone { .. } | Event::CompactFailed { .. } + ))); assert!(events.iter().any(|event| matches!( event, Event::InternalWorker { worker, .. } if matches!(&worker.kind, protocol::InternalWorkerKind::Service { kind } if kind == "compaction") )), "compactor activity must be projected through the parent stream"); - let completed = events - .iter() - .find_map(|event| match event { - Event::CompactDone { lifecycle } => Some(lifecycle), - _ => None, - }) - .expect("completed lifecycle"); - assert_eq!(completed.compaction_id, starts[0].compaction_id); - assert_eq!(completed.revision, 3); - assert_eq!(completed.summary.as_deref(), Some("summary")); - assert_eq!(completed.state, protocol::CompactionLifecycleState::Done); - let done_index = events - .iter() - .position(|event| matches!(event, Event::CompactDone { .. })) - .expect("done event"); - let removed_index = events - .iter() - .position(|event| matches!(event, Event::InternalWorkerRemoved { .. })) - .expect("terminal compactor session must be released"); - assert!( - done_index < removed_index, - "terminal lifecycle precedes release fence" - ); - // CompactDone carries the new Segment ID; the Session ID is unchanged. - let new_id_in_event = events.iter().find_map(|e| match e { - Event::CompactDone { lifecycle } => lifecycle - .new_segment_id - .as_deref() - .and_then(|value| uuid::Uuid::parse_str(value).ok()), - _ => None, - }); - assert!(new_id_in_event.is_some(), "CompactDone missing"); - assert_eq!(new_id_in_event.unwrap(), worker.segment_id()); + let active_entries = worker + .store() + .read_all(worker.session_id(), worker.segment_id()) + .unwrap(); + assert!(!active_entries.iter().any(|entry| matches!( + entry, + LogEntry::Extension { domain, .. } if domain == "yoi.compaction" + ))); } #[tokio::test] -async fn mid_turn_compact_success_broadcasts_start_and_done() { +async fn request_threshold_compact_publishes_runtime_progress() { // Path: `do_compact_and_resume` via PreRequestAction::Yield. // // Sequence of LLM calls the mock will serve: @@ -748,36 +739,20 @@ async fn mid_turn_compact_success_broadcasts_start_and_done() { worker.run_text("second").await.unwrap(); let events = drain(&mut rx); - let kinds: Vec<&str> = events - .iter() - .map(|e| match e { - Event::CompactStart { .. } => "start", - Event::CompactDone { .. } => "done", - Event::CompactFailed { .. } => "failed", - _ => "other", - }) - .collect(); + assert!(events.iter().any(|event| matches!( + event, + Event::CompactionProgress { compaction: Some(progress) } + if progress.phase == protocol::CompactionPhase::Committing + ))); assert!( - kinds.contains(&"start") && kinds.contains(&"done"), - "expected CompactStart + CompactDone in {kinds:?}" + events + .iter() + .any(|event| matches!(event, Event::CompactionProgress { compaction: None })) ); - assert!( - !kinds.contains(&"failed"), - "unexpected CompactFailed in {kinds:?}" - ); - - let new_id_in_event = events.iter().find_map(|e| match e { - Event::CompactDone { lifecycle } => lifecycle - .new_segment_id - .as_deref() - .and_then(|value| uuid::Uuid::parse_str(value).ok()), - _ => None, - }); - assert_eq!(new_id_in_event, Some(worker.segment_id())); } #[tokio::test] -async fn pre_run_compact_failure_broadcasts_start_and_failed() { +async fn pre_run_compact_failure_clears_runtime_progress() { // Only the first run has a response. Compaction will run the // compact worker which immediately exhausts the mock → failure. let client = MockClient::new(vec![single_text_events("hi")]); @@ -789,31 +764,28 @@ async fn pre_run_compact_failure_broadcasts_start_and_failed() { worker.run_text("first").await.unwrap(); let _ = drain(&mut rx); - // Best-effort: returns Ok(()) even on failure, but emits CompactFailed. + // Best-effort: returns Ok(()) even on failure and clears runtime progress. worker.try_pre_run_compact().await; let events = drain(&mut rx); - let kinds: Vec<&str> = events - .iter() - .map(|e| match e { - Event::CompactStart { .. } => "start", - Event::CompactDone { .. } => "done", - Event::CompactFailed { .. } => "failed", - _ => "other", - }) - .collect(); + assert!(events.iter().any(|event| matches!( + event, + Event::CompactionProgress { compaction: Some(progress) } + if progress.phase == protocol::CompactionPhase::Preparing + ))); assert!( - kinds.contains(&"start") && kinds.contains(&"failed"), - "expected CompactStart + CompactFailed in {kinds:?}" - ); - assert!( - !kinds.contains(&"done"), - "unexpected CompactDone in {kinds:?}" + events + .iter() + .any(|event| matches!(event, Event::CompactionProgress { compaction: None })) ); + assert!(events.iter().all(|event| !matches!( + event, + Event::CompactStart { .. } | Event::CompactDone { .. } | Event::CompactFailed { .. } + ))); } #[tokio::test] -async fn manual_compact_cancel_terminalizes_before_returning_idle() { +async fn manual_compact_cancel_clears_progress_before_returning_idle() { let worker = make_worker_with_manifest(POST_RUN_MANIFEST_TOML, BlockingCompactClient::new()).await; let runtime_tmp = tempfile::tempdir().unwrap(); @@ -856,7 +828,9 @@ async fn manual_compact_cancel_terminalizes_before_returning_idle() { .await .expect("timeout waiting for compact start") .expect("event"), - Event::CompactStart { .. } + Event::CompactionProgress { + compaction: Some(_) + } ) { break; } @@ -875,9 +849,7 @@ async fn manual_compact_cancel_terminalizes_before_returning_idle() { .expect("timeout waiting for compact cancellation") .expect("event") { - Event::CompactFailed { lifecycle } - if lifecycle.state == protocol::CompactionLifecycleState::Interrupted => - { + Event::CompactionProgress { compaction: None } => { saw_interrupted = true; } Event::WorkerState { snapshot } @@ -904,7 +876,9 @@ async fn manual_compact_cancel_terminalizes_before_returning_idle() { .await .expect("timeout waiting for second compact start") .expect("event"), - Event::CompactStart { .. } + Event::CompactionProgress { + compaction: Some(_) + } ) { break; } @@ -922,9 +896,7 @@ async fn manual_compact_cancel_terminalizes_before_returning_idle() { .expect("timeout waiting for shutdown") .expect("event") { - Event::CompactFailed { lifecycle } - if lifecycle.state == protocol::CompactionLifecycleState::Interrupted => - { + Event::CompactionProgress { compaction: None } => { interrupted_before_shutdown = true; } Event::Shutdown => { @@ -944,7 +916,7 @@ async fn manual_compact_cancel_terminalizes_before_returning_idle() { } #[tokio::test] -async fn controller_compact_method_emits_start_and_done() { +async fn controller_compact_method_publishes_progress_and_clear() { let client = MockClient::new(vec![ text_events_with_usage("hi", 1000), write_summary_tool_use_events("manual-summary", "manual compact summary"), @@ -991,14 +963,12 @@ async fn controller_compact_method_emits_start_and_done() { .expect("timeout waiting for compact events") .expect("event") { - Event::CompactStart { .. } => saw_start = true, - Event::CompactDone { .. } => { + Event::CompactionProgress { + compaction: Some(_), + } => saw_start = true, + Event::CompactionProgress { compaction: None } => { break; } - Event::CompactFailed { lifecycle } => panic!( - "manual compact failed: {}", - lifecycle.error.as_deref().unwrap_or("unknown error") - ), _ => {} } } diff --git a/web/workspace/src/lib/generated/protocol.ts b/web/workspace/src/lib/generated/protocol.ts index 2b1ba468..621fa566 100644 --- a/web/workspace/src/lib/generated/protocol.ts +++ b/web/workspace/src/lib/generated/protocol.ts @@ -75,6 +75,10 @@ export type CompactionLifecycle = { schema_version: number, compaction_id: strin */ started_at_ms: number, ended_at_ms?: number | null, summary?: string | null, error?: string | null, new_segment_id?: string | null, }; +export type CompactionPhase = "preparing" | "summarizing" | "committing"; + +export type CompactionTrigger = "manual" | "pre_run" | "request_threshold"; + export type UploadedFileAvailability = "available" | "unavailable" | "integrity_failed"; export type UploadedFileRef = { artifact_id: string, file_name: string, media_type: string, created_at_ms: number, availability: UploadedFileAvailability, byte_len: number, sha256: string, source_entry_id?: string | null, }; @@ -113,7 +117,7 @@ export type RewindSummary = { truncated_to_entries: number, discarded_entries: n export type InFlightBlock = { "kind": "text", text: string, finished?: boolean, } | { "kind": "thinking", text: string, finished?: boolean, } | { "kind": "tool_call", id: string, name: string, args: string, state?: InFlightToolCallState, }; -export type InFlightCompaction = { schema_version: number, compaction_id: string, revision: number, internal_worker: InternalWorkerRef | null, started_at_ms: number, }; +export type InFlightCompaction = { phase: CompactionPhase, started_at_ms: number, trigger: CompactionTrigger, }; export type InFlightSnapshot = { blocks?: Array, commands?: Array, /** @@ -318,4 +322,4 @@ in_flight?: InFlightSnapshot, * Parent-owned Internal Worker sessions visible to this client. * Service-private Internal Workers are deliberately excluded. */ -internal_workers?: Array, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "internal_worker_removed", "data": { worker: InternalWorkerRef, revision: number, } } | { "event": "segment_rotated", "data": { session: SessionSnapshot, } } | { "event": "worker_state", "data": { snapshot: WorkerStateSnapshot, } } | { "event": "command_acknowledged", "data": { acknowledgement: WorkerCommandAcknowledgement, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array, } } | { "event": "rewind_applied", "data": { session: SessionSnapshot, input: Array, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compact_start", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_done", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_failed", "data": { lifecycle: CompactionLifecycle, } } | { "event": "shutdown" }; +internal_workers?: Array, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "internal_worker_removed", "data": { worker: InternalWorkerRef, revision: number, } } | { "event": "segment_rotated", "data": { session: SessionSnapshot, } } | { "event": "worker_state", "data": { snapshot: WorkerStateSnapshot, } } | { "event": "command_acknowledged", "data": { acknowledgement: WorkerCommandAcknowledgement, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array, } } | { "event": "rewind_applied", "data": { session: SessionSnapshot, input: Array, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compaction_progress", "data": { compaction: InFlightCompaction | null, } } | { "event": "compact_start", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_done", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_failed", "data": { lifecycle: CompactionLifecycle, } } | { "event": "shutdown" }; diff --git a/web/workspace/src/lib/workspace/console/model.test.ts b/web/workspace/src/lib/workspace/console/model.test.ts index e2c29985..4e5b118a 100644 --- a/web/workspace/src/lib/workspace/console/model.test.ts +++ b/web/workspace/src/lib/workspace/console/model.test.ts @@ -1084,19 +1084,18 @@ Deno.test("snapshot restores running compaction without staged content", () => { snapshot.data.in_flight = { blocks: [], compaction: { - schema_version: 3, - compaction_id: "compaction-snapshot", - revision: 1, - internal_worker: null, + phase: "summarizing", started_at_ms: 1_000, + trigger: "manual", }, }; const projection = projectConsole([{ eventId: "snapshot", event: snapshot }]); assertEquals(projection.lines.length, 1); - assertEquals(projection.lines[0].id, "compaction-compaction-snapshot"); - assertEquals(projection.lines[0].compaction?.state, "running"); + assertEquals(projection.lines[0].id, "compaction-runtime"); + assertEquals(projection.lines[0].streaming, true); + assertEquals(projection.lines[0].body, "compacting · summarizing"); assertEquals(projection.lines[0].body.includes("staged"), false); }); diff --git a/web/workspace/src/lib/workspace/console/model.ts b/web/workspace/src/lib/workspace/console/model.ts index e44dcc0a..810ea3a9 100644 --- a/web/workspace/src/lib/workspace/console/model.ts +++ b/web/workspace/src/lib/workspace/console/model.ts @@ -720,16 +720,25 @@ function compactionActivity( function applyInFlightCompaction( projection: ConsoleProjection, - progress: InFlightCompaction, + progress: InFlightCompaction | null, ): ConsoleProjection { - return applyCompactionLifecycle(projection, { - ...progress, - state: "running", - ended_at_ms: null, - summary: null, - error: null, - new_segment_id: null, - }); + const id = "compaction-runtime"; + const lines = projection.lines.filter((line) => line.id !== id); + if (!progress) return { ...projection, lines }; + return { + ...projection, + lines: [ + ...lines, + { + id, + kind: "status", + title: "Compaction", + body: `compacting · ${progress.phase.replaceAll("_", " ")}`, + source: "event", + streaming: true, + }, + ], + }; } function applyCompactionLifecycle( @@ -1110,6 +1119,8 @@ export function applyProtocolEvent( // These are protocol/status/control events. TUI Console does not append // them to the conversation surface; browser Console should not either. break; + case "compaction_progress": + return applyInFlightCompaction(next, event.data.compaction ?? null); case "compact_start": case "compact_done": case "compact_failed":