From 055db1895faf3bd61342ff231e6d93f5bc76d92d Mon Sep 17 00:00:00 2001 From: Kayshen-X Date: Sat, 8 Aug 2026 22:35:42 +0800 Subject: [PATCH] fix(web): request-time epoch snapshots, AI write admission, drain without ceiling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- .../src/runtime/conflict_tests.rs | 6 +- .../op-collab-host/src/runtime/local_edit.rs | 28 ++++- crates/op-collab-host/src/runtime/tests.rs | 16 +-- .../op-host-services/src/web_canvas_server.rs | 2 +- .../src/web_canvas_server/collab_state.rs | 17 ++- .../web_canvas_server/collab_state_tests.rs | 63 ++++++---- .../web_canvas_server/connection_ai_routes.rs | 3 + .../src/web_canvas_server/online_run_loop.rs | 63 +++++----- .../online_shutdown_tests.rs | 25 ++++ .../src/web_canvas_server/tenant.rs | 8 ++ .../op-host-services/src/web_chat_standard.rs | 117 +++++++++++++++--- .../src/web_chat_standard_error.rs | 6 + .../src/web_chat_standard_tests.rs | 79 +++++++++++- crates/op-host-web/src/web_credential_sync.rs | 90 +++++++++++--- crates/op-host-web/src/web_settings.rs | 9 ++ 15 files changed, 423 insertions(+), 109 deletions(-) diff --git a/crates/op-collab-host/src/runtime/conflict_tests.rs b/crates/op-collab-host/src/runtime/conflict_tests.rs index 1bda7623c..56fe1c41e 100644 --- a/crates/op-collab-host/src/runtime/conflict_tests.rs +++ b/crates/op-collab-host/src/runtime/conflict_tests.rs @@ -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. diff --git a/crates/op-collab-host/src/runtime/local_edit.rs b/crates/op-collab-host/src/runtime/local_edit.rs index b5c1e1222..4dd67ef95 100644 --- a/crates/op-collab-host/src/runtime/local_edit.rs +++ b/crates/op-collab-host/src/runtime/local_edit.rs @@ -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, + } } } } diff --git a/crates/op-collab-host/src/runtime/tests.rs b/crates/op-collab-host/src/runtime/tests.rs index 7ac258506..32ccc20d4 100644 --- a/crates/op-collab-host/src/runtime/tests.rs +++ b/crates/op-collab-host/src/runtime/tests.rs @@ -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 { diff --git a/crates/op-host-services/src/web_canvas_server.rs b/crates/op-host-services/src/web_canvas_server.rs index d06153b9a..2c833b116 100644 --- a/crates/op-host-services/src/web_canvas_server.rs +++ b/crates/op-host-services/src/web_canvas_server.rs @@ -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, }; diff --git a/crates/op-host-services/src/web_canvas_server/collab_state.rs b/crates/op-host-services/src/web_canvas_server/collab_state.rs index d2c8232d9..0b71e701e 100644 --- a/crates/op-host-services/src/web_canvas_server/collab_state.rs +++ b/crates/op-host-services/src/web_canvas_server/collab_state.rs @@ -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 } } diff --git a/crates/op-host-services/src/web_canvas_server/collab_state_tests.rs b/crates/op-host-services/src/web_canvas_server/collab_state_tests.rs index 62928bc78..0778b6407 100644 --- a/crates/op-host-services/src/web_canvas_server/collab_state_tests.rs +++ b/crates/op-host-services/src/web_canvas_server/collab_state_tests.rs @@ -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() +} diff --git a/crates/op-host-services/src/web_canvas_server/connection_ai_routes.rs b/crates/op-host-services/src/web_canvas_server/connection_ai_routes.rs index 819938e98..8a6534726 100644 --- a/crates/op-host-services/src/web_canvas_server/connection_ai_routes.rs +++ b/crates/op-host-services/src/web_canvas_server/connection_ai_routes.rs @@ -92,6 +92,9 @@ pub(super) fn serve_ai_route( 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}"))) diff --git a/crates/op-host-services/src/web_canvas_server/online_run_loop.rs b/crates/op-host-services/src/web_canvas_server/online_run_loop.rs index 14fdb55da..eba38f875 100644 --- a/crates/op-host-services/src/web_canvas_server/online_run_loop.rs +++ b/crates/op-host-services/src/web_canvas_server/online_run_loop.rs @@ -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) -> 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) -> 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) { + drain_write_barrier(barrier); +} + +fn drain_write_barrier(barrier: &Arc) { 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(); diff --git a/crates/op-host-services/src/web_canvas_server/online_shutdown_tests.rs b/crates/op-host-services/src/web_canvas_server/online_shutdown_tests.rs index 1064e3c13..5f6715038 100644 --- a/crates/op-host-services/src/web_canvas_server/online_shutdown_tests.rs +++ b/crates/op-host-services/src/web_canvas_server/online_shutdown_tests.rs @@ -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 diff --git a/crates/op-host-services/src/web_canvas_server/tenant.rs b/crates/op-host-services/src/web_canvas_server/tenant.rs index 6a6c12cde..8c2bd7279 100644 --- a/crates/op-host-services/src/web_canvas_server/tenant.rs +++ b/crates/op-host-services/src/web_canvas_server/tenant.rs @@ -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); diff --git a/crates/op-host-services/src/web_chat_standard.rs b/crates/op-host-services/src/web_chat_standard.rs index e1e330c8d..a66a945b6 100644 --- a/crates/op-host-services/src/web_chat_standard.rs +++ b/crates/op-host-services/src/web_chat_standard.rs @@ -158,16 +158,49 @@ fn parse_chat_attachments(value: Option<&Value>) -> Vec { .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, + 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>, 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( out: &mut W, req: WebStandardTurnRequest, state: &Mutex, 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( 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( } 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, hub: &SseHub, + write_barrier: Option<&crate::web_canvas_server::WriteBarrier>, ) -> Result { 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( provider: &dyn ChatProvider, state: &Mutex, 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( { (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( req: WebStandardTurnRequest, snapshot: EditorState, provider: Box, - state: &Mutex, - hub: &SseHub, model: Option, + 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( // the user picked instead of needing a second key. let provider_arc: Arc = 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( struct WebDesignDocSink<'a> { state: &'a Mutex, 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, hub: &'a SseHub, mirror: EditorState) -> Self { - Self { state, hub, mirror } + fn new( + state: &'a Mutex, + 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 diff --git a/crates/op-host-services/src/web_chat_standard_error.rs b/crates/op-host-services/src/web_chat_standard_error.rs index 1739cf4ba..9d43f55c7 100644 --- a/crates/op-host-services/src/web_chat_standard_error.rs +++ b/crates/op-host-services/src/web_chat_standard_error.rs @@ -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") } diff --git a/crates/op-host-services/src/web_chat_standard_tests.rs b/crates/op-host-services/src/web_chat_standard_tests.rs index 8c7b135ea..64bb9c77e 100644 --- a/crates/op-host-services/src/web_chat_standard_tests.rs +++ b/crates/op-host-services/src/web_chat_standard_tests.rs @@ -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() +} diff --git a/crates/op-host-web/src/web_credential_sync.rs b/crates/op-host-web/src/web_credential_sync.rs index bc7e5d99b..cfdafca2d 100644 --- a/crates/op-host-web/src/web_credential_sync.rs +++ b/crates/op-host-web/src/web_credential_sync.rs @@ -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 = 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 = 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 = 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 = Rc::new(|status, _| { - if !identity_is_current() { + let issued_epoch = crate::identity_epoch::epoch(); + let on_response: Rc = 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::::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" + ); + } +} diff --git a/crates/op-host-web/src/web_settings.rs b/crates/op-host-web/src/web_settings.rs index 569a50fe5..d551ae6e8 100644 --- a/crates/op-host-web/src/web_settings.rs +++ b/crates/op-host-web/src/web_settings.rs @@ -88,6 +88,15 @@ pub(crate) fn reload_for_active_partition // 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(); }