Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 30 additions & 15 deletions crates/connetto-web/src/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2349,7 +2349,8 @@ fn schedule_recovery_event(
/// What the recovery loop serves while the connection sits idle, in the retry sleep and
/// during the connect.
///
/// Everything but a frame, which needs the server.
/// Tab work that the replica or the hub can answer stays live offline, so a tab can boot,
/// subscribe, mutate and ask content questions before the worker reaches the server.
fn recovery_serves_idle(event: &HubEvent) -> bool {
matches!(
event,
Expand All @@ -2362,6 +2363,7 @@ fn recovery_serves_idle(event: &HubEvent) -> bool {
| HubEvent::ForgetRetired(_, _)
| HubEvent::RefusedContent(_)
| HubEvent::RetryRefused(_, _)
| HubEvent::Frame(_, _)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Fan out offline mutations to subscribed sibling tabs

When one tab mutates a synced table during recovery, this newly served Frame is applied to the worker replica by handle_synced_mutation, but that path does not enqueue a LivePatch for the other tabs and ordinarily relies on the server echo to do so. Consequently, while the upstream remains offline, an already-subscribed sibling stays stale, whereas a tab subscribing afterward receives the new row from serve_snapshot; the tabs only converge after reconnect. The offline frame path should keep existing subscribers consistent with snapshots of the same replica.

Useful? React with 👍 / 👎.

| HubEvent::Internal(_, _, _)
)
}
Expand Down Expand Up @@ -4354,10 +4356,14 @@ mod tests {
assert!(deferred.is_empty());
}

/// Chapter 18's idle column: everything but a frame is served where it arrives while
/// the connection sits idle, because a frame is the only request that needs the server.
/// Chapter 18's idle column (amended R69-F): every event that can be answered from
/// local state is served where it arrives during idle recovery, including Frames. A
/// tab mutation commits to durable pending and replays exactly once on attach under
/// the R67 watermark, mirroring the native client's offline behaviour. The only events
/// that wait for the upstream are those in `recovery_interrupts_attach` not in this
/// predicate, because the attach phase owns the connection for replay.
#[wasm_bindgen_test]
fn only_a_frame_waits_for_the_upstream_while_the_connection_is_idle() {
fn a_frame_is_served_into_pending_during_idle_recovery() {
use connetto_core::messages::{ControlMessage, Ping};
use connetto_core::traits::IncomingFrame;
use tokio::sync::mpsc::unbounded_channel;
Expand All @@ -4373,7 +4379,7 @@ mod tests {
),
Some(HubEvent::Attached(_, _))
),
"a tab attaching only writes hub state, and its announce has its own deadline"
"a tab attach writes hub state and is served"
);
assert!(matches!(
schedule_recovery_event(&mut deferred, HubEvent::Gone(1), recovery_serves_idle),
Expand All @@ -4383,19 +4389,28 @@ mod tests {
schedule_recovery_event(&mut deferred, HubEvent::Kill(1), recovery_serves_idle),
Some(HubEvent::Kill(1))
));
// Frames are now served during idle recovery: a tab mutation that arrives while
// the upstream is down commits to the replica and to _connetto_pending, then
// replays exactly once when the upstream attaches, the same contract as native.
assert!(
schedule_recovery_event(
&mut deferred,
HubEvent::Frame(
1,
IncomingFrame::Control(ControlMessage::Ping(Ping { nonce: 1 })),
matches!(
schedule_recovery_event(
&mut deferred,
HubEvent::Frame(
1,
IncomingFrame::Control(ControlMessage::Ping(Ping { nonce: 1 })),
),
recovery_serves_idle,
),
recovery_serves_idle,
)
.is_none(),
"a frame needs the server, so it waits"
Some(HubEvent::Frame(_, _))
),
"a frame is served from local state during idle recovery"
);
assert_eq!(
deferred.len(),
0,
"nothing is queued while the deque was empty"
);
assert_eq!(deferred.len(), 1, "and only the frame is queued");
}

/// Chapter 18's attach column: a departure or a kill unsubscribes through the
Expand Down
15 changes: 14 additions & 1 deletion crates/connetto-web/src/workers/boot/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,8 @@ pub struct DbWorkerConfig {
pub(crate) content_namespace: Option<&'static str>,
/// The transport every content transfer runs under, carrying the idle bound.
pub(crate) content_http: connetto_file_client::BrowserHttp,
/// Broadcast channel opening the first upstream connect attempt.
pub(crate) connect_gate: Option<&'static str>,
/// Schema version presented to the server at handshake.
pub(crate) schema_version: connetto_core::SchemaVersion,
/// Custom SQLite functions registered on every connection before any DDL.
Expand Down Expand Up @@ -200,6 +202,7 @@ impl DbWorkerConfig {
hub_meta_name: "",
content_namespace: None,
content_http: connetto_file_client::BrowserHttp::new(),
connect_gate: None,
schema_version,
sql_functions: connetto_client::SqlFunctions::default(),
policy_tables: connetto_client::PolicyTables::new(),
Expand Down Expand Up @@ -273,6 +276,13 @@ impl DbWorkerConfig {
self
}

/// Hold the worker offline until this channel receives any message.
#[must_use]
pub fn with_connect_gate(mut self, connect_gate: &'static str) -> Self {
self.connect_gate = Some(connect_gate);
self
}

/// How long a content transfer may stay silent before it is aborted.
///
/// The bound covers every phase of a request on this device, and a
Expand Down Expand Up @@ -659,7 +669,10 @@ where
let replica_key = replica::resolve_replica_key(&key_store, &spec).await?;
let login = spec.login.take();
let client_config = replica::build_boot_client_config(config, login, &spec);
let transport = replica::try_connect_upstream(config.ws_url).await;
let transport = match config.connect_gate {
Some(_) => None,
None => replica::try_connect_upstream(config.ws_url).await,
};
let (mut worker, content_root_key) =
replica::open_boot_replica(transport, &spec, config, &client_config, replica_key).await?;
replica::subscribe_and_boot(&mut worker, config).await?;
Expand Down
49 changes: 44 additions & 5 deletions crates/connetto-web/src/workers/boot/services.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,15 @@
use std::cell::RefCell;
use std::cell::{Cell, RefCell};
use std::rc::Rc;

use connetto_client::ConnettoConnection;
use connetto_client::reconnect::{ReconnectPolicy, Sleeper, TransportFactory};
use connetto_core::messages::SubscriptionSpec;
use connetto_file_client::{BrowserStore, ContentArchive};
use tokio::sync::mpsc::UnboundedReceiver;
use wasm_bindgen::JsCast;
use wasm_bindgen::closure::Closure;
use wasm_bindgen_futures::spawn_local;
use web_sys::{BroadcastChannel, MessageEvent};

use super::super::helpers::content_store_namespace;
use super::BootError;
Expand All @@ -19,6 +22,32 @@ thread_local! {
static DB_ALIVE: RefCell<Option<locks::HeldLock>> = const { RefCell::new(None) };
}

struct ConnectGate {
open: Rc<Cell<bool>>,
_channel: BroadcastChannel,
_listener: Closure<dyn FnMut(MessageEvent)>,
}

fn install_connect_gate(name: &'static str) -> Result<Rc<ConnectGate>, BootError> {
let channel =
BroadcastChannel::new(name).map_err(|err| super::super::IntakeError::ChannelOpen {
operation: "connect gate",
detail: format!("{err:?}"),
})?;
let open = Rc::new(Cell::new(false));
let listener = {
let open = Rc::clone(&open);
Closure::<dyn FnMut(MessageEvent)>::new(move |_event: MessageEvent| {
open.set(true);
})
};
channel.set_onmessage(Some(listener.as_ref().unchecked_ref()));
Ok(Rc::new(ConnectGate {
open,
_channel: channel,
_listener: listener,
}))
}
/// Install storage and custody, carry out any outstanding data wipe, and
/// reserve this boot's database slots.
///
Expand Down Expand Up @@ -117,11 +146,21 @@ pub(super) async fn start_boot_services<Id>(
content_root_key: Option<[u8; 32]>,
) -> Result<Option<bool>, BootError> {
let ws_url = config.ws_url;
let connect_gate = match config.connect_gate {
Some(name) => Some(install_connect_gate(name)?),
None => None,
};
let reconnect = HubReconnect {
factory: move || async move {
BrowserSocket::connect(ws_url)
.await
.map_err(|err| err.to_string())
factory: move || {
let connect_gate = connect_gate.clone();
async move {
if connect_gate.as_ref().is_some_and(|gate| !gate.open.get()) {
return Err("connect gate closed".to_owned());
}
BrowserSocket::connect(ws_url)
.await
.map_err(|err| err.to_string())
}
},
sleeper: super::super::intake::sleep,
policy: ReconnectPolicy::default(),
Expand Down
8 changes: 5 additions & 3 deletions docs/architecture/18-file-handling.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# 18: File handling

**Status**: normative for the decisions it records. R64 the file core, R65 the file server, R66 the connetto seam, R67 the native client and R68 the browser client are built. R69 the demos is in progress, its executable half, demo schemas, tab and worker content protocol and browser-stack wiring built 2026-09-17 as pull requests #28 to #31, with the demo surfaces and the offline and two-viewer proofs open. R79 the peer link and R87 the quotas are not built. Every statement carries **Decided (RN)** or an **Amended (RN)** beside it, where `RN` is the phase in `plans/master-implementation-plan.md` that owns it, and that phase's section records each decision with its rejected alternatives. Chapter 07 is the historical record of the thinking that preceded these decisions and defers to this chapter wherever the two disagree.
**Status**: normative for the decisions it records. R64 the file core, R65 the file server, R66 the connetto seam, R67 the native client and R68 the browser client are built. R69 the demos is in progress, its executable half, demo schemas, tab and worker content protocol and browser-stack wiring built 2026-09-17 as pull requests #28 to #31, and its offline stage, reconnect upload and two-viewer refusal proofs built 2026-09-18 with the recovery table's Frame row amended to match, leaving the demo surfaces open. R79 the peer link and R87 the quotas are not built. Every statement carries **Decided (RN)** or an **Amended (RN)** beside it, where `RN` is the phase in `plans/master-implementation-plan.md` that owns it, and that phase's section records each decision with its rejected alternatives. Chapter 07 is the historical record of the thinking that preceded these decisions and defers to this chapter wherever the two disagree.

---

Expand Down Expand Up @@ -56,10 +56,12 @@ Client content divides into two classes. Unsent content, authored here and not y

**Amended (R68, 2026-09-12): what the worker relay serves in every situation it can be in.** R68 made the recovery loop serve requests it had refused on `main`, and the list of situations it handles was never written down, so sixteen review findings were each one line of this table. The table is the list. A cell says whether the request is served where it arrives, held until the upstream is back, or refused because the hub has ended.

**Amended (R69-F, 2026-09-18): the Frame row for the sleeping and connecting situations.** R69's offline-stage feature requires a tab to complete its full protocol (handshake, subscribe, mutation) against the relay while the upstream is absent. Every Frame the relay receives during idle recovery can be answered from local state: a handshake reads the tab's watermark from the replica, a subscribe reads the snapshot from the replica, and a mutation commits to the replica and to `_connetto_pending`, then replays exactly once when the upstream first attaches, under the R67 exactly-once watermark. This is the same contract as the native client's offline write path, and serving Frames during idle recovery makes the browser relay a true local peer rather than a proxy that requires the server to be reachable. The `connect_gate` default of `None` does not change any other row of the table or any other column of the Frame row.

| Request | Connected | Sleeping before retry | Connecting | Attaching | Replaying subscriptions | Closed |
|---|---|---|---|---|---|---|
| `Attached` | Served | Served | Served | Served, no interruption | Served, no interruption | Refused |
| `Frame` | Served | Held | Held | Held | Held | Refused |
| `Frame` | Served | Served | Served | Held | Held | Refused |
| `Gone` | Served | Served | Served | Held | Held | Refused |
| `Kill` | Served | Served | Served | Held | Held | Refused |
| `Unsynced` | Served | Served | Served | Served, interrupts | Served, interrupts | Refused |
Expand All @@ -69,7 +71,7 @@ Client content divides into two classes. Unsent content, authored here and not y
| `RefusedContent` | Served | Served | Served | Served, interrupts | Served, interrupts | Refused |
| `RetryRefused` | Served | Served | Served | Served, interrupts | Served, interrupts | Refused |

A frame is held in every recovering situation because a write or a subscribe needs the server, and it is the only request that does. An attachment is served everywhere because registering a tab writes hub state and nothing else, and the tab's own announce gives up after fifteen seconds, so holding it fails a healthy tab for the length of an outage. A departure and a kill unsubscribe upstream, so they are served while the connection is idle and held while the attach owns it, where nothing waits on them. The six requests that answer from the replica interrupt an attach rather than wait for it, and the interrupted attach resumes on the installed transport afterwards.
A frame is held while the attach or the subscription replay owns the connection, because those phases replay the full upstream state and a concurrent frame would race their ordering. During idle recovery (sleeping before retry, connecting) every Frame is served from local state into durable pending: the relay answers it from the replica or from hub state alone, commits mutations to `_connetto_pending`, and they replay exactly once on attach under the R67 watermark. This is native parity: the original rationale that a write or a subscribe needs the server is corrected by this amendment. The relay has everything it needs to serve any Frame offline, and holding them during a transient outage makes a healthy tab wait out the reconnect for no reason. An attachment is served everywhere because registering a tab writes hub state and nothing else, and the tab's own announce gives up after fifteen seconds, so holding it fails a healthy tab for the length of an outage. A departure and a kill unsubscribe upstream, so they are served while the connection is idle and held while the attach owns it, where nothing waits on them. The six requests that answer from the replica interrupt an attach rather than wait for it, and the interrupted attach resumes on the installed transport afterwards.

One order rule covers every served cell. A request that would be served where it arrives is queued instead while any older request is already queued, and the queue drains in arrival order once the upstream is back, so a kill can never overtake a frame of the tab it kills. Refused means the intake has closed and the hub has ended, so every waiting reply sender drops and its caller reads the hub as gone.

Expand Down
1 change: 1 addition & 0 deletions examples/wasm-smoke/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions examples/wasm-smoke/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ serde_json = "1"
# The full-resync parity test hand-builds snapshot patchsets to feed a fake
# upstream through the relay hub, the same native encoding the wire carries.
sqlite-diff-rs = { version = "0.11", default-features = false }
connetto-file-client = { path = "../../crates/connetto-file-client" }
zstd = "0.13"

# Standalone workspace: wasm-only spike, never enters the root gate.
Expand Down
37 changes: 37 additions & 0 deletions examples/wasm-smoke/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,8 @@ pub mod workers {
pub const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0";
/// The extra upstream subscription the photo flow needs.
pub const PHOTO_QUERY: &str = "SELECT * FROM photos";
/// A test-only gate that keeps the photo worker offline until opened.
pub const PHOTO_CONNECT_CHANNEL: &str = "connetto-photo-connect";
/// The OPFS file holding the DB worker's durable replica.
pub const DB_NAME: &str = "connetto-relay.sqlite";
/// OPFS file for unlock-protocol tests, separate from DB_NAME so the two
Expand Down Expand Up @@ -238,6 +240,41 @@ pub mod workers {
.map_err(JsValue::from)
}

/// DB worker entry point for the offline photo flow test binary.
///
/// # Errors
///
/// A string describing the VFS, upstream connect, or subscribe failure.
#[wasm_bindgen]
pub async fn db_worker_photo_offline_boot() -> Result<(), JsValue> {
connetto_web::logging::init_console();
connetto_web::workers::boot_db_worker::<String>(
&connetto_web::workers::DbWorkerConfig::new(crate::demo_schema_version())
.with_ws_url(DEMO_WS_URL)
.with_replica_db_prefix(DB_NAME)
.with_replica_ddl(DEMO_SQLITE_DDL)
.with_frontend_ddl(DEMO_FRONTEND_DDL)
.with_upstream_sub_id("db-upstream")
.with_upstream_query(DEMO_QUERY)
.with_extra_upstream("db-photos-upstream", PHOTO_QUERY)
.with_hub_meta_name("connetto-hub-meta.sqlite")
.with_content_namespace("connetto-photo-content")
.with_sql_functions(crate::uuidv4_functions())
.with_policy_tables(crate::demo_policy_tables())
.with_caller_function(crate::CALLER_FUNCTION)
.with_auth(Some(connetto_web::auth::WorkerAuthConfig::new(
"http://127.0.0.1:18099",
"dev-idp",
"http://127.0.0.1:18099/dev/landing",
)))
.with_auth_db_name("connetto-auth.sqlite")
.with_connect_gate(PHOTO_CONNECT_CHANNEL),
)
.await
.map(drop)
.map_err(JsValue::from)
}

/// DB worker entry point for the unlock-protocol test binary. Same as
/// `db_worker_boot` except the passkey unlock protocol is enabled.
///
Expand Down
Loading
Loading