fix(editor): clear agent indicators the instant a design turn stops

Stopping a turn or starting a new chat drops the DesignSession, but the
canvas indicators only cleared when the worker thread next touched its
channel and unwound — which lags seconds behind when the worker is
mid-LLM-call. The breathing borders lingered on a canvas the user had
already moved on from.

Give the host (which owns the turn lifecycle) the indicator epoch:
DesignSession::start mints it (clearing the prior turn at once) and the
session's Drop ends it immediately via end_if_epoch. The worker registers
under the same epoch via Orchestrator::with_indicator_epoch; headless and
test callers pass none and the concurrent path mints its own, so run()'s
signature and the DesignRequest construction sites are untouched.

end_if_epoch clears AND retires the epoch (vs clear_if_epoch which only
clears), so a worker still in its add_frame loop when the turn is stopped
can't re-populate the set under the now-stale epoch. It's epoch-scoped, so
a newer turn that already began is left untouched.
This commit is contained in:
Fini 2026-06-10 11:18:38 +08:00
parent 56e55f1159
commit 09e1f5f606
3 changed files with 107 additions and 17 deletions

View file

@ -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();
}
}

View file

@ -63,6 +63,22 @@ pub struct DesignSession {
delta_rx: Receiver<DesignDelta>,
cmd_rx: Receiver<DesignCmdReq>,
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::<DesignDelta>();
let (cmd_tx, cmd_rx) = mpsc::channel::<DesignCmdReq>();
// 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,
}
}
}

View file

@ -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<u64>,
}
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<u64>,
) -> Result<RunSummary, OrchestratorError> {
// -- 进入"已动文档"区 --
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,