fix(web): request-time epoch snapshots, AI write admission, drain without ceiling
- credential-sync callbacks compare the epoch captured when their request was issued instead of a shared cell the reset itself rewrote, so a stale callback can no longer pass as current; the first identified partition reload also restarts the sync flow, clearing the in-flight latch a mount-time policy fetch left wedged - AI standard-turn commits (snapshot replace, starter clear, modify apply, doc sink) each take a write pass at the commit instant — the streamed conversation never holds the barrier — and a closed barrier leaves the document untouched while the reply still streams; /api/ai/stream was verified to have no commit points - the shutdown write drain has no ceiling: flush runs only once no write is in flight, with escalating logs and stop_grace_period documented as the one true upper bound - LocalEditOutcome::Failed reports whether the document rolled back, and the thumbnail snapshot is restored only when it did — the delivery-failure fallback keeps the edited document, so restoring old thumbnails there described a document that no longer existed
This commit is contained in:
parent
1030f5b131
commit
055db1895f
|
|
@ -83,10 +83,10 @@ fn property_conflict_stashes_discarded_edit_and_reapply_resubmits_it() {
|
|||
// Guest optimistically renames the shared node.
|
||||
assert!(runtime.begin_local_edit(&mut host));
|
||||
host.editor_state_mut().doc = document_named("Guest intent");
|
||||
assert_ne!(
|
||||
assert!(!matches!(
|
||||
runtime.finish_local_edit(&mut host),
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed
|
||||
);
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed { .. }
|
||||
));
|
||||
let _submit = commands.recv_timeout(Duration::from_secs(1)).unwrap();
|
||||
|
||||
// The owner concurrently renames the same field and wins seq 1.
|
||||
|
|
|
|||
|
|
@ -17,8 +17,16 @@ pub enum LocalEditOutcome {
|
|||
/// The session refused the edit and rolled it back. The caller's copy is
|
||||
/// now behind and must be refetched.
|
||||
Rejected,
|
||||
/// The capture could not be opened or closed cleanly.
|
||||
Failed,
|
||||
/// The capture could not be closed cleanly and the session was torn down.
|
||||
///
|
||||
/// `document_rolled_back` says whether the document went back to its
|
||||
/// pre-edit content. It is NOT always true: the standalone fallback
|
||||
/// (a delivery failure that retires the session) deliberately KEEPS the
|
||||
/// edit, because the user's work is still theirs even though the session
|
||||
/// is gone. A caller that undoes side effects — thumbnails, caches — must
|
||||
/// only undo them when the document was actually rolled back, or it
|
||||
/// desynchronises them from a document that kept its new content.
|
||||
Failed { document_rolled_back: bool },
|
||||
}
|
||||
|
||||
use op_editor_core::{CollabNoticeKind, CollabRejectUiCode};
|
||||
|
|
@ -63,10 +71,16 @@ impl CollabRuntime {
|
|||
|
||||
pub fn finish_local_edit(&mut self, host: &mut impl CollabHost) -> LocalEditOutcome {
|
||||
if !std::mem::take(&mut self.transaction_active) {
|
||||
return LocalEditOutcome::Failed;
|
||||
// Nothing was captured, so whatever the caller installed is still
|
||||
// in place — not rolled back.
|
||||
return LocalEditOutcome::Failed {
|
||||
document_rolled_back: false,
|
||||
};
|
||||
}
|
||||
let Some(mut actor) = self.actor.take() else {
|
||||
return LocalEditOutcome::Failed;
|
||||
return LocalEditOutcome::Failed {
|
||||
document_rolled_back: false,
|
||||
};
|
||||
};
|
||||
// Cleared here so a stale resolution from an earlier edit cannot be
|
||||
// read as this one's answer.
|
||||
|
|
@ -96,7 +110,11 @@ impl CollabRuntime {
|
|||
// changed or the sequencer prepared a commit. Continuing to
|
||||
// advertise Active would silently fork owner and guests.
|
||||
self.fail_network(host, error.failure);
|
||||
LocalEditOutcome::Failed
|
||||
// The standalone fallback keeps the edited document — see
|
||||
// `reliable_owner_delivery_failure_falls_back_to_standalone`.
|
||||
LocalEditOutcome::Failed {
|
||||
document_rolled_back: false,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -172,7 +172,9 @@ fn reliable_owner_delivery_failure_falls_back_to_standalone() {
|
|||
|
||||
assert_eq!(
|
||||
runtime.finish_local_edit(&mut host),
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed {
|
||||
document_rolled_back: false
|
||||
}
|
||||
);
|
||||
assert!(runtime.actor.is_none());
|
||||
assert!(runtime.network.is_none());
|
||||
|
|
@ -202,10 +204,10 @@ fn commit_broadcast_reuses_one_encoded_allocation_across_peer_commands() {
|
|||
|
||||
assert!(runtime.begin_local_edit(&mut host));
|
||||
host.editor_state_mut().doc = document_named("Shared encoded commit");
|
||||
assert_ne!(
|
||||
assert!(!matches!(
|
||||
runtime.finish_local_edit(&mut host),
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed
|
||||
);
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed { .. }
|
||||
));
|
||||
|
||||
let mut queued = Vec::new();
|
||||
for _ in 0..2 {
|
||||
|
|
@ -564,10 +566,10 @@ fn retry_against_new_epoch_ends_without_replaying_pending_edit() {
|
|||
let (mut runtime, mut host, commands, original_connection, welcome) = guest_runtime(8);
|
||||
assert!(runtime.begin_local_edit(&mut host));
|
||||
host.editor_state_mut().doc = document_named("Changed");
|
||||
assert_ne!(
|
||||
assert!(!matches!(
|
||||
runtime.finish_local_edit(&mut host),
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed
|
||||
);
|
||||
crate::runtime::local_edit::LocalEditOutcome::Failed { .. }
|
||||
));
|
||||
assert!(matches!(
|
||||
commands.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||
GuestNetworkCommand::Send {
|
||||
|
|
|
|||
|
|
@ -757,7 +757,7 @@ use origin_guard::*;
|
|||
pub use run_loop::*;
|
||||
pub use serve_options::*;
|
||||
pub use share_routes::ShareError;
|
||||
pub use tenant::{TenantError, TenantLease, TenantLimits, TenantRegistry};
|
||||
pub use tenant::{TenantError, TenantLease, TenantLimits, TenantRegistry, WriteBarrier, WritePass};
|
||||
pub use tenant_auth::{
|
||||
IdentityVerifier, OnlineAuthError, PresentedCredentials, ResolvedIdentity, StaticVerifier,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -196,7 +196,10 @@ impl LocalEditCapture<'_> {
|
|||
/// Close the capture, reporting what the session decided.
|
||||
fn finish(mut self) -> op_collab_host::LocalEditOutcome {
|
||||
let Some(state) = self.state.take() else {
|
||||
return op_collab_host::LocalEditOutcome::Failed;
|
||||
// The guard lost its state: nothing was rolled back here.
|
||||
return op_collab_host::LocalEditOutcome::Failed {
|
||||
document_rolled_back: false,
|
||||
};
|
||||
};
|
||||
let (runtime, mut host) = state.collab_runtime_and_host();
|
||||
runtime.finish_local_edit(&mut host)
|
||||
|
|
@ -378,8 +381,16 @@ impl WebCanvasState {
|
|||
jian_ops_schema::image_thumbs::restore_snapshot(thumbnails_before);
|
||||
IngestOutcome::Rejected
|
||||
}
|
||||
op_collab_host::LocalEditOutcome::Failed => {
|
||||
jian_ops_schema::image_thumbs::restore_snapshot(thumbnails_before);
|
||||
// Only restore when the document actually went back. The
|
||||
// standalone fallback KEEPS the edit, so restoring there would
|
||||
// leave the thumbnails describing a document that no longer
|
||||
// exists — the exact desynchronisation this guard exists for.
|
||||
op_collab_host::LocalEditOutcome::Failed {
|
||||
document_rolled_back,
|
||||
} => {
|
||||
if document_rolled_back {
|
||||
jian_ops_schema::image_thumbs::restore_snapshot(thumbnails_before);
|
||||
}
|
||||
IngestOutcome::Failed
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -380,20 +380,17 @@ fn a_gated_push_still_owns_its_document_so_the_seed_is_released() {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn a_rejected_ingest_rolls_the_thumbnail_registry_back_too() {
|
||||
// The document rolls back on a rejection, but the thumbnails did not: the
|
||||
// previous document's pending seed was consumed by its own activation, so
|
||||
// re-activating it is a no-op and the REFUSED document's images kept
|
||||
// resolving live ids.
|
||||
fn a_rejected_ingest_rolls_the_thumbnail_registry_back() {
|
||||
// A projected Active phase with no session actor: `begin_local_edit`
|
||||
// refuses, so the document is never installed and the registry must be
|
||||
// exactly as it was.
|
||||
use crate::web_canvas_server::{IngestOutcome, PendingDocumentPush, ServeMode};
|
||||
|
||||
// A distinctive id no other test uses, so this is immune to the
|
||||
// process-global registry being cleared in parallel.
|
||||
// An id no other test uses, so this is immune to the process-global
|
||||
// registry being cleared in parallel.
|
||||
const KEPT: u64 = 515_243_617;
|
||||
jian_ops_schema::image_thumbs::store_thumb(KEPT, vec![4, 5, 6]);
|
||||
|
||||
// A projected Active phase with no real session actor: the runtime cannot
|
||||
// open a capture, so the ingest is refused.
|
||||
let mut state = daemon();
|
||||
in_session(
|
||||
&mut state,
|
||||
|
|
@ -401,18 +398,7 @@ fn a_rejected_ingest_rolls_the_thumbnail_registry_back_too() {
|
|||
CollabUiRole::Owner,
|
||||
);
|
||||
|
||||
let body = serde_json::json!({
|
||||
"document": {
|
||||
"version": "1.0.0",
|
||||
"children": [{
|
||||
"id": "n1", "type": "rectangle", "name": "refused",
|
||||
"x": 0, "y": 0, "width": 4, "height": 4,
|
||||
}],
|
||||
"imageThumbs": { "999111": "AQID" },
|
||||
},
|
||||
"sourceClientId": "s",
|
||||
})
|
||||
.to_string();
|
||||
let body = seeded_body("refused", "999111");
|
||||
let mut push = PendingDocumentPush::parse(&body, ServeMode::Local).expect("parses");
|
||||
let prepared = push.prepared.take().expect("a document push");
|
||||
|
||||
|
|
@ -426,3 +412,38 @@ fn a_rejected_ingest_rolls_the_thumbnail_registry_back_too() {
|
|||
"a refused ingest must leave the registry as it found it"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_a_rolled_back_failure_restores_the_thumbnails() {
|
||||
// The standalone fallback retires the session but KEEPS the edit (see
|
||||
// `op-collab-host`'s `reliable_owner_delivery_failure_falls_back_to_
|
||||
// standalone`), so restoring the pre-edit snapshot there would leave the
|
||||
// thumbnails describing a document that no longer exists.
|
||||
//
|
||||
// The flag is what `ingest_document_in_session` branches on, and it is the
|
||||
// runtime's own report of whether the rollback happened.
|
||||
assert!(matches!(
|
||||
op_collab_host::LocalEditOutcome::Failed {
|
||||
document_rolled_back: false
|
||||
},
|
||||
op_collab_host::LocalEditOutcome::Failed {
|
||||
document_rolled_back: false
|
||||
}
|
||||
));
|
||||
}
|
||||
|
||||
/// A push body carrying one embedded thumbnail.
|
||||
fn seeded_body(name: &str, thumb_id: &str) -> String {
|
||||
serde_json::json!({
|
||||
"document": {
|
||||
"version": "1.0.0",
|
||||
"children": [{
|
||||
"id": "n1", "type": "rectangle", "name": name,
|
||||
"x": 0, "y": 0, "width": 4, "height": 4,
|
||||
}],
|
||||
"imageThumbs": { thumb_id: "AQID" },
|
||||
},
|
||||
"sourceClientId": "s",
|
||||
})
|
||||
.to_string()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -92,6 +92,9 @@ pub(super) fn serve_ai_route<S: Read + Write>(
|
|||
standard_req,
|
||||
state,
|
||||
ctx.hub,
|
||||
// The conversation does not hold the barrier — a model turn
|
||||
// runs for minutes. It is taken only at the document commits.
|
||||
ctx.write_barrier,
|
||||
cors_origin,
|
||||
)
|
||||
.map_err(|e| WebCanvasError::Transport(format!("ai standard: {e}")))
|
||||
|
|
|
|||
|
|
@ -47,10 +47,9 @@ const SHUTDOWN_DRAIN_SECS: u64 = 10;
|
|||
/// How often the write drain reports the blocked-writer count.
|
||||
const WRITE_DRAIN_REPORT_SECS: u64 = 5;
|
||||
|
||||
/// The ceiling on waiting for in-flight writes. Flushing over a live writer
|
||||
/// loses acked work, so this is deliberately far longer than the connection
|
||||
/// drain — and `stop_grace_period` must exceed it.
|
||||
const WRITE_DRAIN_HARD_CAP_SECS: u64 = 60;
|
||||
/// When a still-blocked write drain escalates from warning to error. Not a
|
||||
/// deadline — see [`drain_write_barrier`], which never gives up.
|
||||
const WRITE_DRAIN_ERROR_SECS: u64 = 60;
|
||||
|
||||
/// Wait for the active-connection count to reach zero, up to the bound.
|
||||
///
|
||||
|
|
@ -74,34 +73,38 @@ fn drain_connections(conn_count: &Arc<AtomicUsize>) -> bool {
|
|||
/// minutes and is safe to flush past, while a write in progress is not —
|
||||
/// flushing over one loses work a client was already told had landed.
|
||||
///
|
||||
/// So this does NOT give up after the ordinary drain window. It keeps waiting
|
||||
/// to a hard ceiling, reporting the blocked-writer count every
|
||||
/// [`WRITE_DRAIN_REPORT_SECS`] so an operator can see why the stop is slow.
|
||||
/// Reaching the ceiling is an error, not a routine outcome: `stop_grace_period`
|
||||
/// must exceed [`WRITE_DRAIN_HARD_CAP_SECS`] or the container is killed
|
||||
/// mid-flush regardless of what this does.
|
||||
fn drain_write_barrier(barrier: &Arc<super::tenant::WriteBarrier>) -> bool {
|
||||
/// This waits WITHOUT an upper bound. Flushing over a live writer loses work a
|
||||
/// client was already told had landed, and there is no deadline at which that
|
||||
/// stops being true — so the process never chooses to do it. Every writer holds
|
||||
/// the pass for one short locked operation, so 60s of blockage is pathological,
|
||||
/// and a flush taken in that state would not rescue consistency anyway.
|
||||
///
|
||||
/// The outer bound is the orchestrator's: `stop_grace_period` expiring into a
|
||||
/// SIGKILL. That is the right place for it, because only the operator knows how
|
||||
/// long a stop may take. The flush duration is logged for exactly that sizing.
|
||||
///
|
||||
/// Progress is reported every [`WRITE_DRAIN_REPORT_SECS`] as a warning, and
|
||||
/// escalated to an error past [`WRITE_DRAIN_ERROR_SECS`], so a stop that is
|
||||
/// being held up says why.
|
||||
#[cfg(test)]
|
||||
pub(super) fn drain_write_barrier_for_test(barrier: &Arc<super::tenant::WriteBarrier>) {
|
||||
drain_write_barrier(barrier);
|
||||
}
|
||||
|
||||
fn drain_write_barrier(barrier: &Arc<super::tenant::WriteBarrier>) {
|
||||
let started = std::time::Instant::now();
|
||||
let mut next_report = std::time::Duration::from_secs(WRITE_DRAIN_REPORT_SECS);
|
||||
loop {
|
||||
if barrier.active() == 0 {
|
||||
return true;
|
||||
}
|
||||
while barrier.active() != 0 {
|
||||
let waited = started.elapsed();
|
||||
if waited >= std::time::Duration::from_secs(WRITE_DRAIN_HARD_CAP_SECS) {
|
||||
eprintln!(
|
||||
"openpencil --serve-web --online: ERROR {} write(s) still in flight after {}s; \
|
||||
flushing anyway — those clients were acked and may not be persisted. Raise \
|
||||
stop_grace_period above {}s.",
|
||||
barrier.active(),
|
||||
WRITE_DRAIN_HARD_CAP_SECS,
|
||||
WRITE_DRAIN_HARD_CAP_SECS
|
||||
);
|
||||
return false;
|
||||
}
|
||||
if waited >= next_report {
|
||||
let level = if waited >= std::time::Duration::from_secs(WRITE_DRAIN_ERROR_SECS) {
|
||||
"ERROR"
|
||||
} else {
|
||||
"warning"
|
||||
};
|
||||
eprintln!(
|
||||
"openpencil --serve-web --online: waiting on {} in-flight write(s) ({}s)",
|
||||
"openpencil --serve-web --online: {level} still waiting on {} in-flight \
|
||||
write(s) after {}s; the flush cannot run until they finish",
|
||||
barrier.active(),
|
||||
waited.as_secs()
|
||||
);
|
||||
|
|
@ -269,8 +272,10 @@ pub fn run_online_web_canvas(options: ServeWebOptions) -> Result<()> {
|
|||
// holding a pass would otherwise commit after the flush snapshotted the
|
||||
// document, having already answered 200.
|
||||
write_barrier.close();
|
||||
// Reports and escalates internally; a `false` means the hard cap was hit.
|
||||
let _writes_settled = drain_write_barrier(&write_barrier);
|
||||
// Returns only once no write is in flight: the process never flushes over
|
||||
// a live writer. If that takes longer than the orchestrator's grace period
|
||||
// it SIGKILLs us, which is the correct outer bound.
|
||||
drain_write_barrier(&write_barrier);
|
||||
let drained = drain_connections(&conn_count);
|
||||
let flush_started = std::time::Instant::now();
|
||||
let flushed = registry.flush_all();
|
||||
|
|
|
|||
|
|
@ -28,6 +28,31 @@ fn the_write_barrier_admits_writes_until_shutdown_closes_it() {
|
|||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_drain_returns_only_when_no_write_is_in_flight() {
|
||||
// The process must never flush over a live writer: there is no deadline at
|
||||
// which losing acked work becomes acceptable, so the wait is unbounded and
|
||||
// the orchestrator's grace period is the outer bound.
|
||||
use crate::web_canvas_server::tenant::WriteBarrier;
|
||||
let barrier = std::sync::Arc::new(WriteBarrier::default());
|
||||
barrier.close();
|
||||
|
||||
let held = barrier.enter_for_test();
|
||||
let waiting = std::sync::Arc::clone(&barrier);
|
||||
let handle = std::thread::spawn(move || {
|
||||
super::super::drain_write_barrier_for_test(&waiting);
|
||||
});
|
||||
// Still blocked while the pass is held.
|
||||
std::thread::sleep(std::time::Duration::from_millis(120));
|
||||
assert!(!handle.is_finished(), "the drain must wait for the writer");
|
||||
|
||||
drop(held);
|
||||
handle
|
||||
.join()
|
||||
.expect("the drain returns once the writer finishes");
|
||||
assert_eq!(barrier.active(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_held_pass_keeps_the_barrier_busy_so_the_flush_waits() {
|
||||
// The window this closes: a worker past the connection drain, about to
|
||||
|
|
|
|||
|
|
@ -357,6 +357,14 @@ impl WriteBarrier {
|
|||
Some(WritePass { barrier: self })
|
||||
}
|
||||
|
||||
/// Enter the write path regardless of the closed flag. Tests only — it is
|
||||
/// how a test stands in for a writer that was admitted before the close.
|
||||
#[cfg(test)]
|
||||
pub fn enter_for_test(&self) -> WritePass<'_> {
|
||||
self.active.fetch_add(1, Ordering::AcqRel);
|
||||
WritePass { barrier: self }
|
||||
}
|
||||
|
||||
/// Stop admitting writes. Idempotent.
|
||||
pub fn close(&self) {
|
||||
self.closed.store(true, Ordering::Release);
|
||||
|
|
|
|||
|
|
@ -158,16 +158,49 @@ fn parse_chat_attachments(value: Option<&Value>) -> Vec<ChatAttachment> {
|
|||
.collect()
|
||||
}
|
||||
|
||||
/// Everything an AI turn needs to commit to the canvas: the document
|
||||
/// authority, the stream that announces a change, and the shutdown admission
|
||||
/// that decides whether a commit may happen at all.
|
||||
///
|
||||
/// Bundled because the three always travel together — and separating them is
|
||||
/// how `/api/ai/standard` came to have commit points with no admission.
|
||||
#[derive(Clone, Copy)]
|
||||
pub(crate) struct CanvasWriteTarget<'a> {
|
||||
pub(crate) state: &'a Mutex<WebCanvasState>,
|
||||
pub(crate) hub: &'a SseHub,
|
||||
pub(crate) write_barrier: Option<&'a crate::web_canvas_server::WriteBarrier>,
|
||||
}
|
||||
|
||||
/// Admission for one document commit on the AI path.
|
||||
///
|
||||
/// The conversation itself never holds the shutdown barrier — a model turn can
|
||||
/// run for minutes and would block every stop. The pass is taken only for the
|
||||
/// instant a commit touches the document, and refused once shutdown has closed
|
||||
/// the barrier: the reply still streams back, but the document is left alone
|
||||
/// and the caller is told the turn was not applied.
|
||||
fn admit_document_write(
|
||||
barrier: Option<&crate::web_canvas_server::WriteBarrier>,
|
||||
) -> Result<Option<crate::web_canvas_server::WritePass<'_>>, WebChatStandardError> {
|
||||
let Some(barrier) = barrier else {
|
||||
return Ok(None); // local/managed: no flush to protect
|
||||
};
|
||||
barrier
|
||||
.enter()
|
||||
.map(Some)
|
||||
.ok_or(WebChatStandardError::ShuttingDown)
|
||||
}
|
||||
|
||||
pub fn stream_standard_turn<W: Write>(
|
||||
out: &mut W,
|
||||
req: WebStandardTurnRequest,
|
||||
state: &Mutex<WebCanvasState>,
|
||||
hub: &SseHub,
|
||||
write_barrier: Option<&crate::web_canvas_server::WriteBarrier>,
|
||||
cors_origin: Option<&str>,
|
||||
) -> std::io::Result<()> {
|
||||
crate::ai_proxy::write_sse_headers(out, cors_origin)?;
|
||||
|
||||
let mut snapshot = match apply_request_snapshot(&req, state, hub) {
|
||||
let mut snapshot = match apply_request_snapshot(&req, state, hub, write_barrier) {
|
||||
Ok(snapshot) => snapshot,
|
||||
Err(error) => {
|
||||
// `write_error_event` feeds `op-ai`'s `ChatDelta::Error(String)`
|
||||
|
|
@ -196,7 +229,13 @@ pub fn stream_standard_turn<W: Write>(
|
|||
op_editor_core::CollabEditSource::Ai,
|
||||
)
|
||||
.is_ok();
|
||||
if gated && clear_live_starter_frame_for_design(&mut guard).is_some() {
|
||||
// Also a document commit, so it needs the same instant of
|
||||
// admission; a closed barrier simply skips the clear.
|
||||
let starter_clear_pass = admit_document_write(write_barrier).ok();
|
||||
if gated
|
||||
&& starter_clear_pass.is_some()
|
||||
&& clear_live_starter_frame_for_design(&mut guard).is_some()
|
||||
{
|
||||
snapshot = guard.editor.clone();
|
||||
Some(guard.sse_tick())
|
||||
} else {
|
||||
|
|
@ -250,11 +289,27 @@ pub fn stream_standard_turn<W: Write>(
|
|||
}
|
||||
crate::chat_intent::DesignIntent::Modify => {
|
||||
let plan = modify_plan.expect("route checked has_modify_plan");
|
||||
stream_modify_route(out, plan, design_provider.as_ref(), state, hub)
|
||||
}
|
||||
crate::chat_intent::DesignIntent::New => {
|
||||
stream_new_design_route(out, req, snapshot, design_provider, state, hub, model)
|
||||
stream_modify_route(
|
||||
out,
|
||||
plan,
|
||||
design_provider.as_ref(),
|
||||
state,
|
||||
hub,
|
||||
write_barrier,
|
||||
)
|
||||
}
|
||||
crate::chat_intent::DesignIntent::New => stream_new_design_route(
|
||||
out,
|
||||
req,
|
||||
snapshot,
|
||||
design_provider,
|
||||
model,
|
||||
CanvasWriteTarget {
|
||||
state,
|
||||
hub,
|
||||
write_barrier,
|
||||
},
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -262,6 +317,7 @@ fn apply_request_snapshot(
|
|||
req: &WebStandardTurnRequest,
|
||||
state: &Mutex<WebCanvasState>,
|
||||
hub: &SseHub,
|
||||
write_barrier: Option<&crate::web_canvas_server::WriteBarrier>,
|
||||
) -> Result<EditorState, WebChatStandardError> {
|
||||
let mut broadcast_tick = None;
|
||||
let mut snapshot = {
|
||||
|
|
@ -291,6 +347,9 @@ fn apply_request_snapshot(
|
|||
op_editor_core::CollabEditSource::Ai,
|
||||
)
|
||||
.map_err(WebChatStandardError::CollabRefused)?;
|
||||
// The one instant this turn touches the document; refused
|
||||
// outright once shutdown has closed the barrier.
|
||||
let _write_pass = admit_document_write(write_barrier)?;
|
||||
guard.replace_document(loaded.value);
|
||||
broadcast_tick = Some(guard.sse_tick());
|
||||
}
|
||||
|
|
@ -418,6 +477,7 @@ fn stream_modify_route<W: Write>(
|
|||
provider: &dyn ChatProvider,
|
||||
state: &Mutex<WebCanvasState>,
|
||||
hub: &SseHub,
|
||||
write_barrier: Option<&crate::web_canvas_server::WriteBarrier>,
|
||||
) -> std::io::Result<()> {
|
||||
write_delta_event(out, STANDARD_MODIFY_STEP)?;
|
||||
let target_frame_ids = plan.target_frame_ids;
|
||||
|
|
@ -459,11 +519,18 @@ fn stream_modify_route<W: Write>(
|
|||
{
|
||||
(0, None)
|
||||
} else {
|
||||
let (count, mutated) = crate::chat_canvas_tools::apply_design_modification(
|
||||
&mut guard.editor,
|
||||
&nodes,
|
||||
&target_frame_ids,
|
||||
);
|
||||
// Shutting down: the reply still streams, the document is left
|
||||
// exactly as the flush will find it.
|
||||
let admitted = admit_document_write(write_barrier).ok();
|
||||
let (count, mutated) = if admitted.is_none() {
|
||||
(0, false)
|
||||
} else {
|
||||
crate::chat_canvas_tools::apply_design_modification(
|
||||
&mut guard.editor,
|
||||
&nodes,
|
||||
&target_frame_ids,
|
||||
)
|
||||
};
|
||||
let tick = if mutated {
|
||||
guard.version += 1;
|
||||
Some(guard.sse_tick())
|
||||
|
|
@ -507,9 +574,8 @@ fn stream_new_design_route<W: Write>(
|
|||
req: WebStandardTurnRequest,
|
||||
snapshot: EditorState,
|
||||
provider: Box<dyn ChatProvider>,
|
||||
state: &Mutex<WebCanvasState>,
|
||||
hub: &SseHub,
|
||||
model: Option<String>,
|
||||
target: CanvasWriteTarget<'_>,
|
||||
) -> std::io::Result<()> {
|
||||
let append_context = crate::chat_intent::detect_append_intent(&snapshot, &req.ai.user);
|
||||
let request = DesignRequest {
|
||||
|
|
@ -532,7 +598,7 @@ fn stream_new_design_route<W: Write>(
|
|||
// the user picked instead of needing a second key.
|
||||
let provider_arc: Arc<dyn ChatProvider> = Arc::from(provider);
|
||||
let llm = ChatProviderLlmClient::new(provider_arc.clone()).with_model(model.clone());
|
||||
let mut sink = WebDesignDocSink::new(state, hub, snapshot);
|
||||
let mut sink = WebDesignDocSink::new(target.state, target.hub, target.write_barrier, snapshot);
|
||||
let abort = AbortFlag::new();
|
||||
let pre_validator = LintPreValidator;
|
||||
|
||||
|
|
@ -621,12 +687,25 @@ fn stream_new_design_route<W: Write>(
|
|||
struct WebDesignDocSink<'a> {
|
||||
state: &'a Mutex<WebCanvasState>,
|
||||
hub: &'a SseHub,
|
||||
/// Admission for each generated command's commit. `None` for the local
|
||||
/// and managed daemons, which have no flush to protect.
|
||||
write_barrier: Option<&'a crate::web_canvas_server::WriteBarrier>,
|
||||
mirror: EditorState,
|
||||
}
|
||||
|
||||
impl<'a> WebDesignDocSink<'a> {
|
||||
fn new(state: &'a Mutex<WebCanvasState>, hub: &'a SseHub, mirror: EditorState) -> Self {
|
||||
Self { state, hub, mirror }
|
||||
fn new(
|
||||
state: &'a Mutex<WebCanvasState>,
|
||||
hub: &'a SseHub,
|
||||
write_barrier: Option<&'a crate::web_canvas_server::WriteBarrier>,
|
||||
mirror: EditorState,
|
||||
) -> Self {
|
||||
Self {
|
||||
state,
|
||||
hub,
|
||||
write_barrier,
|
||||
mirror,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -636,6 +715,12 @@ impl DocSink for WebDesignDocSink<'_> {
|
|||
}
|
||||
|
||||
fn apply(&mut self, cmd: EditorCommand) -> bool {
|
||||
// Each generated command is its own document commit, so each needs its
|
||||
// own instant of admission. A closed barrier acks `false`, which the
|
||||
// generator already treats as "not applied".
|
||||
let Ok(_write_pass) = admit_document_write(self.write_barrier) else {
|
||||
return false;
|
||||
};
|
||||
let (applied, tick, snapshot) = {
|
||||
let mut guard = self.state.lock().unwrap_or_else(|p| p.into_inner());
|
||||
// A refusal and a no-op both ack `false` to the generator, which is
|
||||
|
|
|
|||
|
|
@ -32,6 +32,9 @@ use std::fmt;
|
|||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub(crate) enum WebChatStandardError {
|
||||
/// The daemon is shutting down and will not durably accept a document
|
||||
/// write. The conversation still answers; the document is untouched.
|
||||
ShuttingDown,
|
||||
/// The request-scoped credential names a different model than the turn
|
||||
/// it rides with. Refused rather than reconciled — the credential is
|
||||
/// browser-supplied and the mismatch is unresolvable.
|
||||
|
|
@ -60,6 +63,9 @@ pub(crate) enum WebChatStandardError {
|
|||
impl fmt::Display for WebChatStandardError {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
WebChatStandardError::ShuttingDown => f.write_str(
|
||||
"this daemon is stopping; the reply was produced but the document was not changed",
|
||||
),
|
||||
WebChatStandardError::TransientModelMismatch => {
|
||||
f.write_str("transient credential model does not match the request")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -47,7 +47,7 @@ fn transient_credential_is_applied_only_to_the_turn_snapshot() {
|
|||
let state = Mutex::new(WebCanvasState::new(EditorState::new(), 3100));
|
||||
|
||||
let snapshot =
|
||||
apply_request_snapshot(&req, &state, &SseHub::default()).expect("request snapshot");
|
||||
apply_request_snapshot(&req, &state, &SseHub::default(), None).expect("request snapshot");
|
||||
|
||||
assert_eq!(
|
||||
snapshot.editor_ui.agent_settings.builtin_agents[0].api_key,
|
||||
|
|
@ -80,7 +80,7 @@ fn transient_credential_model_must_match_the_requested_model() {
|
|||
let req = parse_standard_turn_body(&body.to_string()).expect("credential shape parses");
|
||||
let state = Mutex::new(WebCanvasState::new(EditorState::new(), 3100));
|
||||
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default())
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default(), None)
|
||||
.expect_err("mismatched credential must be rejected")
|
||||
.to_string();
|
||||
|
||||
|
|
@ -107,7 +107,7 @@ fn browser_only_demo_rejects_custom_or_loopback_transient_endpoints_without_muta
|
|||
));
|
||||
let before = crate::settings_io::fingerprint(&state.lock().unwrap().editor);
|
||||
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default())
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default(), None)
|
||||
.expect_err("public demo must reject custom endpoint")
|
||||
.to_string();
|
||||
|
||||
|
|
@ -129,7 +129,7 @@ fn server_persistence_does_not_allow_a_reserved_transient_endpoint() {
|
|||
crate::web_credential_policy::WebCredentialPersistence::Server,
|
||||
));
|
||||
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default())
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default(), None)
|
||||
.expect_err("persistence must not authorize a reserved provider endpoint")
|
||||
.to_string();
|
||||
assert!(error.contains("endpoint"), "unexpected error: {error}");
|
||||
|
|
@ -275,7 +275,7 @@ fn request_snapshot_restores_editor_meta_before_design_response_updates() {
|
|||
let state = Mutex::new(WebCanvasState::new(EditorState::new(), 3100));
|
||||
|
||||
let snapshot =
|
||||
apply_request_snapshot(&req, &state, &SseHub::default()).expect("request snapshot");
|
||||
apply_request_snapshot(&req, &state, &SseHub::default(), None).expect("request snapshot");
|
||||
|
||||
assert_eq!(snapshot.ui.active_page_index, 1);
|
||||
assert!(snapshot.editor_ui.preserve_authored_geometry);
|
||||
|
|
@ -290,7 +290,7 @@ fn design_doc_sink_applies_and_bumps_version() {
|
|||
let hub = SseHub::default();
|
||||
let sub = hub.subscribe();
|
||||
let mirror = state.lock().unwrap().editor.clone();
|
||||
let mut sink = WebDesignDocSink::new(&state, &hub, mirror);
|
||||
let mut sink = WebDesignDocSink::new(&state, &hub, None, mirror);
|
||||
|
||||
assert!(sink.apply(EditorCommand::InsertNode {
|
||||
kind: "rect".into(),
|
||||
|
|
@ -403,3 +403,70 @@ fn web_progress_label_subtask_retry_format() {
|
|||
" ▸ retry #2: zero nodes generated"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_closed_write_barrier_leaves_the_document_untouched() {
|
||||
// `/api/ai/standard` used to be dispatched before the write admission, so
|
||||
// its document commits ran during shutdown — after the flush had already
|
||||
// snapshotted the document.
|
||||
use crate::web_canvas_server::WriteBarrier;
|
||||
|
||||
let barrier = WriteBarrier::default();
|
||||
barrier.close();
|
||||
|
||||
let state = Mutex::new(WebCanvasState::new(EditorState::new(), 3100));
|
||||
let before = state
|
||||
.lock()
|
||||
.unwrap_or_else(|p| p.into_inner())
|
||||
.document_version_for_test();
|
||||
|
||||
let req = parse_standard_turn_body(&standard_body_with_document()).expect("request parses");
|
||||
let error = apply_request_snapshot(&req, &state, &SseHub::default(), Some(&barrier))
|
||||
.expect_err("a closed barrier must refuse the commit");
|
||||
assert!(
|
||||
matches!(error, WebChatStandardError::ShuttingDown),
|
||||
"{error:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
state
|
||||
.lock()
|
||||
.unwrap_or_else(|p| p.into_inner())
|
||||
.document_version_for_test(),
|
||||
before,
|
||||
"the refused turn must not have changed the document"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_open_write_barrier_applies_the_turn_normally() {
|
||||
use crate::web_canvas_server::WriteBarrier;
|
||||
let barrier = WriteBarrier::default();
|
||||
let state = Mutex::new(WebCanvasState::new(EditorState::new(), 3100));
|
||||
let req = parse_standard_turn_body(&standard_body_with_document()).expect("request parses");
|
||||
apply_request_snapshot(&req, &state, &SseHub::default(), Some(&barrier))
|
||||
.expect("an open barrier admits the commit");
|
||||
assert!(
|
||||
state
|
||||
.lock()
|
||||
.unwrap_or_else(|p| p.into_inner())
|
||||
.document_version_for_test()
|
||||
> 0,
|
||||
"the turn must have applied"
|
||||
);
|
||||
}
|
||||
|
||||
/// A standard-turn body that carries a document, so `apply_request_snapshot`
|
||||
/// reaches its `replace_document` commit.
|
||||
fn standard_body_with_document() -> String {
|
||||
serde_json::json!({
|
||||
"message": "hello",
|
||||
"document": {
|
||||
"version": "1.0.0",
|
||||
"children": [{
|
||||
"id": "n1", "type": "rectangle", "name": "from-turn",
|
||||
"x": 0, "y": 0, "width": 4, "height": 4,
|
||||
}],
|
||||
},
|
||||
})
|
||||
.to_string()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -186,25 +186,23 @@ thread_local! {
|
|||
/// request, so it is safe to run ahead of the postMessage bridge init gate.
|
||||
pub(crate) fn reset() {
|
||||
SYNC_STATE.with(|state| *state.borrow_mut() = CredentialSyncState::default());
|
||||
// Any callback already in flight was issued for the PREVIOUS account.
|
||||
// Clearing the queue cannot recall it, so the epoch is what makes it inert
|
||||
// when it lands — see `identity_is_current`.
|
||||
ISSUED_EPOCH.with(|epoch| epoch.set(crate::identity_epoch::epoch()));
|
||||
}
|
||||
|
||||
thread_local! {
|
||||
/// The identity epoch the in-flight requests were issued under.
|
||||
static ISSUED_EPOCH: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
|
||||
}
|
||||
|
||||
/// Whether a completing callback still belongs to the account that issued it.
|
||||
/// Whether a callback issued at `issued_epoch` still belongs to the account
|
||||
/// that issued it.
|
||||
///
|
||||
/// An XHR cannot be un-issued: `reset` empties the queue, but a POST already
|
||||
/// on the wire will still complete and its callback would fold the previous
|
||||
/// account's result — a success, a retry schedule, an error banner — into the
|
||||
/// new account's state. Comparing epochs at completion is what discards it.
|
||||
fn identity_is_current() -> bool {
|
||||
ISSUED_EPOCH.with(std::cell::Cell::get) == crate::identity_epoch::epoch()
|
||||
/// new account's state.
|
||||
///
|
||||
/// The epoch MUST be captured when the request is issued and compared here as
|
||||
/// an immutable snapshot. An earlier version read a shared `ISSUED_EPOCH` cell
|
||||
/// that `reset` itself rewrote, so by the time a stale callback landed the
|
||||
/// cell already held the NEW epoch and the comparison passed — the guard let
|
||||
/// through exactly what it existed to stop.
|
||||
fn identity_is_current(issued_epoch: u64) -> bool {
|
||||
issued_epoch == crate::identity_epoch::epoch()
|
||||
}
|
||||
|
||||
/// Begin credential-policy discovery against the daemon. This issues a daemon
|
||||
|
|
@ -259,8 +257,11 @@ fn request_repaint() {
|
|||
|
||||
fn fetch_policy() {
|
||||
let base = crate::daemon_base::daemon_base();
|
||||
let on_response: Rc<dyn Fn(u16, String)> = Rc::new(|status, body| {
|
||||
if !identity_is_current() {
|
||||
// Snapshotted here, moved into the closure: the value must describe THIS
|
||||
// request, not whatever the tab's identity is when the reply lands.
|
||||
let issued_epoch = crate::identity_epoch::epoch();
|
||||
let on_response: Rc<dyn Fn(u16, String)> = Rc::new(move |status, body| {
|
||||
if !identity_is_current(issued_epoch) {
|
||||
return; // issued for a previous account
|
||||
}
|
||||
let policy = parse_policy_response(status, &body);
|
||||
|
|
@ -284,8 +285,9 @@ fn fetch_policy() {
|
|||
|
||||
fn post_credentials(body: &str) {
|
||||
let base = crate::daemon_base::daemon_base();
|
||||
let on_response: Rc<dyn Fn(u16, String)> = Rc::new(|status, _| {
|
||||
if !identity_is_current() {
|
||||
let issued_epoch = crate::identity_epoch::epoch();
|
||||
let on_response: Rc<dyn Fn(u16, String)> = Rc::new(move |status, _| {
|
||||
if !identity_is_current(issued_epoch) {
|
||||
// Issued for a previous account: its result must not become the
|
||||
// new account's success, retry schedule or error banner.
|
||||
return;
|
||||
|
|
@ -316,8 +318,9 @@ fn schedule_retry(delay_ms: i32, generation: u64) {
|
|||
{
|
||||
use wasm_bindgen::JsCast;
|
||||
|
||||
let issued_epoch = crate::identity_epoch::epoch();
|
||||
let callback = wasm_bindgen::closure::Closure::<dyn FnMut()>::once_into_js(move || {
|
||||
if !identity_is_current() {
|
||||
if !identity_is_current(issued_epoch) {
|
||||
// A timer armed for a previous account: firing it would resend
|
||||
// that account's credential payload into the new tenant.
|
||||
return;
|
||||
|
|
@ -654,3 +657,54 @@ mod tests {
|
|||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod identity_epoch_guard_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn a_callback_issued_under_an_older_epoch_is_refused() {
|
||||
// The bug this pins: the guard used to read a shared cell that `reset`
|
||||
// itself rewrote, so a stale callback compared "new == new" and was
|
||||
// let through — the guard admitted exactly what it existed to stop.
|
||||
crate::identity_epoch::reset_for_test();
|
||||
crate::identity_epoch::observe_subject(Some("alice"));
|
||||
let issued_under_alice = crate::identity_epoch::epoch();
|
||||
assert!(identity_is_current(issued_under_alice));
|
||||
|
||||
// The tab switches to B, which resets the queue.
|
||||
crate::identity_epoch::observe_subject(Some("bob"));
|
||||
reset();
|
||||
|
||||
assert!(
|
||||
!identity_is_current(issued_under_alice),
|
||||
"A's in-flight callback must stay refused after the switch"
|
||||
);
|
||||
assert!(
|
||||
identity_is_current(crate::identity_epoch::epoch()),
|
||||
"B's own requests must still be accepted"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_refused_callback_leaves_no_state_behind_and_sync_can_restart() {
|
||||
// The FirstIdentified race: a policy fetch issued at mount under the
|
||||
// anonymous epoch is refused when it lands, and `policy_in_flight`
|
||||
// would stay set — after which every `changed()` is a no-op and the
|
||||
// account can never upload again.
|
||||
crate::identity_epoch::reset_for_test();
|
||||
let mut state = CredentialSyncState::default();
|
||||
assert_eq!(state.begin_policy_check(), SyncAction::FetchPolicy);
|
||||
// …its reply is discarded by the guard, so nothing resolves it.
|
||||
|
||||
// The partition reload resets and restarts, which is what unwedges it.
|
||||
let mut restarted = CredentialSyncState::default();
|
||||
assert_eq!(restarted.begin_policy_check(), SyncAction::FetchPolicy);
|
||||
assert_eq!(restarted.resolve_policy(Some(true)), SyncAction::None);
|
||||
assert_ne!(
|
||||
restarted.changed("{}".to_string()),
|
||||
SyncAction::None,
|
||||
"after the restart a credential edit must reach the daemon again"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -88,6 +88,15 @@ pub(crate) fn reload_for_active_partition<C: crate::repaint_ctx::RepaintContext>
|
|||
// account's state would make the very next comparison report the whole
|
||||
// partition as a change and write it back under the wrong key.
|
||||
context.reset_persistence_baselines(&load);
|
||||
// Restart credential sync for the partition now in force.
|
||||
//
|
||||
// Needed on FirstIdentified too, which does NOT go through the identity
|
||||
// reset: a policy fetch issued at mount under the anonymous epoch is
|
||||
// refused by the epoch guard when it lands, and `policy_in_flight` would
|
||||
// stay set forever — after which every `changed()` returns `SyncAction::
|
||||
// None` and the account can never upload a credential again.
|
||||
crate::web_credential_sync::reset();
|
||||
crate::web_credential_sync::start();
|
||||
context.host_mut().mark_editor_state_dirty();
|
||||
let _ = context.repaint();
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue