feat(web): wire browser collab UI to the daemon /api/collab service
collab_sync drives the already-complete web collab surface against the PR3 daemon routes, mirroring web_auth_sync: - adaptive re-arming timer (400ms idle / 150ms in session); collabSeq rides the existing version probe and never triggers a document pull - pending panel actions post as CollabActionWire with single-flight and 409 collab-busy retry; presence uploads throttle to 100ms - session-time pushes that would mint collaboration-invalid local node ids fail closed until the daemon echoes them (namespace handshake is the PR5+ follow-up); mount treats the daemon's sync-reset 409 as completion instead of a wasted retry
This commit is contained in:
parent
08a26548fe
commit
ebc681881d
|
|
@ -252,6 +252,18 @@ impl CollabStateWire {
|
|||
}
|
||||
}
|
||||
|
||||
/// Read the collaboration sequence out of a `GET /api/mcp/version` body.
|
||||
///
|
||||
/// The daemon answers `{"version":N,"collabSeq":M}` on that one probe, so a
|
||||
/// client polls once and routes each number to the loop that cares: `version`
|
||||
/// drives the document fetch, `collabSeq` drives this projection. A body
|
||||
/// without the field — an older daemon — reads as `None`, which simply leaves
|
||||
/// collaboration idle.
|
||||
pub fn parse_collab_seq_probe(body: &str) -> Option<u64> {
|
||||
let value: serde_json::Value = serde_json::from_str(body).ok()?;
|
||||
value.get("collabSeq")?.as_u64()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "collab_wire/tests.rs"]
|
||||
mod tests;
|
||||
|
|
|
|||
|
|
@ -218,3 +218,20 @@ fn local_presence_wire_defaults_its_reserved_client_id() {
|
|||
let json = serde_json::to_string(&decoded).expect("encodes");
|
||||
assert!(!json.contains("clientId"), "{json}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_version_probe_yields_the_collaboration_sequence_separately() {
|
||||
assert_eq!(
|
||||
super::parse_collab_seq_probe(r#"{"version":7,"collabSeq":42}"#),
|
||||
Some(42)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_daemon_without_the_collaboration_field_reads_as_absent() {
|
||||
// An older daemon answers only `{"version":N}`; collaboration then simply
|
||||
// stays idle instead of the client mis-reading the document version as a
|
||||
// projection sequence.
|
||||
assert_eq!(super::parse_collab_seq_probe(r#"{"version":7}"#), None);
|
||||
assert_eq!(super::parse_collab_seq_probe("not json"), None);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -258,13 +258,22 @@ pub(super) fn start_bootstrap_reset(
|
|||
let on_reset: std::rc::Rc<dyn Fn(u16, String)> = {
|
||||
let complete = complete.clone();
|
||||
let base = base.clone();
|
||||
std::rc::Rc::new(move |_status: u16, body: String| {
|
||||
std::rc::Rc::new(move |status: u16, body: String| {
|
||||
// The daemon answers `{"ok":true,...}` for both a fresh reset and a
|
||||
// peer-skipped one (`"skipped":true`) — either is completion.
|
||||
if body.contains("\"ok\":true") {
|
||||
complete();
|
||||
return;
|
||||
}
|
||||
// A live collaboration session owns this document, so the daemon
|
||||
// refuses to reset it (409 `collab-active`). That is a deliberate
|
||||
// answer, not a failure: resetting is exactly the wrong thing to do
|
||||
// to a document peers are editing. Retrying would only ask again
|
||||
// and warn about a daemon that is working correctly.
|
||||
if status == 409 {
|
||||
complete();
|
||||
return;
|
||||
}
|
||||
// Error body / empty (transport failure, or an XHR timeout ->
|
||||
// status 0 + empty body): retry once, then proceed.
|
||||
if retries_left > 0 {
|
||||
|
|
|
|||
|
|
@ -159,6 +159,11 @@ pub(super) async fn mount_ck(canvas_id: String) -> Result<(), JsValue> {
|
|||
// daemon restored (shared with the desktop GUI), then drive login
|
||||
// flows through the daemon's `/api/auth/*` proxy.
|
||||
crate::web_auth_sync::start(&inner);
|
||||
// Collaboration relay. Starts here rather than inside the sync-reset
|
||||
// completion because it neither reads nor writes the document — it drives
|
||||
// `editor_ui.collab` only — and the panel should report availability as
|
||||
// soon as the daemon can answer, not one bootstrap later.
|
||||
crate::collab_sync::start(&inner);
|
||||
// 4. Reset the daemon's transient sync document, THEN emit the managed
|
||||
// `ready` reply and start the live-sync ticks. The reset must complete
|
||||
// FIRST for two reasons:
|
||||
|
|
|
|||
510
crates/op-host-web/src/collab_sync.rs
Normal file
510
crates/op-host-web/src/collab_sync.rs
Normal file
|
|
@ -0,0 +1,510 @@
|
|||
//! Drive the daemon's collaboration runtime from the web shell.
|
||||
//!
|
||||
//! The wasm bundle carries the collaboration *panel* but no transport: a relay
|
||||
//! session needs device credentials only a native process can hold, so the
|
||||
//! session lives in the daemon and this module is the wire between them. It is
|
||||
//! the collaboration twin of [`crate::web_auth_sync`] and follows its shape —
|
||||
//! one interval, `thread_local` latches, every write through the shared
|
||||
//! `RepaintContext`.
|
||||
//!
|
||||
//! Three jobs per tick, in this order:
|
||||
//!
|
||||
//! 1. drain the panel's queued [`CollabUiAction`] to `POST /api/collab/action`;
|
||||
//! 2. pull `GET /api/collab/state` when the projection moved, and install it;
|
||||
//! 3. publish the local cursor to `POST /api/collab/presence`.
|
||||
//!
|
||||
//! ## Two counters, two meanings
|
||||
//!
|
||||
//! The daemon publishes `documentRevision` and `collabSeq` separately.
|
||||
//! [`crate::live_sync_glue`]'s version poll already runs every 400 ms, so it
|
||||
//! hands the observed `collabSeq` to [`note_version_probe`] instead of this
|
||||
//! module opening a second probe. A `collabSeq` bump pulls *state*; it must
|
||||
//! never pull the document, or every remote cursor move would refetch the whole
|
||||
//! canvas.
|
||||
//!
|
||||
//! ## Availability
|
||||
//!
|
||||
//! Only ever what the daemon projected. Before the first successful `/state`
|
||||
//! the panel keeps its default `Unavailable`, so a shell that cannot reach a
|
||||
//! daemon offers nothing rather than offering a session it cannot start.
|
||||
|
||||
use std::cell::{Cell, RefCell};
|
||||
use std::collections::HashSet;
|
||||
use std::rc::Rc;
|
||||
|
||||
use op_editor_core::collab_wire::{
|
||||
CollabActionWire, CollabLocalPresenceWire, CollabPointWire, CollabStateWire,
|
||||
};
|
||||
use op_editor_core::{collab_routes, CollabConnectionPhase, CollabUiAction};
|
||||
|
||||
use crate::live_sync;
|
||||
use crate::repaint_ctx::RepaintContext;
|
||||
|
||||
/// Tick cadence with no session — matched to the document version poll, since
|
||||
/// with nothing running there is nothing to be responsive about.
|
||||
const IDLE_TICK_MS: i32 = 400;
|
||||
/// Tick cadence once a session exists. Presence and admission prompts are the
|
||||
/// interactive parts of collaboration; 150 ms keeps them from feeling laggy
|
||||
/// without approaching a per-frame request rate.
|
||||
const SESSION_TICK_MS: i32 = 150;
|
||||
/// Floor between presence posts. The daemon throttles again on its side to one
|
||||
/// frame interval, so anything under this is spent bytes for no visible gain.
|
||||
const PRESENCE_MIN_INTERVAL_MS: f64 = 100.0;
|
||||
|
||||
thread_local! {
|
||||
/// Latest `collabSeq` seen on the wire, from either probe.
|
||||
static OBSERVED_SEQ: Cell<u64> = const { Cell::new(0) };
|
||||
/// `collabSeq` of the projection currently installed in the panel.
|
||||
/// `None` until the first `/state` lands, which is what forces that first
|
||||
/// pull regardless of the counter.
|
||||
static APPLIED_SEQ: Cell<Option<u64>> = const { Cell::new(None) };
|
||||
/// One `/state` request in flight at a time.
|
||||
static STATE_BUSY: Cell<bool> = const { Cell::new(false) };
|
||||
/// An action already posted and not yet answered.
|
||||
static ACTION_BUSY: Cell<bool> = const { Cell::new(false) };
|
||||
/// An action the daemon refused with `collab-busy`, kept for the next tick.
|
||||
/// Losing it would silently drop something the user clicked.
|
||||
static ACTION_RETRY: RefCell<Option<CollabUiAction>> = const { RefCell::new(None) };
|
||||
/// `performance.now()` of the last presence post.
|
||||
static LAST_PRESENCE_MS: Cell<f64> = const { Cell::new(f64::NEG_INFINITY) };
|
||||
/// Last cursor actually sent, so an unmoved cursor costs nothing.
|
||||
static LAST_PRESENCE_POINT: Cell<Option<(f64, f64)>> = const { Cell::new(None) };
|
||||
/// Node ids present in the document the daemon last handed us.
|
||||
///
|
||||
/// The browser mints ids from a local sequential counter, which is exactly
|
||||
/// what an active session cannot accept — see [`push_blocked_by_session`].
|
||||
static DAEMON_NODE_IDS: RefCell<Option<HashSet<op_editor_core::NodeId>>> =
|
||||
const { RefCell::new(None) };
|
||||
}
|
||||
|
||||
/// Feed the shared version probe's `collabSeq` in.
|
||||
///
|
||||
/// `live_sync_glue` already polls `GET /api/mcp/version` every 400 ms and the
|
||||
/// daemon answers both counters there, so collaboration rides that request
|
||||
/// rather than opening a second one.
|
||||
pub(crate) fn note_version_probe(body: &str) {
|
||||
if let Some(seq) = op_editor_core::collab_wire::parse_collab_seq_probe(body) {
|
||||
OBSERVED_SEQ.set(seq);
|
||||
}
|
||||
}
|
||||
|
||||
/// Record the node ids of a document just applied from the daemon.
|
||||
pub(crate) fn note_daemon_document(state: &op_editor_core::EditorState) {
|
||||
DAEMON_NODE_IDS.with(|ids| *ids.borrow_mut() = Some(document_node_ids(state)));
|
||||
}
|
||||
|
||||
/// Whether a live session forbids pushing the current local document.
|
||||
///
|
||||
/// The browser has no owner-assigned id namespace: it mints `n<counter>` from a
|
||||
/// local sequential allocator, and two peers creating a node in the same moment
|
||||
/// would mint the same id. The collaboration protocol replays those ids
|
||||
/// verbatim, so a colliding pair silently forks the document — the failure this
|
||||
/// refuses to produce.
|
||||
///
|
||||
/// The check is deliberately whole-document rather than per-gesture: draw,
|
||||
/// duplicate, paste, group and import all mint through the same counter, and
|
||||
/// gating the single push covers every one of them without a guard at each
|
||||
/// call site. Any id the daemon has not seen blocks the push; edits to existing
|
||||
/// nodes (move, restyle, delete) carry no new ids and go through untouched.
|
||||
///
|
||||
/// The local node stays on screen until the next pull replaces it with the
|
||||
/// daemon's document, so the divergence is bounded and self-healing.
|
||||
pub(crate) fn push_blocked_by_session(state: &op_editor_core::EditorState) -> bool {
|
||||
if state.editor_ui.collab.phase != CollabConnectionPhase::Active {
|
||||
return false;
|
||||
}
|
||||
DAEMON_NODE_IDS.with(|known| {
|
||||
let known = known.borrow();
|
||||
let Some(known) = known.as_ref() else {
|
||||
// No daemon document seen yet in this session; refuse rather than
|
||||
// guess, since the pull that would settle it is one tick away.
|
||||
return true;
|
||||
};
|
||||
document_node_ids(state)
|
||||
.iter()
|
||||
.any(|id| !known.contains(id))
|
||||
})
|
||||
}
|
||||
|
||||
/// Node ids in a document, through `op-editor-core`'s own walker so pages and
|
||||
/// nodes share the one collision domain the allocator uses.
|
||||
fn document_node_ids(state: &op_editor_core::EditorState) -> HashSet<op_editor_core::NodeId> {
|
||||
op_editor_core::collect_document_ids(&state.doc)
|
||||
}
|
||||
|
||||
/// Wire the collaboration relay onto the mounted shell. Called once from mount.
|
||||
pub(crate) fn start<C: RepaintContext + 'static>(inner: &Rc<RefCell<C>>) {
|
||||
let base = crate::daemon_base::daemon_base();
|
||||
let inner = inner.clone();
|
||||
let tick: Rc<dyn Fn()> = Rc::new(move || {
|
||||
drain_pending_action(&inner, &base);
|
||||
maybe_pull_state(&inner, &base);
|
||||
maybe_push_presence(&inner, &base);
|
||||
});
|
||||
schedule(tick, IDLE_TICK_MS);
|
||||
}
|
||||
|
||||
/// Re-arming timeout rather than a fixed interval: the cadence has to follow
|
||||
/// the session phase, and re-reading it at each arming is what makes a session
|
||||
/// starting between two ticks speed the loop up immediately.
|
||||
fn schedule(tick: Rc<dyn Fn()>, delay_ms: i32) {
|
||||
use wasm_bindgen::closure::Closure;
|
||||
use wasm_bindgen::JsCast;
|
||||
|
||||
let Some(window) = web_sys::window() else {
|
||||
return;
|
||||
};
|
||||
let closure = Closure::once_into_js(move || {
|
||||
tick();
|
||||
schedule(tick, current_cadence_ms());
|
||||
});
|
||||
let _ = window
|
||||
.set_timeout_with_callback_and_timeout_and_arguments_0(closure.unchecked_ref(), delay_ms);
|
||||
}
|
||||
|
||||
/// Cadence for the next arming, from the phase the last projection installed.
|
||||
fn current_cadence_ms() -> i32 {
|
||||
match APPLIED_SEQ.get() {
|
||||
// Before the first projection there is nothing to be responsive to.
|
||||
None => IDLE_TICK_MS,
|
||||
Some(_) if SESSION_LIVE.get() => SESSION_TICK_MS,
|
||||
Some(_) => IDLE_TICK_MS,
|
||||
}
|
||||
}
|
||||
|
||||
thread_local! {
|
||||
/// Whether the installed projection is anything other than Idle. Read by
|
||||
/// the scheduler, which has no access to the host.
|
||||
static SESSION_LIVE: Cell<bool> = const { Cell::new(false) };
|
||||
}
|
||||
|
||||
fn drain_pending_action<C: RepaintContext + 'static>(inner: &Rc<RefCell<C>>, base: &str) {
|
||||
if ACTION_BUSY.get() {
|
||||
return;
|
||||
}
|
||||
let action = ACTION_RETRY
|
||||
.with(|slot| slot.borrow_mut().take())
|
||||
.or_else(|| {
|
||||
inner
|
||||
.try_borrow_mut()
|
||||
.ok()?
|
||||
.host_mut()
|
||||
.editor_state_mut()
|
||||
.editor_ui
|
||||
.collab
|
||||
.take_pending_action()
|
||||
});
|
||||
let Some(action) = action else {
|
||||
return;
|
||||
};
|
||||
let Some(wire) = wire_action(&action) else {
|
||||
// No wire form: nothing the daemon could do with it. Dropping is
|
||||
// correct — re-queueing would spin forever.
|
||||
return;
|
||||
};
|
||||
let Ok(body) = serde_json::to_string(&wire) else {
|
||||
return;
|
||||
};
|
||||
ACTION_BUSY.set(true);
|
||||
let retry_action = action.clone();
|
||||
let started = live_sync::post_json_with_status(
|
||||
&format!("{base}{}", collab_routes::ACTION),
|
||||
&body,
|
||||
Rc::new(move |status, response| {
|
||||
ACTION_BUSY.set(false);
|
||||
if status == 409 && response.contains("collab-busy") {
|
||||
// The daemon's single action slot was still full. Hold the
|
||||
// action for the next tick instead of dropping what the user
|
||||
// asked for.
|
||||
ACTION_RETRY.with(|slot| *slot.borrow_mut() = Some(retry_action.clone()));
|
||||
}
|
||||
}),
|
||||
);
|
||||
if !started {
|
||||
ACTION_BUSY.set(false);
|
||||
ACTION_RETRY.with(|slot| *slot.borrow_mut() = Some(action));
|
||||
}
|
||||
}
|
||||
|
||||
/// Map the panel's action onto its wire form.
|
||||
///
|
||||
/// Exhaustive by construction so a new `CollabUiAction` variant fails to
|
||||
/// compile here rather than silently becoming a no-op at runtime.
|
||||
fn wire_action(action: &CollabUiAction) -> Option<CollabActionWire> {
|
||||
let key = |k: &op_editor_core::CollabAdmissionRequestKey| k.as_str().to_owned();
|
||||
Some(match action {
|
||||
CollabUiAction::OpenCreate => CollabActionWire::OpenCreate,
|
||||
CollabUiAction::Start => CollabActionWire::Start,
|
||||
CollabUiAction::StartLan => CollabActionWire::StartLan,
|
||||
CollabUiAction::SetRelayRegion { region } => CollabActionWire::SetRelayRegion {
|
||||
region: (*region).into(),
|
||||
},
|
||||
CollabUiAction::OpenJoin => CollabActionWire::OpenJoin,
|
||||
CollabUiAction::BeginDiscovery => CollabActionWire::BeginDiscovery,
|
||||
CollabUiAction::JoinDiscovered { discovery_id } => CollabActionWire::JoinDiscovered {
|
||||
discovery_id: discovery_id.clone(),
|
||||
},
|
||||
CollabUiAction::JoinAddress { endpoint } => CollabActionWire::JoinAddress {
|
||||
endpoint: endpoint.clone(),
|
||||
},
|
||||
CollabUiAction::Cancel => CollabActionWire::Cancel,
|
||||
CollabUiAction::Retry => CollabActionWire::Retry,
|
||||
CollabUiAction::Leave => CollabActionWire::Leave,
|
||||
CollabUiAction::DiscardPending => CollabActionWire::DiscardPending,
|
||||
CollabUiAction::ReapplyDiscarded => CollabActionWire::ReapplyDiscarded,
|
||||
CollabUiAction::SaveAsFork => CollabActionWire::SaveAsFork,
|
||||
CollabUiAction::ApproveAdmissionEditor { request_key } => {
|
||||
CollabActionWire::ApproveAdmissionEditor {
|
||||
request_key: key(request_key),
|
||||
}
|
||||
}
|
||||
CollabUiAction::ApproveAdmissionViewer { request_key } => {
|
||||
CollabActionWire::ApproveAdmissionViewer {
|
||||
request_key: key(request_key),
|
||||
}
|
||||
}
|
||||
CollabUiAction::RejectAdmission { request_key } => CollabActionWire::RejectAdmission {
|
||||
request_key: key(request_key),
|
||||
},
|
||||
CollabUiAction::ConfirmOwnerIdentity { request_key } => {
|
||||
CollabActionWire::ConfirmOwnerIdentity {
|
||||
request_key: key(request_key),
|
||||
}
|
||||
}
|
||||
CollabUiAction::RejectOwnerIdentity { request_key } => {
|
||||
CollabActionWire::RejectOwnerIdentity {
|
||||
request_key: key(request_key),
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
fn maybe_pull_state<C: RepaintContext + 'static>(inner: &Rc<RefCell<C>>, base: &str) {
|
||||
if STATE_BUSY.get() {
|
||||
return;
|
||||
}
|
||||
let applied = APPLIED_SEQ.get();
|
||||
let observed = OBSERVED_SEQ.get();
|
||||
// Pull when the projection moved, when nothing has ever been pulled, or
|
||||
// continuously while a session runs — presence and admission prompts change
|
||||
// faster than the shared version probe reports.
|
||||
let due = applied != Some(observed) || applied.is_none() || SESSION_LIVE.get();
|
||||
if !due {
|
||||
return;
|
||||
}
|
||||
STATE_BUSY.set(true);
|
||||
let inner = inner.clone();
|
||||
let started = live_sync::get_with_status(
|
||||
&format!("{base}{}", collab_routes::STATE),
|
||||
Rc::new(move |status, body| {
|
||||
STATE_BUSY.set(false);
|
||||
if status != 200 {
|
||||
return;
|
||||
}
|
||||
let Ok(wire) = serde_json::from_str::<CollabStateWire>(&body) else {
|
||||
return;
|
||||
};
|
||||
apply_state(&inner, &wire);
|
||||
}),
|
||||
);
|
||||
if !started {
|
||||
STATE_BUSY.set(false);
|
||||
}
|
||||
}
|
||||
|
||||
fn apply_state<C: RepaintContext + 'static>(inner: &Rc<RefCell<C>>, wire: &CollabStateWire) {
|
||||
let Ok(mut context) = inner.try_borrow_mut() else {
|
||||
return;
|
||||
};
|
||||
let now_ms = now_ms();
|
||||
let state = context.host_mut().editor_state_mut();
|
||||
// Everything lands through the projection's own `set_*` API, which
|
||||
// re-runs each sanitising constructor — the wire is data, not state.
|
||||
wire.apply_to(&mut state.editor_ui.collab, now_ms);
|
||||
APPLIED_SEQ.set(Some(wire.collab_seq));
|
||||
SESSION_LIVE.set(state.editor_ui.collab.phase != CollabConnectionPhase::Idle);
|
||||
context.host_mut().mark_editor_state_dirty();
|
||||
let _ = context.repaint();
|
||||
}
|
||||
|
||||
fn maybe_push_presence<C: RepaintContext + 'static>(inner: &Rc<RefCell<C>>, base: &str) {
|
||||
if !SESSION_LIVE.get() {
|
||||
return;
|
||||
}
|
||||
let now = now_ms_f64();
|
||||
if now - LAST_PRESENCE_MS.get() < PRESENCE_MIN_INTERVAL_MS {
|
||||
return;
|
||||
}
|
||||
let Ok(context) = inner.try_borrow() else {
|
||||
return;
|
||||
};
|
||||
let (width, height) = context.viewport_size();
|
||||
let cursor = context.host().last_cursor_doc_point(width, height);
|
||||
drop(context);
|
||||
|
||||
if LAST_PRESENCE_POINT.get() == cursor {
|
||||
return;
|
||||
}
|
||||
let wire = CollabLocalPresenceWire {
|
||||
cursor: cursor.map(|(x, y)| CollabPointWire { x, y }),
|
||||
client_id: None,
|
||||
};
|
||||
let Ok(body) = serde_json::to_string(&wire) else {
|
||||
return;
|
||||
};
|
||||
LAST_PRESENCE_MS.set(now);
|
||||
LAST_PRESENCE_POINT.set(cursor);
|
||||
let _ = live_sync::post_json(&format!("{base}{}", collab_routes::PRESENCE), &body, None);
|
||||
}
|
||||
|
||||
fn now_ms_f64() -> f64 {
|
||||
web_sys::window()
|
||||
.and_then(|window| window.performance())
|
||||
.map(|performance| performance.now())
|
||||
.unwrap_or(0.0)
|
||||
}
|
||||
|
||||
fn now_ms() -> u64 {
|
||||
now_ms_f64().max(0.0) as u64
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn doc_with(ids: &[&str]) -> op_editor_core::EditorState {
|
||||
let children: Vec<serde_json::Value> = ids
|
||||
.iter()
|
||||
.map(|id| {
|
||||
serde_json::json!({
|
||||
"type": "rectangle", "id": id,
|
||||
"x": 0, "y": 0, "width": 4, "height": 4
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
let mut state = op_editor_core::EditorState::starter();
|
||||
state.doc = serde_json::from_value(serde_json::json!({
|
||||
"version": "1.0",
|
||||
"children": children,
|
||||
}))
|
||||
.expect("valid document");
|
||||
state
|
||||
}
|
||||
|
||||
fn reset_latches() {
|
||||
DAEMON_NODE_IDS.with(|ids| *ids.borrow_mut() = None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_idle_shell_never_blocks_a_push() {
|
||||
reset_latches();
|
||||
let state = doc_with(&["n100", "n101"]);
|
||||
assert_eq!(
|
||||
state.editor_ui.collab.phase,
|
||||
CollabConnectionPhase::Idle,
|
||||
"precondition"
|
||||
);
|
||||
assert!(
|
||||
!push_blocked_by_session(&state),
|
||||
"a shell with no session must sync exactly as it did before"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_active_session_blocks_a_push_that_invents_a_node_id() {
|
||||
reset_latches();
|
||||
let mut state = doc_with(&["n100"]);
|
||||
note_daemon_document(&state);
|
||||
state
|
||||
.editor_ui
|
||||
.collab
|
||||
.set_phase(CollabConnectionPhase::Active);
|
||||
|
||||
// Editing what the daemon already knows is fine.
|
||||
assert!(!push_blocked_by_session(&state));
|
||||
|
||||
// Minting a new local id is not: the browser has no owner-assigned
|
||||
// namespace, so this id could collide with a peer's.
|
||||
let mut grown = doc_with(&["n100", "n101"]);
|
||||
grown
|
||||
.editor_ui
|
||||
.collab
|
||||
.set_phase(CollabConnectionPhase::Active);
|
||||
assert!(push_blocked_by_session(&grown));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deleting_a_node_during_a_session_is_not_blocked() {
|
||||
reset_latches();
|
||||
let seed = doc_with(&["n100", "n101"]);
|
||||
note_daemon_document(&seed);
|
||||
|
||||
let mut shrunk = doc_with(&["n100"]);
|
||||
shrunk
|
||||
.editor_ui
|
||||
.collab
|
||||
.set_phase(CollabConnectionPhase::Active);
|
||||
assert!(
|
||||
!push_blocked_by_session(&shrunk),
|
||||
"removals carry no new ids and must keep syncing"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_active_session_blocks_until_the_first_daemon_document_arrives() {
|
||||
reset_latches();
|
||||
let mut state = doc_with(&["n100"]);
|
||||
state
|
||||
.editor_ui
|
||||
.collab
|
||||
.set_phase(CollabConnectionPhase::Active);
|
||||
assert!(
|
||||
push_blocked_by_session(&state),
|
||||
"with no daemon document to compare against, refusing is the safe answer"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn every_ui_action_has_a_wire_form() {
|
||||
let key = op_editor_core::CollabAdmissionRequestKey::new("req-1").expect("valid key");
|
||||
for action in [
|
||||
CollabUiAction::OpenCreate,
|
||||
CollabUiAction::Start,
|
||||
CollabUiAction::StartLan,
|
||||
CollabUiAction::OpenJoin,
|
||||
CollabUiAction::BeginDiscovery,
|
||||
CollabUiAction::Cancel,
|
||||
CollabUiAction::Retry,
|
||||
CollabUiAction::Leave,
|
||||
CollabUiAction::DiscardPending,
|
||||
CollabUiAction::ReapplyDiscarded,
|
||||
CollabUiAction::SaveAsFork,
|
||||
CollabUiAction::JoinAddress {
|
||||
endpoint: "1.2.3.4:5".into(),
|
||||
},
|
||||
CollabUiAction::JoinDiscovered {
|
||||
discovery_id: "d".into(),
|
||||
},
|
||||
CollabUiAction::ApproveAdmissionEditor {
|
||||
request_key: key.clone(),
|
||||
},
|
||||
CollabUiAction::RejectOwnerIdentity { request_key: key },
|
||||
] {
|
||||
let wire = wire_action(&action).expect("every action maps");
|
||||
assert!(serde_json::to_string(&wire).is_ok(), "{action:?}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn admission_keys_survive_the_round_trip_to_the_daemon() {
|
||||
let key = op_editor_core::CollabAdmissionRequestKey::new("req-abc_1").expect("valid");
|
||||
let wire = wire_action(&CollabUiAction::RejectAdmission {
|
||||
request_key: key.clone(),
|
||||
})
|
||||
.expect("maps");
|
||||
assert_eq!(
|
||||
wire.clone().into_ui_action().expect("revalidates"),
|
||||
CollabUiAction::RejectAdmission { request_key: key }
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -50,6 +50,9 @@ mod tooltip_pump;
|
|||
mod live_sync_glue;
|
||||
#[cfg(feature = "canvaskit")]
|
||||
mod web_auth_sync;
|
||||
// Daemon collaboration relay (action drain + projection pull + presence).
|
||||
#[cfg(feature = "canvaskit")]
|
||||
mod collab_sync;
|
||||
// Shared daemon base-URL resolution (page origin when served by the daemon,
|
||||
// localhost fallback for the dev smoke page).
|
||||
#[cfg(feature = "canvaskit")]
|
||||
|
|
|
|||
|
|
@ -307,6 +307,10 @@ fn poll_version<C: RepaintContext + 'static>(
|
|||
let fetch_busy = fetch_busy.clone();
|
||||
let last_selection_key = last_selection_key.clone();
|
||||
let on_version: Rc<dyn Fn(String)> = Rc::new(move |body: String| {
|
||||
// The daemon answers both counters here. Hand the collaboration one to
|
||||
// its own loop rather than opening a second probe; a `collabSeq` bump
|
||||
// must never reach the document fetch below.
|
||||
crate::collab_sync::note_version_probe(&body);
|
||||
let Some(version) = WebSyncClient::parse_version_probe(&body) else {
|
||||
return; // daemon down / non-JSON error body — retry next tick
|
||||
};
|
||||
|
|
@ -405,6 +409,10 @@ fn apply_document_response<C: RepaintContext + 'static>(
|
|||
})
|
||||
.unwrap_or(false);
|
||||
if applied {
|
||||
// Every id in this document came from the daemon, so it is the set an
|
||||
// active session will accept a push against. Recording it here is what
|
||||
// lets the push gate tell an edited node from an invented one.
|
||||
crate::collab_sync::note_daemon_document(inner_ref.host().editor_state());
|
||||
// Baseline = OUR serialization of the just-applied document, so the
|
||||
// push tick compares apples to apples (serde normalization differs
|
||||
// from the daemon's wire bytes).
|
||||
|
|
@ -529,6 +537,12 @@ fn push_document_if_changed<C: RepaintContext + 'static>(
|
|||
// from the scalar baseline, not mistaken for a content change.
|
||||
reasons.editor_meta
|
||||
};
|
||||
// An active collaboration session cannot sequence a node id this browser
|
||||
// minted from its local counter. Hold the push rather than fork the shared
|
||||
// document; the next pull replaces the local node with the daemon's copy.
|
||||
if crate::collab_sync::push_blocked_by_session(inner.borrow().host().editor_state()) {
|
||||
return;
|
||||
}
|
||||
if !should_push {
|
||||
// Bytes already match the daemon baseline (e.g. a generation-only
|
||||
// replace — same content, new generation): no push to send, but the
|
||||
|
|
|
|||
|
|
@ -207,3 +207,31 @@ impl WidgetHost {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl WidgetHost {
|
||||
/// The last known cursor position in document coordinates, for the
|
||||
/// collaboration presence relay.
|
||||
///
|
||||
/// `None` before the pointer has ever entered the canvas or while it sits
|
||||
/// over the chrome — peers should see no cursor rather than one parked at
|
||||
/// the origin.
|
||||
pub(crate) fn last_cursor_doc_point(&self, width: f32, height: f32) -> Option<(f64, f64)> {
|
||||
use op_editor_ui::widgets::host_canvas_geometry as canvas_geometry;
|
||||
|
||||
if !canvas_geometry::over_canvas(
|
||||
&self.editor_state,
|
||||
self.last_cursor_x,
|
||||
self.last_cursor_y,
|
||||
width,
|
||||
height,
|
||||
) {
|
||||
return None;
|
||||
}
|
||||
let point = canvas_geometry::canvas_doc_point_unclamped(
|
||||
&self.editor_state,
|
||||
self.last_cursor_x,
|
||||
self.last_cursor_y,
|
||||
);
|
||||
Some((f64::from(point.x), f64::from(point.y)))
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue