diff --git a/crates/op-editor-core/src/agent_indicators.rs b/crates/op-editor-core/src/agent_indicators.rs index 179855337..e8d2d21ea 100644 --- a/crates/op-editor-core/src/agent_indicators.rs +++ b/crates/op-editor-core/src/agent_indicators.rs @@ -137,6 +137,23 @@ pub fn clear_if_epoch(epoch: u64) { } } +/// End the run identified by `epoch`: if it's still the active run, clear +/// its indicators *and retire the epoch* (bump it), so any registration +/// still in flight from that run's worker no-ops instead of re-populating +/// the set we just cleared. A no-op if a newer run already took over. +/// +/// This is what the host calls when a turn is stopped: `clear_if_epoch` +/// alone would clear, but a worker mid `add_frame` loop (scaffold just +/// built, first cancel-detecting channel op not yet reached) would re-add +/// under the unchanged epoch. Retiring the epoch closes that window. +pub fn end_if_epoch(epoch: u64) { + let mut r = REGISTRY.lock().unwrap(); + if r.epoch == epoch { + r.clear_maps(); + r.epoch += 1; + } +} + /// `true` while any node / frame indicator is active — the paint loop /// uses this to keep requesting redraws so the breathing animates. pub fn is_active() -> bool { @@ -204,5 +221,27 @@ mod tests { clear_if_epoch(e2); assert!(snapshot().frames.is_empty()); assert!(!is_active()); + + // end_if_epoch retires the epoch: after the host ends a run, a + // worker still mid-registration under it can't re-populate. + let e3 = begin(); + add_frame(e3, "f3", "#FFD93D", "Pixel"); + end_if_epoch(e3); + assert!(snapshot().frames.is_empty(), "end_if_epoch clears"); + add_frame(e3, "late", "#FFD93D", "Pixel"); // in-flight registration + assert!( + snapshot().frames.is_empty(), + "registration under a retired epoch no-ops" + ); + + // end_if_epoch on a stale epoch must not touch the live run. + let e4 = begin(); + add_frame(e4, "f4", "#6C5CE7", "Echo"); + end_if_epoch(e3); + assert!( + snapshot().frames.contains_key("f4"), + "end_if_epoch ignores a stale epoch" + ); + clear(); } } diff --git a/crates/op-host-desktop/src/design_session.rs b/crates/op-host-desktop/src/design_session.rs index 7c3428371..90b4fa2cf 100644 --- a/crates/op-host-desktop/src/design_session.rs +++ b/crates/op-host-desktop/src/design_session.rs @@ -63,6 +63,22 @@ pub struct DesignSession { delta_rx: Receiver, cmd_rx: Receiver, finished: bool, + /// Agent-team canvas-indicator epoch for this turn (see + /// [`op_editor_core::agent_indicators`]). Minted on `start`; the + /// [`Drop`] impl clears it so stop / new-chat wipes the breathing + /// borders immediately instead of waiting for the worker to unwind. + indicator_epoch: u64, +} + +impl Drop for DesignSession { + fn drop(&mut self) { + // Stop / new-chat / completion all drop the session. Clear this + // turn's agent-team indicators right away and retire the epoch — + // epoch-scoped, so if a newer turn already began it keeps its own, + // and a worker still mid-registration for this turn can't re-add + // after we've cleared. + op_editor_core::agent_indicators::end_if_epoch(self.indicator_epoch); + } } /// Progress / completion events emitted by the worker. @@ -120,6 +136,13 @@ impl DesignSession { let (delta_tx, delta_rx) = mpsc::channel::(); let (cmd_tx, cmd_rx) = mpsc::channel::(); + // Mint this turn's indicator epoch on the UI thread BEFORE spawning + // the worker. `begin` clears any prior turn's borders immediately + // and hands us the epoch so Drop can target exactly this run — the + // worker registers its frames under the same epoch via + // `with_indicator_epoch` below. + let indicator_epoch = op_editor_core::agent_indicators::begin(); + thread::Builder::new() .name("op-design-turn".into()) .spawn(move || { @@ -138,14 +161,18 @@ impl DesignSession { let mut on_progress = move |p: Progress| { let _ = delta_tx_for_progress.send(DesignDelta::Progress(p)); }; - let summary = shared_runtime().block_on(Orchestrator::new().run( - request, - &mut sink, - &llm, - &mut on_progress, - &abort, - &providers, - )); + let summary = shared_runtime().block_on( + Orchestrator::new() + .with_indicator_epoch(indicator_epoch) + .run( + request, + &mut sink, + &llm, + &mut on_progress, + &abort, + &providers, + ), + ); let _ = delta_tx.send(DesignDelta::Done(summary)); }) .expect("spawn op-design-turn thread"); @@ -154,6 +181,7 @@ impl DesignSession { delta_rx, cmd_rx, finished: false, + indicator_epoch, } } @@ -207,6 +235,9 @@ impl DesignSession { delta_rx, cmd_rx, finished: false, + // No real turn behind these channels — epoch 0 never matches a + // live run, so the Drop clear is a harmless no-op. + indicator_epoch: 0, } } } diff --git a/crates/op-orchestrator/src/run.rs b/crates/op-orchestrator/src/run.rs index ac8aab7c7..5ef7fcfee 100644 --- a/crates/op-orchestrator/src/run.rs +++ b/crates/op-orchestrator/src/run.rs @@ -42,14 +42,30 @@ use crate::variables::{rollback, seed_commands, snapshot_plan_vars}; use futures::StreamExt; use op_editor_core::{EditorCommand, NodeId, PenNodeExt}; -/// 设计编排器。S3a 阶段无构造期配置 —— 保留 struct 以便 S3b/S3c -/// 挂选项。 +/// 设计编排器。 #[derive(Debug, Default, Clone, Copy)] -pub struct Orchestrator; +pub struct Orchestrator { + /// Run epoch for the agent-team canvas indicators. The host owns the + /// design-turn lifecycle, so it mints the epoch (`agent_indicators:: + /// begin`) and clears via `clear_if_epoch` the instant the turn is + /// stopped — registration in the concurrent path must run under that + /// same epoch. `None` for headless / test callers, which then let the + /// concurrent path mint its own epoch. + agent_indicator_epoch: Option, +} impl Orchestrator { pub fn new() -> Self { - Orchestrator + Self::default() + } + + /// Adopt a host-owned indicator epoch. The host clears with + /// `agent_indicators::clear_if_epoch(epoch)` on stop / new-chat, so + /// the concurrent path registers under this epoch instead of minting + /// its own — otherwise the host couldn't target the right run. + pub fn with_indicator_epoch(mut self, epoch: u64) -> Self { + self.agent_indicator_epoch = Some(epoch); + self } /// 跑一次完整编排。见 spec §4 数据流。 @@ -102,6 +118,7 @@ impl Orchestrator { on_progress, abort, providers, + self.agent_indicator_epoch, ) .await; } @@ -416,6 +433,7 @@ async fn run_concurrent_path( on_progress: &mut dyn FnMut(Progress), abort: &AbortFlag, providers: &ValidationProviders<'_>, + host_epoch: Option, ) -> Result { // -- 进入"已动文档"区 -- sink.begin_undo_batch(); @@ -485,11 +503,13 @@ async fn run_concurrent_path( // Tag each group's root frame with a distinct agent identity so the // canvas can paint a per-agent breathing border while the team works. let identities = crate::agent_identity::assign_agent_identities(actual_root_ids.len()); - // `begin` bumps the run epoch and clears any prior indicators. The - // guard below clears on every exit path (finish / error / cancelled - // worker), but only while this run is still the active epoch — a newer - // run that started in the meantime keeps its own indicators. - let epoch = op_editor_core::agent_indicators::begin(); + // Adopt the host-minted epoch when present — the host already called + // `begin` (bumping the epoch + clearing the prior run) at turn start so + // it can clear immediately on stop. Headless / test callers pass `None` + // and we mint our own. Either way the guard below clears on every exit + // path (finish / error / cancelled worker), but only while this run is + // still the active epoch — a newer run keeps its own indicators. + let epoch = host_epoch.unwrap_or_else(op_editor_core::agent_indicators::begin); for (root_id, identity) in actual_root_ids.iter().zip(identities.iter()) { op_editor_core::agent_indicators::add_frame( epoch,