From 2d6e59571cfba3f8ea3182d8731a74fb9f5560a9 Mon Sep 17 00:00:00 2001 From: LucaCappelletti94 Date: Fri, 18 Sep 2026 11:42:00 +0200 Subject: [PATCH 1/3] Record R69 E as built in the plan row, section prose and chapter status --- docs/architecture/18-file-handling.md | 2 +- plans/master-implementation-plan.md | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/architecture/18-file-handling.md b/docs/architecture/18-file-handling.md index 509590c7..e6cdbd57 100644 --- a/docs/architecture/18-file-handling.md +++ b/docs/architecture/18-file-handling.md @@ -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, 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. +**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, and the web demos' photos surfaces built 2026-09-18, leaving the desktop demo surface 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. --- diff --git a/plans/master-implementation-plan.md b/plans/master-implementation-plan.md index a9610064..8f6da127 100644 --- a/plans/master-implementation-plan.md +++ b/plans/master-implementation-plan.md @@ -249,7 +249,7 @@ Execution order and nothing else. Status, blockers, landing dates and what each | R87 storage quotas and deployment ceilings | NOT STARTED, raised 2026-09-08 | R65, which is done | no | | R67 native file client | **DONE** (2026-09-08) | nothing. `connetto-file-client`: the `std::fs` encrypted chunk store, the manifests and outbox in the replica committed with the entry row, the outbox walk with a boot integrity pass, the resolver over a `LocalContentSource` list, the query-shaped pin surface with a whole-file fetch, and `tidy_content`. Thirteen decisions recorded above, three of them defects found by grounding: the tier cannot be atomic with the replica, `MemStore` answered empty bytes for an absent chunk, and a double-quoted pin column silently became a string literal. 28 tests, the offline photo case among them, end to end against a real Postgres, a real file server on a socket and two real devices | no | | R68 browser file client | **DONE** (2026-09-10) | nothing. Worker-owned encrypted OPFS with memory fallback, browser fetch, reference-counted object URLs, and version 3 archives that restore unsent content through the production worker relay. The offline photo survives export, import under another key, local display and later upload. The browser stack passed and the full release suite passed 738 tests with 3 skipped | no | -| R69 files in every demo | IN PROGRESS (A, B, C, browser-stack wiring and F done by 2026-09-18), designed (2026-09-12) | nothing. A is #28, B is #29, C is #30, and the browser stack boots the executable's file half with the online photo flow proven by `photo_flow.rs` and F's offline and two-viewer proofs by `photo_offline.rs` and `photo_visibility.rs` in this pull request. D and E remain | no | +| R69 files in every demo | IN PROGRESS (A, B, C, browser-stack wiring, F and E done by 2026-09-18), designed (2026-09-12) | nothing. A is #28, B is #29, C is #30, and the browser stack boots the executable's file half with the online photo flow proven by `photo_flow.rs` and F's offline and two-viewer proofs by `photo_offline.rs` and `photo_visibility.rs`. E is built in the pull request adding `photo_flow.rs` to both web demos, and D remains | no | | R70 backup and restore story | NOT STARTED | nothing | no | | R71 Linux key custody survives reboot | NOT STARTED | nothing for grounding. One custody decision to take with the maintainer at execution | no | | R72 clock discipline (X6) | NOT STARTED | nothing | no | @@ -4733,7 +4733,7 @@ The offline photo case runs in headless Chrome end to end, and an export taken o ## R69: files in every demo -**Status.** IN PROGRESS (2026-09-17), **designed 2026-09-12** with the maintainer. A is #28, B is #29 and C is #30, and #31 landed the browser-stack half of A with `examples/wasm-smoke/tests/photo_flow.rs` driving a photo through stage, commit, resolve and fetch in headless Chrome, so the executable's file half now boots for real and the online round trip is proven in CI. F is built in the pull request that adds `examples/wasm-smoke/tests/photo_offline.rs` and `photo_visibility.rs`, the stage that completes behind a connect gate and uploads on attach, the second viewer refused by the server's mint under the owner-only policy, and the recovery table's Frame row amended to match. D and E remain. The phase as first written assumed three things that did not exist, found by reading every module of the file stack against the two steps: no `main` ran the file server's router (only `crates/connetto-file-client/tests/it/offline_photo.rs` binds one), the shipped `connetto-server` executable passed `NoSigner` (`bin/connetto-server.rs`) so no deployment could mint a ticket, and no demo schema declared the metadata table or the two SQL functions the file server's preflight requires. Six decisions below close them, each with what was rejected, and the steps are seven pull requests with their dependencies stated so three of them run at once. +**Status.** IN PROGRESS (2026-09-17), **designed 2026-09-12** with the maintainer. A is #28, B is #29 and C is #30, and #31 landed the browser-stack half of A with `examples/wasm-smoke/tests/photo_flow.rs` driving a photo through stage, commit, resolve and fetch in headless Chrome, so the executable's file half now boots for real and the online round trip is proven in CI. F is built in the pull request that adds `examples/wasm-smoke/tests/photo_offline.rs` and `photo_visibility.rs`, the stage that completes behind a connect gate and uploads on attach, the second viewer refused by the server's mint under the owner-only policy, and the recovery table's Frame row amended to match. E is built in the pull request that adds a photos surface and a `photo_flow.rs` suite to both web demos, their schemas byte-identical to wasm-smoke's with the handshake guard as the suite's first test. D remains. The phase as first written assumed three things that did not exist, found by reading every module of the file stack against the two steps: no `main` ran the file server's router (only `crates/connetto-file-client/tests/it/offline_photo.rs` binds one), the shipped `connetto-server` executable passed `NoSigner` (`bin/connetto-server.rs`) so no deployment could mint a ticket, and no demo schema declared the metadata table or the two SQL functions the file server's preflight requires. Six decisions below close them, each with what was rejected, and the steps are seven pull requests with their dependencies stated so three of them run at once. **Blocked on nothing.** R64 to R68 are done, and the twelve questions raised before this phase (`plans/open-questions-before-r69.md`) were settled, and the nine needing code were built 2026-09-14 to 2026-09-16 as pull requests #18 to #27: the teardown list, the anonymous boot, the two permanent upload outcomes, the silence-bounded transfer, the streaming archive, the memory fallback, the keyring index, the span-chained log capture and the cursor-bounded silence assertion, each recorded in chapters 13, 14, 18 or `open-questions.md`. From 156c6c8c0b45ee85fe80122111b61ff34acbd42f Mon Sep 17 00:00:00 2001 From: LucaCappelletti94 Date: Fri, 18 Sep 2026 11:42:09 +0200 Subject: [PATCH 2/3] Give both web demos a photos surface and register their browser suites --- .../src/bin/connetto-browser-stack.rs | 6 + examples/dioxus-web-demo/Cargo.lock | 88 ++- examples/dioxus-web-demo/Cargo.toml | 23 +- examples/dioxus-web-demo/build.rs | 8 +- examples/dioxus-web-demo/policies.sql | 26 + examples/dioxus-web-demo/schema.sql | 38 +- examples/dioxus-web-demo/src/lib.rs | 114 ++++ examples/dioxus-web-demo/src/main.rs | 244 +++++++- examples/dioxus-web-demo/tests/photo_flow.rs | 524 ++++++++++++++++++ examples/wasm-smoke/Cargo.lock | 12 +- examples/yew-web-demo/Cargo.lock | 125 ++++- examples/yew-web-demo/Cargo.toml | 21 + examples/yew-web-demo/build.rs | 8 +- examples/yew-web-demo/policies.sql | 26 + examples/yew-web-demo/schema.sql | 38 +- examples/yew-web-demo/src/lib.rs | 114 ++++ examples/yew-web-demo/src/main.rs | 240 +++++++- examples/yew-web-demo/tests/photo_flow.rs | 524 ++++++++++++++++++ 18 files changed, 2065 insertions(+), 114 deletions(-) create mode 100644 examples/dioxus-web-demo/policies.sql create mode 100644 examples/dioxus-web-demo/src/lib.rs create mode 100644 examples/dioxus-web-demo/tests/photo_flow.rs create mode 100644 examples/yew-web-demo/policies.sql create mode 100644 examples/yew-web-demo/src/lib.rs create mode 100644 examples/yew-web-demo/tests/photo_flow.rs diff --git a/crates/connetto-test-harness/src/bin/connetto-browser-stack.rs b/crates/connetto-test-harness/src/bin/connetto-browser-stack.rs index a96075be..1365e8cb 100644 --- a/crates/connetto-test-harness/src/bin/connetto-browser-stack.rs +++ b/crates/connetto-test-harness/src/bin/connetto-browser-stack.rs @@ -465,6 +465,12 @@ async fn run_default_browser_suites(services: &Services, shard: Option) - for test in test_files(&["examples", "wasm-smoke", "tests"])? { suite_args.push(per_test_args("examples/wasm-smoke", test)); } + for test in test_files(&["examples", "yew-web-demo", "tests"])? { + suite_args.push(per_test_args("examples/yew-web-demo", test)); + } + for test in test_files(&["examples", "dioxus-web-demo", "tests"])? { + suite_args.push(per_test_args("examples/dioxus-web-demo", test)); + } let total = suite_args.len(); if let Some(shard) = shard { diff --git a/examples/dioxus-web-demo/Cargo.lock b/examples/dioxus-web-demo/Cargo.lock index 6d85dd49..5d26a793 100644 --- a/examples/dioxus-web-demo/Cargo.lock +++ b/examples/dioxus-web-demo/Cargo.lock @@ -328,6 +328,12 @@ dependencies = [ "serde", ] +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cc" version = "1.4.5" @@ -566,16 +572,21 @@ dependencies = [ "connetto-client", "connetto-core", "connetto-dioxus", + "connetto-file-client", + "connetto-file-core", "connetto-web", "diesel", "dioxus", + "futures-channel", "js-sys", "libsqlite3-sys", "pg2sqlite", "rosetta-uuid", + "serde_json", "tracing", "wasm-bindgen", "wasm-bindgen-futures", + "wasm-bindgen-test", "web-sys", ] @@ -621,6 +632,7 @@ dependencies = [ "connetto-client", "connetto-core", "connetto-file-client", + "connetto-file-core", "diesel", "diesel-sqlite-session", "futures-channel", @@ -1608,7 +1620,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2795,6 +2807,16 @@ dependencies = [ "unicase", ] +[[package]] +name = "minicov" +version = "0.3.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4869b6a491569605d66d3952bcdf03df789e5b536e5f0cf7758a7f08a55ae24d" +dependencies = [ + "cc", + "walkdir", +] + [[package]] name = "minimal-lexical" version = "0.2.1" @@ -2906,6 +2928,15 @@ dependencies = [ "minimal-lexical", ] +[[package]] +name = "nu-ansi-term" +version = "0.50.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "num-bigint" version = "0.4.8" @@ -2993,6 +3024,12 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "oorandom" +version = "11.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" + [[package]] name = "opaque-debug" version = "0.3.1" @@ -3254,7 +3291,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3622,7 +3659,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -4225,10 +4262,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -4865,6 +4902,45 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "wasm-bindgen-test" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45863ef0bef521c12124eb39d9a38513c47db75f22e503c06beacd40afeb35db" +dependencies = [ + "async-trait", + "cast", + "js-sys", + "libm", + "minicov", + "nu-ansi-term", + "num-traits", + "oorandom", + "serde", + "serde_json", + "wasm-bindgen", + "wasm-bindgen-futures", + "wasm-bindgen-test-macro", + "wasm-bindgen-test-shared", +] + +[[package]] +name = "wasm-bindgen-test-macro" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8c89dcab8b516b6b603baca9d550b7282d68fcc7f367e3956cff7ebf406a3f12" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.5", +] + +[[package]] +name = "wasm-bindgen-test-shared" +version = "0.2.128" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f37b4f992cebe528ef34964ae69681ac0fe7080071e7298e46008f9d380302af" + [[package]] name = "wasm-streams" version = "0.4.2" @@ -4913,7 +4989,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/examples/dioxus-web-demo/Cargo.toml b/examples/dioxus-web-demo/Cargo.toml index 9f276ce3..decd1128 100644 --- a/examples/dioxus-web-demo/Cargo.toml +++ b/examples/dioxus-web-demo/Cargo.toml @@ -7,6 +7,9 @@ rust-version = "1.88" license = "MIT" publish = false +[lib] +crate-type = ["cdylib", "rlib"] + [dependencies] # The browser platform: transports, leader election, locks, the relay hub, # and the DB worker orchestration. The demo supplies its schema and baked @@ -16,6 +19,8 @@ connetto-web = { path = "../../crates/connetto-web" } # Default features off: this is a wasm build with a browser transport. connetto-client = { path = "../../crates/connetto-client", default-features = false } connetto-core = { path = "../../crates/connetto-core" } +# FileId and MimeClass for the photo staging closure. +connetto-file-core = { path = "../../crates/connetto-file-core" } # The one live-query hook, bound to the component lifecycle. connetto-dioxus = { path = "../../crates/connetto-dioxus" } diesel = { version = "2", features = ["sqlite"] } @@ -32,7 +37,9 @@ rosetta-uuid = { version = "0.1.3", features = ["diesel", "sqlite"] } # `File`, `FileList` and `HtmlInputElement` back the import file picker. web-sys = { version = "0.3", features = [ "Blob", + "BlobPropertyBag", "BroadcastChannel", + "DedicatedWorkerGlobalScope", "Document", "Element", "File", @@ -40,15 +47,29 @@ web-sys = { version = "0.3", features = [ "HtmlAnchorElement", "HtmlElement", "HtmlInputElement", - "MessageEvent", + "Request", + "RequestInit", + "Response", + "UrlSearchParams", "Url", "Window", + "Worker", + "WorkerGlobalScope", + "WorkerOptions", + "WorkerType", "console", ] } # Structured logging (R12): this demo emits through `tracing` and installs the # developer-console destination in `main`. tracing = { version = "0.1", default-features = false, features = ["std"] } +[dev-dependencies] +wasm-bindgen-test = "0.3" +futures-channel = "0.3" +serde_json = "1" +js-sys = "0.3" +connetto-file-client = { path = "../../crates/connetto-file-client" } + # The build script translates schema.sql and frontend.sql through pg2sqlite and # bakes a template per tier, so the app ships every tier's schema pre-applied # and never executes DDL at startup. Same pipeline as the smoke crate. diff --git a/examples/dioxus-web-demo/build.rs b/examples/dioxus-web-demo/build.rs index 666c0923..35aab877 100644 --- a/examples/dioxus-web-demo/build.rs +++ b/examples/dioxus-web-demo/build.rs @@ -110,13 +110,17 @@ fn write_policy_tables(documents: &[&str], views: &[String], out: &std::path::Pa fn main() { println!("cargo::rerun-if-changed=schema.sql"); println!("cargo::rerun-if-changed=frontend.sql"); + println!("cargo::rerun-if-changed=policies.sql"); let out_dir = std::path::PathBuf::from(std::env::var("OUT_DIR").expect("cargo sets OUT_DIR")); - let synced_views = translate(&["schema.sql"], &out_dir.join("replica-ddl.sql")); + let synced_views = translate( + &["schema.sql", "policies.sql"], + &out_dir.join("replica-ddl.sql"), + ); translate(&["frontend.sql"], &out_dir.join("frontend-ddl.sql")); // The synced tier only: the local tier is a separate database, attached // under its own schema, and the check the map feeds reads `main`. write_policy_tables( - &["schema.sql"], + &["schema.sql", "policies.sql"], &synced_views, &out_dir.join("replica-tables.rs"), ); diff --git a/examples/dioxus-web-demo/policies.sql b/examples/dioxus-web-demo/policies.sql new file mode 100644 index 00000000..f170e654 --- /dev/null +++ b/examples/dioxus-web-demo/policies.sql @@ -0,0 +1,26 @@ +-- The row-level security the backend enforces on the synced tables, kept apart +-- from schema.sql because the two reach the server as separate documents: +-- schema.sql feeds CONNETTO_PG_DDL and is what clients sync, this file feeds +-- CONNETTO_PG_POLICIES and is what the authorization model is derived from. +-- Apply both to Postgres, this one last, after schema.sql, the file server +-- DDL, roles.sql and content.sql. +-- build.rs translates the pair together, which is what splits the replica's +-- orders into a backing table, a view of the logical name, and INSTEAD OF +-- triggers. The caller is read from app.user_id, which the server binds per +-- transaction and the replica answers with the registered current_app_user() +-- function, so both ends compare against the same identity. +ALTER TABLE orders ENABLE ROW LEVEL SECURITY; +CREATE POLICY orders_p ON orders USING (owner_id = current_setting('app.user_id', true)); -- NOSONAR S1192, SQL DDL has no constants for the caller setting the three policies share + +-- The same shape on the composite-key table, so its replica half is split the +-- same way and its INSTEAD OF triggers have to match a row on two key columns +-- rather than one. +ALTER TABLE order_lines ENABLE ROW LEVEL SECURITY; +CREATE POLICY order_lines_p ON order_lines USING (owner_id = current_setting('app.user_id', true)); + +-- The photo rows are visible to the owner of the order they hang off, and the +-- owner is repeated on the row as it is on order_lines, so the comparison +-- settles from the row itself. The file server's visibility function consults +-- exactly this table under the caller's identity. +ALTER TABLE photos ENABLE ROW LEVEL SECURITY; +CREATE POLICY photos_p ON photos USING (owner_id = current_setting('app.user_id', true)); diff --git a/examples/dioxus-web-demo/schema.sql b/examples/dioxus-web-demo/schema.sql index 689b4e2c..143b0874 100644 --- a/examples/dioxus-web-demo/schema.sql +++ b/examples/dioxus-web-demo/schema.sql @@ -1,31 +1,33 @@ -- The one source of truth for the demo: the Postgres dialect schema the --- backend owns. build.rs translates this through pg2sqlite and bakes the --- replica template database the app ships. The connetto-server for this demo --- must be started with this same schema in CONNETTO_PG_DDL and must list --- every synced table in CONNETTO_WRITABLE. Apply in this order: this file, +-- backend owns. build.rs translates this through pg2sqlite into the SQLite DDL +-- a first boot applies. The connetto-server for this demo must be started with +-- this same schema in CONNETTO_PG_DDL and must list every synced table in +-- CONNETTO_WRITABLE. Apply in this order: this file, -- connetto_file_server::DEPLOYMENT_DDL, roles.sql (the non-owner role --- required by CONNETTO_READER_URL), then content.sql. --- The server also requires CONNETTO_AUTH, CONNETTO_AUTH_BIND, and the --- CONNETTO_OIDC_* variables written by the dev IdP (see dev_idp.rs). +-- required by CONNETTO_READER_URL), content.sql, then policies.sql. -- The key default is load-bearing on the client rather than here: build.rs -- translates it through pg2sqlite into the replica's own DEFAULT (uuidv4()), -- which mints the key when a local write omits it. Both ends mint version 4. -- The quantity is non-null because every client schema already declares it so. -CREATE TABLE orders (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), quantity BIGINT NOT NULL CHECK (quantity >= 0)); +-- owner_id carries who a row belongs to, which policies.sql compares against +-- the caller. It has no default: pg2sqlite maps current_setting only inside a +-- policy expression, so a default naming the caller would translate into a +-- call the replica cannot resolve, and every write names the owner instead. +CREATE TABLE orders (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), owner_id TEXT NOT NULL, quantity BIGINT NOT NULL CHECK (quantity >= 0)); -- The lines of an order, keyed by the order and the line number together. It is --- the one table here whose key spans two columns, which the replica's own schema --- and every key connetto encodes on the wire have to carry as a pair. -CREATE TABLE order_lines ( - order_id UUID NOT NULL REFERENCES orders(id), - line_no INTEGER NOT NULL, - quantity BIGINT NOT NULL CHECK (quantity >= 0), - PRIMARY KEY (order_id, line_no) -); +-- the one table here whose key spans two columns, and the translation below +-- splits it like any policy-bearing table, so its INSTEAD OF triggers match a +-- row on both key columns rather than one. owner_id repeats rather than being +-- read through the parent order, so the policy settles from the row itself, +-- which is what keeps the change path free of a round trip. +CREATE TABLE order_lines (order_id UUID NOT NULL REFERENCES orders(id), line_no INTEGER NOT NULL, owner_id TEXT NOT NULL, quantity BIGINT NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)); -- The photo entry: metadata for one file's bytes, attached to an order. -- content_id is the BLAKE3 identity the file server stores and serves under, -- and content_state stays null until the file server's commit writes -- `available`, so the placeholder condition is "not available" and the --- availability flip arrives as an ordinary synced column change. -CREATE TABLE photos (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), order_id UUID NOT NULL REFERENCES orders(id), content_id BYTEA NOT NULL, content_state TEXT); +-- availability flip arrives as an ordinary synced column change. owner_id +-- repeats from the parent order as it does on order_lines, so the visibility +-- policy settles from the row itself. +CREATE TABLE photos (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), order_id UUID NOT NULL REFERENCES orders(id), owner_id TEXT NOT NULL, content_id BYTEA NOT NULL, content_state TEXT); diff --git a/examples/dioxus-web-demo/src/lib.rs b/examples/dioxus-web-demo/src/lib.rs new file mode 100644 index 00000000..4dd2507e --- /dev/null +++ b/examples/dioxus-web-demo/src/lib.rs @@ -0,0 +1,114 @@ +//! Browser photo worker infrastructure for the dioxus web demo. +//! +//! This lib target exposes the `db_worker_photo_boot` wasm-bindgen entry point +//! for browser tests that drive the photo flow through the real content routes. +//! The main binary (`src/main.rs`) is the actual Dioxus application. + +use wasm_bindgen::JsValue; +use wasm_bindgen::prelude::wasm_bindgen; + +include!(concat!(env!("OUT_DIR"), "/replica-tables.rs")); + +/// Schema SQL this build was compiled against (matches the browser-stack server). +pub const SCHEMA_SQL: &str = include_str!("../schema.sql"); + +/// The synced replica schema (worker replica, policy-split by build.rs). +pub const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); + +/// The local tier schema (device-private, attached, never synced). +pub const DEMO_FRONTEND_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")); + +/// The tab mirror schema: both tiers in the tab's main schema. +pub const DEMO_TAB_DDL: &str = concat!( + include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")), + "\n", + include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")), +); + +/// The demo server the DB worker connects upstream to. +pub const DEMO_WS_URL: &str = "ws://127.0.0.1:7777/"; + +/// The upstream subscription the DB worker registers. +pub const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; + +/// The extra upstream subscription for photos. +pub const PHOTO_QUERY: &str = "SELECT * FROM photos"; + +/// The OPFS file base for the worker's durable synced replica. +pub const DB_NAME: &str = "connetto-photo-dioxus.sqlite"; + +/// The registered caller-identity function connetto installs on every connection. +pub const CALLER_FUNCTION: &str = "current_app_user"; + +/// The schema version this build was compiled against. +#[must_use] +pub fn demo_schema_version() -> connetto_core::SchemaVersion { + connetto_core::SchemaVersion::from_source(SCHEMA_SQL) +} + +// The uuidv4 SQL function registered on every connection so the orders +// and photos DEFAULT (uuidv4()) mints a UUID on local writes. +#[diesel::declare_sql_function] +extern "SQL" { + /// Client-authored primary key: a 16-byte UUID v4, stored as a BLOB. + fn uuidv4() -> diesel::sql_types::Binary; +} + +/// The registrar connetto installs on every connection it opens for this app. +#[must_use] +pub fn uuidv4_functions() -> connetto_client::SqlFunctions { + connetto_client::SqlFunctions::new().with(std::sync::Arc::new( + |conn: &mut diesel::SqliteConnection| { + uuidv4_utils::register_impl_with_behavior( + conn, + diesel::sqlite::SqliteFunctionBehavior::INNOCUOUS, + rosetta_uuid::Uuid::new_v4, + ) + }, + )) +} + +/// The policy table map for this build, for `ClientConfig::with_policy_tables`. +#[must_use] +pub fn demo_policy_tables() -> connetto_client::PolicyTables { + connetto_client::PolicyTables::from_translation( + POLICY_TABLES.iter().copied(), + POLICY_VIEWS.iter().copied(), + ) +} + +/// DB worker entry point: boot the connetto DB tier with the photo config. +/// +/// The test's blob worker bootstrap imports this crate's wasm module and awaits this. +/// +/// # Errors +/// +/// A string describing the VFS, upstream connect, or subscribe failure. +#[wasm_bindgen] +pub async fn db_worker_photo_boot() -> Result<(), JsValue> { + connetto_web::logging::init_console(); + connetto_web::workers::boot_db_worker::( + &connetto_web::workers::DbWorkerConfig::new(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-photo-dioxus-hub-meta.sqlite") + .with_content_namespace("connetto-photo-content") + .with_sql_functions(uuidv4_functions()) + .with_policy_tables(demo_policy_tables()) + .with_caller_function(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-photo-dioxus-auth.sqlite"), + ) + .await + .map(drop) + .map_err(JsValue::from) +} diff --git a/examples/dioxus-web-demo/src/main.rs b/examples/dioxus-web-demo/src/main.rs index a8c9b9ee..58926897 100644 --- a/examples/dioxus-web-demo/src/main.rs +++ b/examples/dioxus-web-demo/src/main.rs @@ -27,6 +27,7 @@ //! `CONNETTO_READER_URL`, `DATABASE_URL`, `CONNETTO_BIND`, `CONNETTO_WRITABLE`, and //! `CONNETTO_PG_DDL_FILE`, then `dx serve --port 9912` from this directory. +use std::collections::HashMap; use std::rc::Rc; use std::{ cell::RefCell, @@ -40,8 +41,9 @@ use connetto_client::{ }; use connetto_core::messages::FatalErrorReason; use connetto_dioxus::use_live; +use connetto_file_core::{FileId, MimeClass}; use connetto_web::{ - MessageTransport, + MessageTransport, TabContent, TabResolved, auth::{LogoutOutcome, WorkerAuthConfig, deliver_login_code, request_logout, request_unsynced}, leader, locks, unlock::{AccountChoice, serve_account_choice}, @@ -64,18 +66,21 @@ const DEMO_WS_URL: &str = "ws://127.0.0.1:7777/"; /// yields the version the server advertises, so this build presents a matching /// version at handshake and is not rejected as stale. const SCHEMA_SQL: &str = include_str!("../schema.sql"); -/// The synced replica schema (worker first boot). Matches `schema.sql`. -const DEMO_SQLITE_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ - CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ - CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), content_id BLOB NOT NULL, content_state TEXT) STRICT;"; -/// The tab mirror schema: both tiers in the tab's main schema, because every -/// relayed patch applies to main. The hub, not the tab, keeps the tiers apart. -const DEMO_TAB_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ - CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ - CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), content_id BLOB NOT NULL, content_state TEXT) STRICT; \ +/// The synced replica schema (worker first boot, policy-split by build.rs from schema.sql + +/// policies.sql). The tab mirror uses a simpler non-split DDL below. +const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); +/// The tab mirror schema: simple tables in the tab's main schema. The hub does not use +/// policy views on the tab side; the server's CDC already filters rows to the user's identity. +const DEMO_TAB_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, owner_id TEXT NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ + CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, owner_id TEXT NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ + CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), owner_id TEXT NOT NULL, content_id BLOB NOT NULL, content_state TEXT) STRICT; \ CREATE TABLE notes (id INTEGER PRIMARY KEY NOT NULL, body TEXT) STRICT;"; /// The upstream subscription the worker registers. const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; +/// The extra upstream subscription for photos. +const PHOTO_QUERY: &str = "SELECT * FROM photos"; +/// SQLite function name a translated policy calls for the caller identity. +const CALLER_FUNCTION: &str = "current_app_user"; /// The OPFS file holding the worker's durable synced replica (base name; the /// worker appends the identity hash so each account gets its own encrypted file). const DB_NAME: &str = "connetto-relay.sqlite"; @@ -104,6 +109,7 @@ const EXPORT_FILE_NAME: &str = "connetto-local-data.zip"; diesel::table! { orders (id) { id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, quantity -> diesel::sql_types::BigInt, } } @@ -119,6 +125,7 @@ diesel::table! { photos (id) { id -> rosetta_uuid::sql_types::Uuid, order_id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, content_id -> diesel::sql_types::Binary, content_state -> Nullable, } @@ -129,6 +136,7 @@ diesel::table! { #[diesel(check_for_backend(diesel::sqlite::Sqlite))] struct Order { id: rosetta_uuid::Uuid, + owner_id: String, quantity: i64, } @@ -140,6 +148,17 @@ struct Note { body: String, } +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = photos)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Photo { + id: rosetta_uuid::Uuid, + order_id: rosetta_uuid::Uuid, + owner_id: String, + content_id: Vec, + content_state: Option, +} + // The synced key generator: `orders.id` bakes to `DEFAULT (uuidv4())`, so a // tab write omits the id and this registered function mints it. The impl is // `rosetta_uuid::Uuid::new_v4`, the same strongly typed key the `orders` @@ -409,18 +428,18 @@ async fn run_db_worker() -> Result<(), JsValue> { .with_frontend_ddl(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(uuidv4_functions()) .with_policy_tables(PolicyTables::from_translation( POLICY_TABLES.iter().copied(), POLICY_VIEWS.iter().copied(), )) + .with_caller_function(CALLER_FUNCTION) .with_auth(auth) .with_auth_db_name(AUTH_DB_NAME) - // Gate the replica with a passkey. leader::join installs serve_unlock, - // so the tab handler is in place before the worker can ask. .with_unlock(true) - // Ask the tab which account to sign in as when more than one is stored. .with_pick_account(true), ) .await?; @@ -474,6 +493,8 @@ fn glue_url() -> String { /// liveness lock (dropped on unmount, so the worker reaps this tab). struct Boot { client: ConnettoClient, + /// The tab's content lane, split off the transport before the client took it. + content: Rc>, /// Shared rather than owned, because the gate ceremony is awaited and the /// handle has to outlive the borrow that reaches it. membership: Rc, @@ -505,19 +526,15 @@ async fn boot_window() -> Result { let tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; let wire = format!("connetto-wire-{client_id}-boot"); workers::announce_tab(&wire).await?; - let transport = + let mut transport = MessageTransport::::with_peer_liveness(&wire, workers::DB_ALIVE_LOCK) .map_err(|err| JsValue::from_str(&err.to_string()))?; + let content = Rc::new(TabContent::new(&mut transport)); let config = ClientConfig::new(client_id.clone()) .with_schema_version(Some(connetto_core::SchemaVersion::from_source(SCHEMA_SQL))) .with_sql_functions(uuidv4_functions()) - .with_policy_tables(PolicyTables::from_translation( - POLICY_TABLES.iter().copied(), - POLICY_VIEWS.iter().copied(), - )) - // A low threshold so the free-up-space affordance reclaims after a - // modest deletion, rather than only once the freelist is a quarter of - // the file. Trimming still runs only when the pass is called. + // No with_policy_tables: the tab mirror uses the simple non-split DDL and the + // server's CDC already filters rows to the authenticated user's identity. .with_trim_threshold(5); let conn = ConnettoConnection::connect( transport, @@ -540,6 +557,7 @@ async fn boot_window() -> Result { spawn_local(pump); Ok(Boot { client, + content, membership: Rc::new(membership), _tab_lock: tab_lock, }) @@ -650,6 +668,10 @@ struct PickerActive(Signal); /// The account the user chose in the picker; set by button click, consumed by the chooser. #[derive(Clone, Copy)] struct AccountAnswer(Signal>); + +/// The tab's content lane, set once at boot and consumed by Dashboard. +#[derive(Clone, Copy)] +struct ContentSlot(Signal>>>); // Dioxus components are PascalCase by convention; the `rsx!` call sites name // them as elements, so keep the component name and silence the lint. #[allow(non_snake_case)] @@ -675,6 +697,7 @@ fn App() -> Element { let mut all_accounts: Signal> = use_signal(Vec::new); let mut picker_active: Signal = use_signal(|| false); let mut account_answer: Signal> = use_signal(|| None); + let mut content_slot: Signal>>> = use_signal(|| None); use_context_provider(|| client_slot); use_context_provider(|| status); @@ -686,6 +709,7 @@ fn App() -> Element { use_context_provider(|| AllAccounts(all_accounts)); use_context_provider(|| PickerActive(picker_active)); use_context_provider(|| AccountAnswer(account_answer)); + use_context_provider(|| ContentSlot(content_slot)); // Register the account chooser before the worker boots. For a single stored // credential the chooser returns immediately without blocking the UI. For @@ -820,6 +844,7 @@ fn App() -> Element { if let Ok(level) = workers::request_custody().await { custody_level.set(Some(level)); } + content_slot.set(Some(Rc::clone(&boot.content))); let mut events = boot.client.events(); client_slot.set(Some(boot.client.clone())); status.set("connected".to_owned()); @@ -1255,18 +1280,27 @@ fn Dashboard() -> Element { .read() .clone() .expect("Dashboard mounts only once the client is ready"); + let content = use_context::() + .0 + .read() + .clone() + .expect("Dashboard mounts after content is ready"); + let identity = use_context::>>() + .read() + .clone() + .unwrap_or_default(); let orders = use_live::<_, _, Order>(&client, orders::table.order(orders::id)); let notes = use_live::<_, _, Note>(&client, notes::table.order(notes::id)); + let photos = use_live::<_, _, Photo>(&client, photos::table.order(photos::id)); let order_rows = orders.value().read().clone(); let note_rows = notes.value().read().clone(); - // Aggregates are computed from the live rows: the relay hub does not serve - // aggregate subscriptions, and the tab mirror already holds every row the - // subscription covers, so the derived totals converge as the rows do. + let photo_rows = photos.value().read().clone(); let order_count = order_rows.len(); let order_sum: i64 = order_rows.iter().map(|row| row.quantity).sum(); let note_count = note_rows.len(); + let photo_count = photo_rows.len(); let order_view: Vec<(rosetta_uuid::Uuid, i64)> = order_rows .iter() .map(|row| (row.id, row.quantity)) @@ -1275,13 +1309,52 @@ fn Dashboard() -> Element { .iter() .map(|row| (row.id, row.body.clone())) .collect(); + // Resolved URLs for available photos; refetched each time the set of available + // photos changes so stale tickets never wedge the panel. + let mut photo_urls: Signal> = use_signal(HashMap::new); + { + let content = Rc::clone(&content); + let available: Vec<(rosetta_uuid::Uuid, Vec)> = photo_rows + .iter() + .filter(|p| p.content_state.as_deref() == Some("available")) + .map(|p| (p.id, p.content_id.clone())) + .collect(); + let key: Vec = available.iter().map(|(id, _)| *id).collect(); + use_effect(move || { + let _key = key.clone(); + let available = available.clone(); + let content = Rc::clone(&content); + spawn(async move { + let mut new_urls = HashMap::new(); + for (id, content_id) in available { + if let Ok(bytes) = <[u8; 32]>::try_from(content_id.as_slice()) { + let file_id = FileId::from_bytes(bytes); + if let TabResolved::Remote { url } = content.resolve(file_id).await { + new_urls.insert(id, url); + } + } + } + photo_urls.set(new_urls); + }); + }); + } + let photo_view: Vec<(rosetta_uuid::Uuid, String, Option)> = photo_rows + .iter() + .map(|p| { + let state = p + .content_state + .clone() + .unwrap_or_else(|| "pending".to_owned()); + let url = photo_urls.read().get(&p.id).cloned(); + (p.id, state, url) + }) + .collect(); let orders_error = orders.error().read().clone(); let notes_error = notes.error().read().clone(); + let photos_error = photos.error().read().clone(); let mut note_text = use_signal(String::new); - - // R26: the last export's outcome, so a failed one is not silent. let mut export_status: Signal> = use_signal(|| None); let mut import_status: Signal> = use_signal(|| None); @@ -1291,7 +1364,9 @@ fn Dashboard() -> Element { { let client = client.clone(); use_effect(move || { - let _covered = orders.value().read().len() + notes.value().read().len(); + let _covered = orders.value().read().len() + + notes.value().read().len() + + photos.value().read().len(); let client = client.clone(); spawn(async move { footprint.set(replica_footprint(&client).await); @@ -1305,6 +1380,10 @@ fn Dashboard() -> Element { let newest_order = order_rows.last().map(|order| order.id); let add_order_client = client.clone(); + let add_order_identity = identity.clone(); + let stage_client = client.clone(); + let stage_content = Rc::clone(&content); + let stage_identity = identity.clone(); let save_note_client = client; rsx! { @@ -1314,15 +1393,19 @@ fn Dashboard() -> Element { p { "count {order_count}, total quantity {order_sum}. Converges across every window through Postgres." } div { class: "row", button { + disabled: add_order_identity.is_empty(), onclick: move |_| { let client = add_order_client.clone(); + let identity = add_order_identity.clone(); spawn(async move { - // The DEFAULT mints the id, so the insert omits it. let quantity = fresh_quantity(); let result = client .with_conn(move |conn| { diesel::insert_into(orders::table) - .values(orders::quantity.eq(quantity)) + .values(( + orders::owner_id.eq(identity.as_str()), + orders::quantity.eq(quantity), + )) .execute(conn.conn()) }) .await; @@ -1370,6 +1453,107 @@ fn Dashboard() -> Element { } } } + div { class: "pane", + h2 { "photos " span { class: "badge synced", "synced" } } + p { "count {photo_count}. Pick an image file to stage and upload through the content route." } + div { class: "row", + input { + id: "photo-file-input", + r#type: "file", + accept: "image/*", + disabled: stage_identity.is_empty(), + onchange: move |_| { + let window = web_sys::window().expect("window"); + let input: web_sys::HtmlInputElement = window + .document().expect("document") + .get_element_by_id("photo-file-input").expect("photo input") + .unchecked_into(); + if let Some(files) = input.files() + && let Some(file) = files.get(0) + { + let client = stage_client.clone(); + let content = Rc::clone(&stage_content); + let identity = stage_identity.clone(); + spawn(async move { + let blob: web_sys::Blob = file.into(); + let si = identity.clone(); + let result = content + .stage(&blob, MimeClass::Jpeg, &client, move |conn, file_id| { + conn.transaction(|conn| { + let before_orders: std::collections::HashSet = + orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(si.as_str()), + orders::quantity.eq(1_i64), + )) + .execute(conn)?; + let order_id = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .find(|id| !before_orders.contains(id)) + .expect("order minted"); + let before_photos: std::collections::HashSet = + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(photos::table) + .values(( + photos::order_id.eq(order_id), + photos::owner_id.eq(si.as_str()), + photos::content_id.eq(file_id.as_bytes().to_vec()), + photos::content_state.eq::>(None), + )) + .execute(conn)?; + Ok::<_, diesel::result::Error>( + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .find(|id| !before_photos.contains(id)) + .expect("photo minted"), + ) + }) + }) + .await; + match result { + Ok(_) => { client.replay_pending().await.ok(); } + Err(err) => { + tracing::error!(error = %err, "photo stage failed"); + } + } + }); + } + } + } + } + if let Some(err) = photos_error { + p { style: "color:#b00;", "photos error: {err}" } + } + table { + thead { tr { th { "id" } th { "state" } th { "image" } } } + tbody { + for (id, state, url) in photo_view { + tr { key: "{id}", + td { "{id}" } + td { {state} } + td { + if let Some(url) = url { + img { src: url, style: "max-height:60px;" } + } + } + } + } + } + } + } div { class: "pane", h2 { "notes " span { class: "badge local", "device-only" } } p { "count {note_count}. Converges across this device's windows through the DB worker, never the server." } @@ -1424,7 +1608,7 @@ fn Dashboard() -> Element { } div { class: "pane", h2 { "retention " span { class: "badge trim", "R15" } } - p { "Replica mirror: {pages} pages (~{kb} KB), {free} free to reclaim. Covered rows: {order_count + note_count}." } + p { "Replica mirror: {pages} pages (~{kb} KB), {free} free to reclaim. Covered rows: {order_count + note_count + photo_count}." } p { "Ending a subscription evicts the rows no live subscription still covers, and the trimming pass hands the freed pages back to storage." } div { class: "row", button { diff --git a/examples/dioxus-web-demo/tests/photo_flow.rs b/examples/dioxus-web-demo/tests/photo_flow.rs new file mode 100644 index 00000000..90097ef5 --- /dev/null +++ b/examples/dioxus-web-demo/tests/photo_flow.rs @@ -0,0 +1,524 @@ +//! Photo surface browser tests for the dioxus web demo. +//! +//! Drives the real demo boot path (`db_worker_photo_boot`) through the browser +//! stack. Two tests run in this binary and are serialized by a shared Web Lock +//! to avoid OPFS conflicts between concurrent workers. +//! +//! COUPLING: `SchemaVersion` is a hash of the schema SQL source string. +//! `schema.sql` in this workspace MUST be byte-identical to +//! `examples/wasm-smoke/schema.sql`, which is what the browser stack server +//! is compiled against. The first test (`a_tab_order_reaches_the_replica…`) +//! acts as the permanent drift detector: it fails at the WebSocket handshake +//! the moment the two schemas diverge, long before any photo assertion runs. +//! Maintain the identity with `cp examples/wasm-smoke/schema.sql +//! examples/dioxus-web-demo/schema.sql` whenever wasm-smoke's schema changes. + +#![cfg(target_arch = "wasm32")] + +use connetto_client::dsl::Watchable; +use connetto_client::{ + ClientConfig, ClientEvent, ConnettoClient, ConnettoConnection, Grant, LiveQuery, Replica, +}; +use connetto_dioxus_web_demo::{ + CALLER_FUNCTION, DEMO_TAB_DDL, demo_policy_tables, demo_schema_version, uuidv4_functions, +}; +use connetto_file_core::{FileId, MimeClass}; +use connetto_web::auth::{ + Acquired, BrowserAuthenticator, IdbKeyStore, LOGIN_CHANNEL, LoginMessage, RefreshStore, + WorkerAuthConfig, deliver_login_code, +}; +use connetto_web::storage::{ReplicaStorage, device_key}; +use connetto_web::{MessageTransport, TabContent, TabResolved, locks, workers}; +use diesel::prelude::*; +use futures_channel::oneshot; +use js_sys::{Array, Uint8Array}; +use wasm_bindgen::JsCast; +use wasm_bindgen::prelude::*; +use wasm_bindgen_futures::{JsFuture, spawn_local}; +use wasm_bindgen_test::{wasm_bindgen_test, wasm_bindgen_test_configure}; +use web_sys::{ + BroadcastChannel, DedicatedWorkerGlobalScope, MessageEvent, Request, RequestInit, Worker, +}; + +wasm_bindgen_test_configure!(run_in_dedicated_worker); + +// Auth server coordinates matching the browser stack setup. +const AUTH_BASE: &str = "http://127.0.0.1:18099"; +const AUTH_LANDING: &str = "http://127.0.0.1:18099/dev/landing"; +const AUTH_PROVIDER: &str = "dev-idp"; +const AUTH_USERNAME: &str = "startup"; + +// Serializes tests in this binary to prevent OPFS worker conflicts. +const SUITE_LOCK: &str = "connetto-dioxus-photo-suite"; + +// --- Diesel table schema matching DEMO_TAB_DDL (policy-split, uses logical names) --- + +diesel::table! { + orders (id) { + id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, + quantity -> diesel::sql_types::BigInt, + } +} + +diesel::table! { + photos (id) { + id -> rosetta_uuid::sql_types::Uuid, + order_id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, + content_id -> diesel::sql_types::Binary, + content_state -> diesel::sql_types::Nullable, + } +} + +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = orders)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Order { + id: rosetta_uuid::Uuid, + owner_id: String, + quantity: i64, +} + +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = photos)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Photo { + id: rosetta_uuid::Uuid, + order_id: rosetta_uuid::Uuid, + owner_id: String, + content_id: Vec, + content_state: Option, +} + +// --- Helpers --- + +fn stage(msg: &str) { + web_sys::console::log_1(&msg.into()); +} + +fn relay_worker_breadcrumbs() { + use std::sync::atomic::{AtomicBool, Ordering}; + static INSTALLED: AtomicBool = AtomicBool::new(false); + if INSTALLED.swap(true, Ordering::Relaxed) { + return; + } + if let Ok(ch) = BroadcastChannel::new("connetto-debug") { + let cb = + wasm_bindgen::closure::Closure::::new(|e: MessageEvent| { + web_sys::console::log_1(&e.data()); + }); + ch.set_onmessage(Some(cb.as_ref().unchecked_ref())); + cb.forget(); + std::mem::forget(ch); + } +} + +fn play_the_tab() { + use std::sync::atomic::{AtomicBool, Ordering}; + static INSTALLED: AtomicBool = AtomicBool::new(false); + if INSTALLED.swap(true, Ordering::Relaxed) { + return; + } + let ch = BroadcastChannel::new(LOGIN_CHANNEL).expect("login channel"); + let cb = wasm_bindgen::closure::Closure::::new(|e: MessageEvent| { + let Some(text) = e.data().as_string() else { + return; + }; + let Ok(LoginMessage::Request { url }) = serde_json::from_str::(&text) else { + return; + }; + spawn_local(async move { + let (code, state) = walk_login(&url).await; + deliver_login_code(&code, &state).expect("deliver login code"); + }); + }); + ch.set_onmessage(Some(cb.as_ref().unchecked_ref())); + cb.forget(); + std::mem::forget(ch); +} + +async fn global_fetch_str(url: &str) -> web_sys::Response { + let promise = js_sys::global() + .dyn_into::() + .map(|w| w.fetch_with_str(url)) + .unwrap_or_else(|_| web_sys::window().expect("window").fetch_with_str(url)); + JsFuture::from(promise) + .await + .expect("fetch") + .dyn_into() + .expect("Response") +} + +async fn global_fetch_req(req: &Request) -> web_sys::Response { + let promise = js_sys::global() + .dyn_into::() + .map(|w| w.fetch_with_request(req)) + .unwrap_or_else(|_| web_sys::window().expect("window").fetch_with_request(req)); + JsFuture::from(promise) + .await + .expect("fetch req") + .dyn_into() + .expect("Response") +} + +async fn walk_login(login_url: &str) -> (String, String) { + let resp = global_fetch_str(login_url).await; + let form_url = resp.url(); + let init = RequestInit::new(); + init.set_method("POST"); + init.set_body(&JsValue::from_str(&format!("username={AUTH_USERNAME}"))); + let req = Request::new_with_str_and_init(&form_url, &init).expect("login request"); + req.headers() + .set("content-type", "application/x-www-form-urlencoded") + .expect("content-type header"); + let resp = global_fetch_req(&req).await; + assert!( + resp.ok(), + "login chain ended at {} with status {}", + resp.url(), + resp.status() + ); + let final_url = resp.url(); + let parsed = web_sys::Url::new(&final_url).expect("parse final url"); + let params = parsed.search_params(); + ( + params + .get("code") + .unwrap_or_else(|| panic!("no code in {final_url}")), + params + .get("state") + .unwrap_or_else(|| panic!("no state in {final_url}")), + ) +} + +async fn mint_session() -> (String, String) { + let storage = ReplicaStorage::install().await; + let keys = IdbKeyStore::open().await.expect("open key store"); + let device = device_key(&keys).await.expect("device key"); + static N: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + let n = N.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let db = format!("dioxus-photo-mint-{n}.sqlite"); + let store = RefreshStore::open(&storage.db_url(&db), &device).expect("refresh store"); + let auth = BrowserAuthenticator::new( + WorkerAuthConfig::new(AUTH_BASE, AUTH_PROVIDER, AUTH_LANDING), + None, + ); + let pending = match auth + .acquire::(&store) + .await + .expect("acquire session") + { + Acquired::NeedLogin(p) => p, + Acquired::Access(_) => panic!("fresh store cannot refresh silently"), + }; + let (code, state) = walk_login(&pending.login_url).await; + let session = auth + .complete::(&pending, &code, &state, &store) + .await + .expect("complete login"); + drop(store); + storage.delete_db(&db).ok(); + (session.access_token, session.user_id) +} + +fn glue_url() -> String { + let found = js_sys::eval( + r#"performance.getEntriesByType("resource").map(e=>e.name).find(n=>n.endsWith("_bg.wasm"))"#, + ) + .expect("resource entries") + .as_string() + .expect("wasm resource entry"); + let base = found.strip_suffix("_bg.wasm").expect("wasm suffix"); + format!("{base}.js") +} + +fn spawn_photo_worker(glue_url: &str) -> Worker { + let wasm_url = glue_url + .strip_suffix(".js") + .map_or_else(|| format!("{glue_url}_bg.wasm"), |b| format!("{b}_bg.wasm")); + let src = format!( + "const ch=new BroadcastChannel('connetto-debug');\n\ + try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n ch.postMessage('dioxus photo worker ready');\n await mod.db_worker_photo_boot();\n}}catch(e){{\n ch.postMessage('dioxus photo worker FAILED: '+e);\n throw e;\n}}", + g = glue_url, + w = wasm_url, + ); + let parts = Array::of1(&JsValue::from_str(&src)); + let opts = web_sys::BlobPropertyBag::new(); + opts.set_type("text/javascript"); + let blob = + web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("bootstrap blob"); + let url = web_sys::Url::create_object_url_with_blob(&blob).expect("object url"); + let worker_opts = web_sys::WorkerOptions::new(); + worker_opts.set_type(web_sys::WorkerType::Module); + worker_opts.set_name("connetto-dioxus-photo"); + let w = Worker::new_with_options(&url, &worker_opts).expect("spawn worker"); + web_sys::Url::revoke_object_url(&url).ok(); + w +} + +fn photo_bytes() -> Vec { + const SEED: [u8; 16] = [ + 0x10, 0x32, 0x54, 0x76, 0x98, 0xba, 0xdc, 0xfe, 0x01, 0x23, 0x45, 0x67, 0x89, 0xab, 0xcd, + 0xef, + ]; + let mut v = Vec::with_capacity(4096); + while v.len() < 4096 { + v.extend_from_slice(&SEED); + } + v.truncate(4096); + v +} + +fn blob_of(bytes: &[u8]) -> web_sys::Blob { + let arr = Uint8Array::from(bytes); + web_sys::Blob::new_with_u8_array_sequence(&Array::of1(&arr)).expect("blob from bytes") +} + +async fn fetch_bytes(url: &str) -> Vec { + let scope = js_sys::global() + .dyn_into::() + .expect("dedicated worker"); + let resp: web_sys::Response = JsFuture::from(scope.fetch_with_str(url)) + .await + .expect("fetch content") + .dyn_into() + .expect("Response"); + assert!( + resp.ok(), + "content fetch at {} returned {}", + resp.url(), + resp.status() + ); + let buf = JsFuture::from(resp.array_buffer().expect("array_buffer promise")) + .await + .expect("array buffer"); + Uint8Array::new(&buf).to_vec() +} + +async fn connect_tab( + client_id: &str, + token: String, + identity: &str, +) -> ( + TabContent, + ConnettoConnection>, +) { + let wire = format!("connetto-wire-{client_id}"); + workers::announce_tab(&wire).await.expect("announce tab"); + let mut transport = MessageTransport::::new(&wire).expect("transport"); + let content = TabContent::new(&mut transport); + let config = ClientConfig::new(client_id.to_owned()) + .with_login(Some(Grant::new(token))) + .with_schema_version(Some(demo_schema_version())) + .with_sql_functions(uuidv4_functions()) + .with_policy_tables(demo_policy_tables()) + .with_caller(CALLER_FUNCTION, identity); + let conn = ConnettoConnection::connect( + transport, + &Replica::in_memory(), + DEMO_TAB_DDL, + &config, + None, + ) + .await + .expect("tab connect"); + (content, conn) +} + +// --- Test 1: alignment proof — version agreement via an orders round trip --- + +#[wasm_bindgen_test] +async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { + relay_worker_breadcrumbs(); + play_the_tab(); + let _serial = locks::hold_lock(SUITE_LOCK).await; + let worker = spawn_photo_worker(&glue_url()); + workers::await_db_worker_ready(&[]) + .await + .expect("db worker ready"); + stage("dioxus demo worker booted (align)"); + + let (token, identity) = mint_session().await; + let client_id = rosetta_uuid::Uuid::new_v4().to_string(); + let _tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; + let (_content, mut conn) = connect_tab(&client_id, token, &identity).await; + conn.subscribe("align-orders", "SELECT * FROM orders") + .await + .expect("orders subscribe"); + loop { + let event = conn.pump_one().await.expect("pump"); + assert_ne!(event, ClientEvent::Closed, "connection closed early"); + if matches!(event, ClientEvent::SnapshotEnd { .. }) { + break; + } + } + stage("orders subscription ready (align)"); + + let (client, pump) = ConnettoClient::with_pump(conn); + let (done_tx, done_rx) = oneshot::channel::<()>(); + spawn_local(async move { + pump.await; + let _ = done_tx.send(()); + }); + + let mut live: LiveQuery = orders::table + .order(orders::id) + .select(Order::as_select()) + .live(&client) + .await + .expect("orders live query"); + + let qty = 77_i64; + let id_str = identity.clone(); + client + .with_conn(move |conn| { + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(id_str.as_str()), + orders::quantity.eq(qty), + )) + .execute(conn.conn()) + }) + .await + .expect("order insert"); + client.replay_pending().await.expect("replay pending"); + stage("order written (align)"); + + loop { + if live.rows().iter().any(|o| o.quantity == qty) { + break; + } + live.changed().await.expect("live refresh"); + } + stage("order arrived at replica — version agreement confirmed"); + + drop(live); + drop(client); + done_rx.await.expect("pump exited"); + worker.terminate(); +} + +// --- Test 2: photo round trip through stage, commit, resolve, and HTTP fetch --- + +#[wasm_bindgen_test] +async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { + relay_worker_breadcrumbs(); + play_the_tab(); + let _serial = locks::hold_lock(SUITE_LOCK).await; + let worker = spawn_photo_worker(&glue_url()); + workers::await_db_worker_ready(&[]) + .await + .expect("db worker ready"); + stage("dioxus demo worker booted (photo)"); + + let (token, identity) = mint_session().await; + let client_id = rosetta_uuid::Uuid::new_v4().to_string(); + let _tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; + let (content, mut conn) = connect_tab(&client_id, token, &identity).await; + conn.subscribe("photo-test-photos", "SELECT * FROM photos") + .await + .expect("photos subscribe"); + loop { + let event = conn.pump_one().await.expect("pump"); + assert_ne!(event, ClientEvent::Closed, "connection closed early"); + if matches!(event, ClientEvent::SnapshotEnd { .. }) { + break; + } + } + stage("photos subscription ready"); + + let (client, pump) = ConnettoClient::with_pump(conn); + let (done_tx, done_rx) = oneshot::channel::<()>(); + spawn_local(async move { + pump.await; + let _ = done_tx.send(()); + }); + + let mut live: LiveQuery = photos::table + .order(photos::id) + .select(Photo::as_select()) + .live(&client) + .await + .expect("photo live query"); + + let bytes = photo_bytes(); + let blob = blob_of(&bytes); + let (file_id, photo_id) = content + .stage(&blob, MimeClass::Jpeg, &client, |conn, file_id| { + conn.transaction(|conn| { + let before_orders: std::collections::HashSet = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(identity.as_str()), + orders::quantity.eq(1_i64), + )) + .execute(conn)?; + let order_id = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .find(|id| !before_orders.contains(id)) + .expect("minted order id"); + let before_photos: std::collections::HashSet = photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(photos::table) + .values(( + photos::order_id.eq(order_id), + photos::owner_id.eq(identity.as_str()), + photos::content_id.eq(file_id.as_bytes().to_vec()), + photos::content_state.eq::>(None), + )) + .execute(conn)?; + Ok::( + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .find(|id| !before_photos.contains(id)) + .expect("minted photo id"), + ) + }) + }) + .await + .expect("stage photo and insert rows"); + client.replay_pending().await.expect("send mutation"); + assert_eq!(file_id, FileId::from_chunks([bytes.as_slice()])); + stage("photo staged"); + + let available = loop { + if let Some(p) = live + .rows() + .iter() + .find(|p| p.id == photo_id && p.content_state.as_deref() == Some("available")) + .cloned() + { + break p; + } + live.changed().await.expect("live refresh"); + }; + assert_eq!(available.owner_id, identity); + assert_eq!(available.content_id, file_id.as_bytes().to_vec()); + stage("photo available"); + + let url = match content.resolve(file_id).await { + TabResolved::Remote { url } => url, + TabResolved::Local { .. } => panic!("uploaded photo must resolve to server"), + TabResolved::Unavailable => panic!("available photo must resolve"), + }; + assert_eq!(fetch_bytes(&url).await, bytes); + stage("photo bytes verified via HTTP fetch"); + + drop(live); + drop(content); + drop(client); + done_rx.await.expect("pump exited"); + worker.terminate(); +} diff --git a/examples/wasm-smoke/Cargo.lock b/examples/wasm-smoke/Cargo.lock index 72cc793c..e7381d26 100644 --- a/examples/wasm-smoke/Cargo.lock +++ b/examples/wasm-smoke/Cargo.lock @@ -800,7 +800,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1932,7 +1932,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2247,7 +2247,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2734,10 +2734,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3367,7 +3367,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/examples/yew-web-demo/Cargo.lock b/examples/yew-web-demo/Cargo.lock index 6d780d16..b53e018e 100644 --- a/examples/yew-web-demo/Cargo.lock +++ b/examples/yew-web-demo/Cargo.lock @@ -104,6 +104,17 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "async-trait" +version = "0.1.92" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "82f6aeea286b8eb4dd3431a1be1b59d290ace00f5bfd8e2a159bc2a05e2c1667" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.5", +] + [[package]] name = "atomic-waker" version = "1.1.2" @@ -240,6 +251,12 @@ dependencies = [ "serde", ] +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cc" version = "1.4.5" @@ -437,6 +454,7 @@ dependencies = [ "connetto-client", "connetto-core", "connetto-file-client", + "connetto-file-core", "diesel", "diesel-sqlite-session", "futures-channel", @@ -479,16 +497,21 @@ version = "0.0.0" dependencies = [ "connetto-client", "connetto-core", + "connetto-file-client", + "connetto-file-core", "connetto-web", "connetto-yew", "diesel", + "futures-channel", "js-sys", "libsqlite3-sys", "pg2sqlite", "rosetta-uuid", + "serde_json", "tracing", "wasm-bindgen", "wasm-bindgen-futures", + "wasm-bindgen-test", "web-sys", "yew", ] @@ -805,7 +828,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1884,6 +1907,16 @@ version = "2.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cf8baf1c55e62ffcace7a9f06f4bd9cd3f0c4beb022d3b367256b91b87513d98" +[[package]] +name = "minicov" +version = "0.3.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4869b6a491569605d66d3952bcdf03df789e5b536e5f0cf7758a7f08a55ae24d" +dependencies = [ + "cc", + "walkdir", +] + [[package]] name = "minimal-lexical" version = "0.2.1" @@ -1921,6 +1954,15 @@ dependencies = [ "minimal-lexical", ] +[[package]] +name = "nu-ansi-term" +version = "0.50.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "num-bigint" version = "0.4.8" @@ -1981,6 +2023,12 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "oorandom" +version = "11.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" + [[package]] name = "opaque-debug" version = "0.3.1" @@ -2237,7 +2285,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2552,7 +2600,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2602,6 +2650,15 @@ version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +[[package]] +name = "same-file" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" +dependencies = [ + "winapi-util", +] + [[package]] name = "seahash" version = "4.1.0" @@ -3040,10 +3097,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3542,6 +3599,16 @@ dependencies = [ "thiserror 2.0.20", ] +[[package]] +name = "walkdir" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29790946404f91d9c5d06f9874efddea1dc06c5efe94541a7d6863108e3a5e4b" +dependencies = [ + "same-file", + "winapi-util", +] + [[package]] name = "want" version = "0.3.1" @@ -3621,6 +3688,45 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "wasm-bindgen-test" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45863ef0bef521c12124eb39d9a38513c47db75f22e503c06beacd40afeb35db" +dependencies = [ + "async-trait", + "cast", + "js-sys", + "libm", + "minicov", + "nu-ansi-term", + "num-traits", + "oorandom", + "serde", + "serde_json", + "wasm-bindgen", + "wasm-bindgen-futures", + "wasm-bindgen-test-macro", + "wasm-bindgen-test-shared", +] + +[[package]] +name = "wasm-bindgen-test-macro" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8c89dcab8b516b6b603baca9d550b7282d68fcc7f367e3956cff7ebf406a3f12" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.5", +] + +[[package]] +name = "wasm-bindgen-test-shared" +version = "0.2.128" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f37b4f992cebe528ef34964ae69681ac0fe7080071e7298e46008f9d380302af" + [[package]] name = "wasm-streams" version = "0.4.2" @@ -3663,6 +3769,15 @@ dependencies = [ "rustls-pki-types", ] +[[package]] +name = "winapi-util" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "windows-core" version = "0.62.2" diff --git a/examples/yew-web-demo/Cargo.toml b/examples/yew-web-demo/Cargo.toml index c7e8f70c..48ff6c83 100644 --- a/examples/yew-web-demo/Cargo.toml +++ b/examples/yew-web-demo/Cargo.toml @@ -7,6 +7,9 @@ rust-version = "1.88" license = "MIT" publish = false +[lib] +crate-type = ["cdylib", "rlib"] + [dependencies] # The browser platform: transports, leader election, locks, the relay hub, # and the DB worker orchestration. The demo supplies its schema and baked @@ -16,6 +19,8 @@ connetto-web = { path = "../../crates/connetto-web" } # Default features off: this is a wasm build with a browser transport. connetto-client = { path = "../../crates/connetto-client", default-features = false } connetto-core = { path = "../../crates/connetto-core" } +# FileId and MimeClass for the photo staging closure. +connetto-file-core = { path = "../../crates/connetto-file-core" } # The one live-query hook, bound to the component lifecycle. connetto-yew = { path = "../../crates/connetto-yew" } diesel = { version = "2", features = ["sqlite"] } @@ -35,7 +40,9 @@ tracing = { version = "0.1", default-features = false, features = ["std"] } # `Event`, `File` and `FileList` back the import file-input change handler. web-sys = { version = "0.3", features = [ "Blob", + "BlobPropertyBag", "BroadcastChannel", + "DedicatedWorkerGlobalScope", "Document", "Element", "Event", @@ -47,12 +54,26 @@ web-sys = { version = "0.3", features = [ "HtmlInputElement", "Location", "MessageEvent", + "Request", + "RequestInit", + "Response", "Storage", + "UrlSearchParams", "Url", "Window", + "Worker", + "WorkerGlobalScope", + "WorkerOptions", + "WorkerType", "console", ] } +[dev-dependencies] +wasm-bindgen-test = "0.3" +futures-channel = "0.3" +serde_json = "1" +connetto-file-client = { path = "../../crates/connetto-file-client" } + # The build script translates schema.sql and frontend.sql through pg2sqlite and # bakes a template per tier, so the app ships every tier's schema pre-applied # and never executes DDL at startup. Same pipeline as the dioxus demo. diff --git a/examples/yew-web-demo/build.rs b/examples/yew-web-demo/build.rs index 666c0923..35aab877 100644 --- a/examples/yew-web-demo/build.rs +++ b/examples/yew-web-demo/build.rs @@ -110,13 +110,17 @@ fn write_policy_tables(documents: &[&str], views: &[String], out: &std::path::Pa fn main() { println!("cargo::rerun-if-changed=schema.sql"); println!("cargo::rerun-if-changed=frontend.sql"); + println!("cargo::rerun-if-changed=policies.sql"); let out_dir = std::path::PathBuf::from(std::env::var("OUT_DIR").expect("cargo sets OUT_DIR")); - let synced_views = translate(&["schema.sql"], &out_dir.join("replica-ddl.sql")); + let synced_views = translate( + &["schema.sql", "policies.sql"], + &out_dir.join("replica-ddl.sql"), + ); translate(&["frontend.sql"], &out_dir.join("frontend-ddl.sql")); // The synced tier only: the local tier is a separate database, attached // under its own schema, and the check the map feeds reads `main`. write_policy_tables( - &["schema.sql"], + &["schema.sql", "policies.sql"], &synced_views, &out_dir.join("replica-tables.rs"), ); diff --git a/examples/yew-web-demo/policies.sql b/examples/yew-web-demo/policies.sql new file mode 100644 index 00000000..f170e654 --- /dev/null +++ b/examples/yew-web-demo/policies.sql @@ -0,0 +1,26 @@ +-- The row-level security the backend enforces on the synced tables, kept apart +-- from schema.sql because the two reach the server as separate documents: +-- schema.sql feeds CONNETTO_PG_DDL and is what clients sync, this file feeds +-- CONNETTO_PG_POLICIES and is what the authorization model is derived from. +-- Apply both to Postgres, this one last, after schema.sql, the file server +-- DDL, roles.sql and content.sql. +-- build.rs translates the pair together, which is what splits the replica's +-- orders into a backing table, a view of the logical name, and INSTEAD OF +-- triggers. The caller is read from app.user_id, which the server binds per +-- transaction and the replica answers with the registered current_app_user() +-- function, so both ends compare against the same identity. +ALTER TABLE orders ENABLE ROW LEVEL SECURITY; +CREATE POLICY orders_p ON orders USING (owner_id = current_setting('app.user_id', true)); -- NOSONAR S1192, SQL DDL has no constants for the caller setting the three policies share + +-- The same shape on the composite-key table, so its replica half is split the +-- same way and its INSTEAD OF triggers have to match a row on two key columns +-- rather than one. +ALTER TABLE order_lines ENABLE ROW LEVEL SECURITY; +CREATE POLICY order_lines_p ON order_lines USING (owner_id = current_setting('app.user_id', true)); + +-- The photo rows are visible to the owner of the order they hang off, and the +-- owner is repeated on the row as it is on order_lines, so the comparison +-- settles from the row itself. The file server's visibility function consults +-- exactly this table under the caller's identity. +ALTER TABLE photos ENABLE ROW LEVEL SECURITY; +CREATE POLICY photos_p ON photos USING (owner_id = current_setting('app.user_id', true)); diff --git a/examples/yew-web-demo/schema.sql b/examples/yew-web-demo/schema.sql index 689b4e2c..143b0874 100644 --- a/examples/yew-web-demo/schema.sql +++ b/examples/yew-web-demo/schema.sql @@ -1,31 +1,33 @@ -- The one source of truth for the demo: the Postgres dialect schema the --- backend owns. build.rs translates this through pg2sqlite and bakes the --- replica template database the app ships. The connetto-server for this demo --- must be started with this same schema in CONNETTO_PG_DDL and must list --- every synced table in CONNETTO_WRITABLE. Apply in this order: this file, +-- backend owns. build.rs translates this through pg2sqlite into the SQLite DDL +-- a first boot applies. The connetto-server for this demo must be started with +-- this same schema in CONNETTO_PG_DDL and must list every synced table in +-- CONNETTO_WRITABLE. Apply in this order: this file, -- connetto_file_server::DEPLOYMENT_DDL, roles.sql (the non-owner role --- required by CONNETTO_READER_URL), then content.sql. --- The server also requires CONNETTO_AUTH, CONNETTO_AUTH_BIND, and the --- CONNETTO_OIDC_* variables written by the dev IdP (see dev_idp.rs). +-- required by CONNETTO_READER_URL), content.sql, then policies.sql. -- The key default is load-bearing on the client rather than here: build.rs -- translates it through pg2sqlite into the replica's own DEFAULT (uuidv4()), -- which mints the key when a local write omits it. Both ends mint version 4. -- The quantity is non-null because every client schema already declares it so. -CREATE TABLE orders (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), quantity BIGINT NOT NULL CHECK (quantity >= 0)); +-- owner_id carries who a row belongs to, which policies.sql compares against +-- the caller. It has no default: pg2sqlite maps current_setting only inside a +-- policy expression, so a default naming the caller would translate into a +-- call the replica cannot resolve, and every write names the owner instead. +CREATE TABLE orders (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), owner_id TEXT NOT NULL, quantity BIGINT NOT NULL CHECK (quantity >= 0)); -- The lines of an order, keyed by the order and the line number together. It is --- the one table here whose key spans two columns, which the replica's own schema --- and every key connetto encodes on the wire have to carry as a pair. -CREATE TABLE order_lines ( - order_id UUID NOT NULL REFERENCES orders(id), - line_no INTEGER NOT NULL, - quantity BIGINT NOT NULL CHECK (quantity >= 0), - PRIMARY KEY (order_id, line_no) -); +-- the one table here whose key spans two columns, and the translation below +-- splits it like any policy-bearing table, so its INSTEAD OF triggers match a +-- row on both key columns rather than one. owner_id repeats rather than being +-- read through the parent order, so the policy settles from the row itself, +-- which is what keeps the change path free of a round trip. +CREATE TABLE order_lines (order_id UUID NOT NULL REFERENCES orders(id), line_no INTEGER NOT NULL, owner_id TEXT NOT NULL, quantity BIGINT NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)); -- The photo entry: metadata for one file's bytes, attached to an order. -- content_id is the BLAKE3 identity the file server stores and serves under, -- and content_state stays null until the file server's commit writes -- `available`, so the placeholder condition is "not available" and the --- availability flip arrives as an ordinary synced column change. -CREATE TABLE photos (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), order_id UUID NOT NULL REFERENCES orders(id), content_id BYTEA NOT NULL, content_state TEXT); +-- availability flip arrives as an ordinary synced column change. owner_id +-- repeats from the parent order as it does on order_lines, so the visibility +-- policy settles from the row itself. +CREATE TABLE photos (id UUID PRIMARY KEY DEFAULT gen_random_uuid(), order_id UUID NOT NULL REFERENCES orders(id), owner_id TEXT NOT NULL, content_id BYTEA NOT NULL, content_state TEXT); diff --git a/examples/yew-web-demo/src/lib.rs b/examples/yew-web-demo/src/lib.rs new file mode 100644 index 00000000..6434255b --- /dev/null +++ b/examples/yew-web-demo/src/lib.rs @@ -0,0 +1,114 @@ +//! Browser photo worker infrastructure for the yew web demo. +//! +//! This lib target exposes the `db_worker_photo_boot` wasm-bindgen entry point +//! for browser tests that drive the photo flow through the real content routes. +//! The main binary (`src/main.rs`) is the actual Yew application. + +use wasm_bindgen::JsValue; +use wasm_bindgen::prelude::wasm_bindgen; + +include!(concat!(env!("OUT_DIR"), "/replica-tables.rs")); + +/// Schema SQL this build was compiled against (matches the browser-stack server). +pub const SCHEMA_SQL: &str = include_str!("../schema.sql"); + +/// The synced replica schema (worker replica, policy-split by build.rs). +pub const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); + +/// The local tier schema (device-private, attached, never synced). +pub const DEMO_FRONTEND_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")); + +/// The tab mirror schema: both tiers in the tab's main schema. +pub const DEMO_TAB_DDL: &str = concat!( + include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")), + "\n", + include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")), +); + +/// The demo server the DB worker connects upstream to. +pub const DEMO_WS_URL: &str = "ws://127.0.0.1:7777/"; + +/// The upstream subscription the DB worker registers. +pub const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; + +/// The extra upstream subscription for photos. +pub const PHOTO_QUERY: &str = "SELECT * FROM photos"; + +/// The OPFS file base for the worker's durable synced replica. +pub const DB_NAME: &str = "connetto-photo-yew.sqlite"; + +/// The registered caller-identity function connetto installs on every connection. +pub const CALLER_FUNCTION: &str = "current_app_user"; + +/// The schema version this build was compiled against. +#[must_use] +pub fn demo_schema_version() -> connetto_core::SchemaVersion { + connetto_core::SchemaVersion::from_source(SCHEMA_SQL) +} + +// The uuidv4 SQL function registered on every connection so the orders +// and photos DEFAULT (uuidv4()) mints a UUID on local writes. +#[diesel::declare_sql_function] +extern "SQL" { + /// Client-authored primary key: a 16-byte UUID v4, stored as a BLOB. + fn uuidv4() -> diesel::sql_types::Binary; +} + +/// The registrar connetto installs on every connection it opens for this app. +#[must_use] +pub fn uuidv4_functions() -> connetto_client::SqlFunctions { + connetto_client::SqlFunctions::new().with(std::sync::Arc::new( + |conn: &mut diesel::SqliteConnection| { + uuidv4_utils::register_impl_with_behavior( + conn, + diesel::sqlite::SqliteFunctionBehavior::INNOCUOUS, + rosetta_uuid::Uuid::new_v4, + ) + }, + )) +} + +/// The policy table map for this build, for `ClientConfig::with_policy_tables`. +#[must_use] +pub fn demo_policy_tables() -> connetto_client::PolicyTables { + connetto_client::PolicyTables::from_translation( + POLICY_TABLES.iter().copied(), + POLICY_VIEWS.iter().copied(), + ) +} + +/// DB worker entry point: boot the connetto DB tier with the photo config. +/// +/// The test's blob worker bootstrap imports this crate's wasm module and awaits this. +/// +/// # Errors +/// +/// A string describing the VFS, upstream connect, or subscribe failure. +#[wasm_bindgen] +pub async fn db_worker_photo_boot() -> Result<(), JsValue> { + connetto_web::logging::init_console(); + connetto_web::workers::boot_db_worker::( + &connetto_web::workers::DbWorkerConfig::new(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-photo-yew-hub-meta.sqlite") + .with_content_namespace("connetto-photo-content") + .with_sql_functions(uuidv4_functions()) + .with_policy_tables(demo_policy_tables()) + .with_caller_function(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-photo-yew-auth.sqlite"), + ) + .await + .map(drop) + .map_err(JsValue::from) +} diff --git a/examples/yew-web-demo/src/main.rs b/examples/yew-web-demo/src/main.rs index 131cda38..b3769469 100644 --- a/examples/yew-web-demo/src/main.rs +++ b/examples/yew-web-demo/src/main.rs @@ -29,6 +29,7 @@ //! the served URL in several windows. use std::cell::{Cell, RefCell}; +use std::collections::HashMap; use std::rc::Rc; use connetto_client::reconnect::ReconnectPolicy; @@ -38,9 +39,12 @@ use connetto_client::{ Replica, }; use connetto_core::custody::Custody; +use connetto_file_core::{FileId, MimeClass}; use connetto_web::auth::{PendingWork, WorkerAuthConfig}; use connetto_web::unlock::{AccountChoice, serve_account_choice}; -use connetto_web::{MessageTransport, deliver_login_code, leader, locks, workers}; +use connetto_web::{ + MessageTransport, TabContent, TabResolved, deliver_login_code, leader, locks, workers, +}; use connetto_yew::use_live; use diesel::prelude::*; use wasm_bindgen::closure::Closure; @@ -64,18 +68,21 @@ const AUTH_ORIGIN: &str = "http://127.0.0.1:18081"; /// yields the version the server advertises, so this build presents a matching /// version at handshake and is not rejected as stale. const SCHEMA_SQL: &str = include_str!("../schema.sql"); -/// The synced replica schema (worker first boot). Matches `schema.sql`. -const DEMO_SQLITE_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ - CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ - CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), content_id BLOB NOT NULL, content_state TEXT) STRICT;"; -/// The tab mirror schema: both tiers in the tab's main schema, because every -/// relayed patch applies to main. The hub, not the tab, keeps the tiers apart. -const DEMO_TAB_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ - CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ - CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), content_id BLOB NOT NULL, content_state TEXT) STRICT; \ +/// The synced replica schema (worker first boot, policy-split by build.rs from schema.sql + +/// policies.sql). The tab mirror uses a simpler non-split DDL below. +const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); +/// The tab mirror schema: simple tables in the tab's main schema. The hub does not use +/// policy views on the tab side; the server's CDC already filters rows by the user's identity. +const DEMO_TAB_DDL: &str = "CREATE TABLE orders (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, owner_id TEXT NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0)) STRICT; \ + CREATE TABLE order_lines (order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), line_no INTEGER NOT NULL, owner_id TEXT NOT NULL, quantity INTEGER NOT NULL CHECK (quantity >= 0), PRIMARY KEY (order_id, line_no)) STRICT; \ + CREATE TABLE photos (id BLOB PRIMARY KEY DEFAULT (uuidv4()) CHECK (length(id) = 16) NOT NULL, order_id BLOB NOT NULL REFERENCES orders(id) CHECK (length(order_id) = 16), owner_id TEXT NOT NULL, content_id BLOB NOT NULL, content_state TEXT) STRICT; \ CREATE TABLE notes (id INTEGER PRIMARY KEY NOT NULL, body TEXT) STRICT;"; /// The upstream subscription the worker registers. const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; +/// The extra upstream subscription for photos. +const PHOTO_QUERY: &str = "SELECT * FROM photos"; +/// SQLite function name a translated policy calls for the caller identity. +const CALLER_FUNCTION: &str = "current_app_user"; /// The OPFS file holding the worker's durable synced replica. const DB_NAME: &str = "connetto-relay.sqlite"; /// The shared leader lock every window of this app races. @@ -98,6 +105,7 @@ const EXPORT_FILE_NAME: &str = "connetto-local-data.zip"; diesel::table! { orders (id) { id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, quantity -> diesel::sql_types::BigInt, } } @@ -113,6 +121,7 @@ diesel::table! { photos (id) { id -> rosetta_uuid::sql_types::Uuid, order_id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, content_id -> diesel::sql_types::Binary, content_state -> Nullable, } @@ -123,6 +132,7 @@ diesel::table! { #[diesel(check_for_backend(diesel::sqlite::Sqlite))] struct Order { id: rosetta_uuid::Uuid, + owner_id: String, quantity: i64, } @@ -134,6 +144,17 @@ struct Note { body: String, } +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = photos)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Photo { + id: rosetta_uuid::Uuid, + order_id: rosetta_uuid::Uuid, + owner_id: String, + content_id: Vec, + content_state: Option, +} + // The synced key generator: `orders.id` bakes to `DEFAULT (uuidv4())`, so a // tab write omits the id and this registered function mints it. The impl is // `rosetta_uuid::Uuid::new_v4`, the same strongly typed key the `orders` schema @@ -308,12 +329,15 @@ async fn run_db_worker() -> Result<(), JsValue> { .with_frontend_ddl(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(uuidv4_functions()) .with_policy_tables(PolicyTables::from_translation( POLICY_TABLES.iter().copied(), POLICY_VIEWS.iter().copied(), )) + .with_caller_function(CALLER_FUNCTION) .with_auth(auth) .with_auth_db_name("connetto-auth.sqlite") .with_unlock(true) @@ -425,6 +449,8 @@ fn glue_url() -> String { /// liveness lock (dropped on unmount, so the worker reaps this tab). struct Boot { client: ConnettoClient, + /// The tab's content lane, split off the transport before the client took it. + content: Rc>, _membership: leader::Membership, _tab_lock: locks::HeldLock, /// Passkey custody level as of this boot, read from the worker after it settled. @@ -458,19 +484,16 @@ async fn boot_window() -> Result { let tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; let wire = format!("connetto-wire-{client_id}-boot"); workers::announce_tab(&wire).await?; - let transport = + let mut transport = MessageTransport::::with_peer_liveness(&wire, workers::DB_ALIVE_LOCK) .map_err(|err| JsValue::from_str(&err.to_string()))?; + let content = Rc::new(TabContent::new(&mut transport)); let config = ClientConfig::new(client_id.clone()) .with_schema_version(Some(connetto_core::SchemaVersion::from_source(SCHEMA_SQL))) .with_sql_functions(uuidv4_functions()) - .with_policy_tables(PolicyTables::from_translation( - POLICY_TABLES.iter().copied(), - POLICY_VIEWS.iter().copied(), - )) - // A low threshold so the free-up-space affordance reclaims after a - // modest deletion, rather than only once the freelist is a quarter of - // the file. Trimming still runs only when the pass is called. + // No with_policy_tables: the tab mirror uses the simple non-split DDL and the + // server's CDC already filters rows to the authenticated user's identity. + // A low threshold so the free-up-space affordance reclaims after a modest deletion. .with_trim_threshold(5); let conn = ConnettoConnection::connect( transport, @@ -493,6 +516,7 @@ async fn boot_window() -> Result { spawn_local(pump); Ok(Boot { client, + content, _membership: membership, _tab_lock: tab_lock, custody, @@ -578,6 +602,16 @@ impl PartialEq for ClientHandle { } } +/// A cheaply clonable, identity-compared content handle for Dashboard props. +#[derive(Clone)] +struct ContentHandle(Rc>); + +impl PartialEq for ContentHandle { + fn eq(&self, other: &Self) -> bool { + Rc::ptr_eq(&self.0, &other.0) + } +} + /// Async helper: queries unsynced work and updates the expiry warning. /// Exits early if `exp_secs` is no longer the latest session. async fn compute_expiry_for_session( @@ -782,6 +816,7 @@ fn use_provider_listener() { #[hook] fn use_boot_window_effect( client: UseStateHandle>, + content: UseStateHandle>, status: UseStateHandle, boot_hold: Rc>>, custody: UseStateHandle, @@ -791,6 +826,7 @@ fn use_boot_window_effect( match boot_window().await { Ok(boot) => { custody.set(boot.custody); + content.set(Some(ContentHandle(Rc::clone(&boot.content)))); let mut events = boot.client.events(); client.set(Some(ClientHandle(Rc::new(boot.client.clone())))); status.set("connected".to_owned()); @@ -811,6 +847,8 @@ fn use_boot_window_effect( #[function_component(App)] fn app() -> Html { let client = use_state(|| None::); + // The content lane split off the transport at boot; held here so Dashboard can use it. + let content = use_state(|| None::); let status = use_state(|| "connecting to the connetto stack".to_owned()); // The boot tokens live as long as the app: dropping them on page close // resigns leadership and frees the tab lock, which reaps this tab. @@ -865,6 +903,7 @@ fn app() -> Html { use_boot_window_effect( client.clone(), + content.clone(), status.clone(), boot_hold.clone(), custody.clone(), @@ -1157,8 +1196,9 @@ fn app() -> Html { } }; - let dashboard = if let Some(handle) = &*client { - html! { } + let dashboard = if let (Some(handle), Some(content_handle)) = (&*client, &*content) { + let id = (*identity).clone().unwrap_or_default(); + html! { } } else { html! {

{ "Connecting to the DB worker and the connetto stack..." }

} }; @@ -1182,20 +1222,55 @@ fn app() -> Html { #[derive(Properties, PartialEq)] struct DashboardProps { client: ClientHandle, + content: ContentHandle, + identity: String, } #[function_component(Dashboard)] fn dashboard(props: &DashboardProps) -> Html { let client = (*props.client.0).clone(); + let identity = props.identity.clone(); let orders = use_live::<_, _, Order>(&client, orders::table.order(orders::id)); let notes = use_live::<_, _, Note>(&client, notes::table.order(notes::id)); + let photos = use_live::<_, _, Photo>(&client, photos::table.order(photos::id)); let order_rows = orders.value(); let note_rows = notes.value(); + let photo_rows = photos.value(); let order_count = order_rows.len(); let order_sum: i64 = order_rows.iter().map(|row| row.quantity).sum(); let note_count = note_rows.len(); + let photo_count = photo_rows.len(); + + // Resolved signed URLs for available photos. Refetched (not cached) whenever the + // set of available photos changes, so a stale ticket never wedges the panel. + let photo_urls = use_state(HashMap::::new); + { + let content = Rc::clone(&props.content.0); + let photo_urls = photo_urls.clone(); + let available: Vec<(rosetta_uuid::Uuid, Vec)> = photo_rows + .iter() + .filter(|p| p.content_state.as_deref() == Some("available")) + .map(|p| (p.id, p.content_id.clone())) + .collect(); + let key: Vec = available.iter().map(|(id, _)| *id).collect(); + use_effect_with(key, move |_| { + spawn_local(async move { + let mut new_urls = HashMap::new(); + for (id, content_id) in available { + if let Ok(bytes) = <[u8; 32]>::try_from(content_id.as_slice()) { + let file_id = FileId::from_bytes(bytes); + if let TabResolved::Remote { url } = content.resolve(file_id).await { + new_urls.insert(id, url); + } + } + } + photo_urls.set(new_urls); + }); + || () + }); + } let note_text = use_state(String::new); @@ -1205,9 +1280,8 @@ fn dashboard(props: &DashboardProps) -> Html { { let client = client.clone(); let footprint = footprint.clone(); - let order_count = order_rows.len(); - let note_count = note_rows.len(); - use_effect_with((order_count, note_count), move |_| { + let covered = order_count + note_count + photo_count; + use_effect_with(covered, move |_| { let client = client.clone(); let footprint = footprint.clone(); spawn_local(async move { @@ -1223,14 +1297,19 @@ fn dashboard(props: &DashboardProps) -> Html { let add_order = { let client = client.clone(); + let identity = identity.clone(); Callback::from(move |_| { let client = client.clone(); + let identity = identity.clone(); spawn_local(async move { let quantity = fresh_quantity(); let result = client .with_conn(move |conn| { diesel::insert_into(orders::table) - .values(orders::quantity.eq(quantity)) + .values(( + orders::owner_id.eq(identity.as_str()), + orders::quantity.eq(quantity), + )) .execute(conn.conn()) }) .await; @@ -1298,6 +1377,84 @@ fn dashboard(props: &DashboardProps) -> Html { }) }; + // Pick a photo file, stage its bytes through the content lane, and insert the rows. + // A new order is created for each photo so the foreign key is always satisfied. + let pick_photo = if identity.is_empty() { + Callback::noop() + } else { + let client = client.clone(); + let content = Rc::clone(&props.content.0); + let identity = identity.clone(); + Callback::from(move |event: Event| { + let input: HtmlInputElement = event.target_unchecked_into(); + if let Some(files) = input.files() + && let Some(file) = files.get(0) + { + let client = client.clone(); + let content = Rc::clone(&content); + let identity = identity.clone(); + spawn_local(async move { + let blob: web_sys::Blob = file.into(); + let stage_identity = identity.clone(); + let result = content + .stage(&blob, MimeClass::Jpeg, &client, move |conn, file_id| { + conn.transaction(|conn| { + let before_orders: std::collections::HashSet = + orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(stage_identity.as_str()), + orders::quantity.eq(1_i64), + )) + .execute(conn)?; + let order_id = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .find(|id| !before_orders.contains(id)) + .expect("order minted"); + let before_photos: std::collections::HashSet = + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(photos::table) + .values(( + photos::order_id.eq(order_id), + photos::owner_id.eq(stage_identity.as_str()), + photos::content_id.eq(file_id.as_bytes().to_vec()), + photos::content_state.eq::>(None), + )) + .execute(conn)?; + Ok::<_, diesel::result::Error>( + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .find(|id| !before_photos.contains(id)) + .expect("photo minted"), + ) + }) + }) + .await; + match result { + Ok(_) => { + client.replay_pending().await.ok(); + } + Err(err) => { + tracing::error!(error = %err, "photo stage failed"); + } + } + }); + } + }) + }; + let tidy = { let client = client.clone(); let footprint = footprint.clone(); @@ -1370,6 +1527,10 @@ fn dashboard(props: &DashboardProps) -> Html { let notes_error = notes.error().map(|err| { html! {

{ format!("notes error: {err}") }

} }); + let photos_error = photos.error().map(|err| { + html! {

{ format!("photos error: {err}") }

} + }); + let urls = (*photo_urls).clone(); html! {
@@ -1377,7 +1538,7 @@ fn dashboard(props: &DashboardProps) -> Html {

{ "orders " }{ "synced" }

{ format!("count {order_count}, total quantity {order_sum}. Converges across every window through Postgres.") }

- + if newest_order.is_some() { } @@ -1398,6 +1559,33 @@ fn dashboard(props: &DashboardProps) -> Html {
+
+

{ "photos " }{ "synced" }

+

{ format!("count {photo_count}. Pick an image file to stage and upload through the content route.") }

+
+ +
+ { photos_error.unwrap_or_default() } + + + + { for photo_rows.iter().map(|row| { + let id = row.id.to_string(); + let state = row.content_state.as_deref().unwrap_or("pending"); + let img = urls.get(&row.id).map(|url| { + html! { } + }); + html! { + + + + + + } + }) } + +
{ "id" }{ "state" }{ "image" }
{ id }{ state }{ img.unwrap_or_default() }
+

{ "notes " }{ "device-only" }

{ format!("count {note_count}. Converges across this device's windows through the DB worker, never the server.") }

@@ -1423,7 +1611,7 @@ fn dashboard(props: &DashboardProps) -> Html {

{ "retention " }{ "R15" }

-

{ format!("Replica mirror: {pages} pages (~{kb} KB), {free} free to reclaim. Covered rows: {}.", order_count + note_count) }

+

{ format!("Replica mirror: {pages} pages (~{kb} KB), {free} free to reclaim. Covered rows: {}.", order_count + note_count + photo_count) }

{ "Ending a subscription evicts the rows no live subscription still covers, and the trimming pass hands the freed pages back to storage." }

diff --git a/examples/yew-web-demo/tests/photo_flow.rs b/examples/yew-web-demo/tests/photo_flow.rs new file mode 100644 index 00000000..4f8dea38 --- /dev/null +++ b/examples/yew-web-demo/tests/photo_flow.rs @@ -0,0 +1,524 @@ +//! Photo surface browser tests for the yew web demo. +//! +//! Drives the real demo boot path (`db_worker_photo_boot`) through the browser +//! stack. Two tests run in this binary and are serialized by a shared Web Lock +//! to avoid OPFS conflicts between concurrent workers. +//! +//! COUPLING: `SchemaVersion` is a hash of the schema SQL source string. +//! `schema.sql` in this workspace MUST be byte-identical to +//! `examples/wasm-smoke/schema.sql`, which is what the browser stack server +//! is compiled against. The first test (`a_tab_order_reaches_the_replica…`) +//! acts as the permanent drift detector: it fails at the WebSocket handshake +//! the moment the two schemas diverge, long before any photo assertion runs. +//! Maintain the identity with `cp examples/wasm-smoke/schema.sql +//! examples/yew-web-demo/schema.sql` whenever wasm-smoke's schema changes. + +#![cfg(target_arch = "wasm32")] + +use connetto_client::dsl::Watchable; +use connetto_client::{ + ClientConfig, ClientEvent, ConnettoClient, ConnettoConnection, Grant, LiveQuery, Replica, +}; +use connetto_file_core::{FileId, MimeClass}; +use connetto_web::auth::{ + Acquired, BrowserAuthenticator, IdbKeyStore, LOGIN_CHANNEL, LoginMessage, RefreshStore, + WorkerAuthConfig, deliver_login_code, +}; +use connetto_web::storage::{ReplicaStorage, device_key}; +use connetto_web::{MessageTransport, TabContent, TabResolved, locks, workers}; +use connetto_yew_web_demo::{ + CALLER_FUNCTION, DEMO_TAB_DDL, demo_policy_tables, demo_schema_version, uuidv4_functions, +}; +use diesel::prelude::*; +use futures_channel::oneshot; +use js_sys::{Array, Uint8Array}; +use wasm_bindgen::JsCast; +use wasm_bindgen::prelude::*; +use wasm_bindgen_futures::{JsFuture, spawn_local}; +use wasm_bindgen_test::{wasm_bindgen_test, wasm_bindgen_test_configure}; +use web_sys::{ + BroadcastChannel, DedicatedWorkerGlobalScope, MessageEvent, Request, RequestInit, Worker, +}; + +wasm_bindgen_test_configure!(run_in_dedicated_worker); + +// Auth server coordinates matching the browser stack setup. +const AUTH_BASE: &str = "http://127.0.0.1:18099"; +const AUTH_LANDING: &str = "http://127.0.0.1:18099/dev/landing"; +const AUTH_PROVIDER: &str = "dev-idp"; +const AUTH_USERNAME: &str = "startup"; + +// Serializes tests in this binary to prevent OPFS worker conflicts. +const SUITE_LOCK: &str = "connetto-yew-photo-suite"; + +// --- Diesel table schema matching DEMO_TAB_DDL (policy-split, uses logical names) --- + +diesel::table! { + orders (id) { + id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, + quantity -> diesel::sql_types::BigInt, + } +} + +diesel::table! { + photos (id) { + id -> rosetta_uuid::sql_types::Uuid, + order_id -> rosetta_uuid::sql_types::Uuid, + owner_id -> diesel::sql_types::Text, + content_id -> diesel::sql_types::Binary, + content_state -> diesel::sql_types::Nullable, + } +} + +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = orders)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Order { + id: rosetta_uuid::Uuid, + owner_id: String, + quantity: i64, +} + +#[derive(Queryable, Selectable, Debug, PartialEq, Clone)] +#[diesel(table_name = photos)] +#[diesel(check_for_backend(diesel::sqlite::Sqlite))] +struct Photo { + id: rosetta_uuid::Uuid, + order_id: rosetta_uuid::Uuid, + owner_id: String, + content_id: Vec, + content_state: Option, +} + +// --- Helpers --- + +fn stage(msg: &str) { + web_sys::console::log_1(&msg.into()); +} + +fn relay_worker_breadcrumbs() { + use std::sync::atomic::{AtomicBool, Ordering}; + static INSTALLED: AtomicBool = AtomicBool::new(false); + if INSTALLED.swap(true, Ordering::Relaxed) { + return; + } + if let Ok(ch) = BroadcastChannel::new("connetto-debug") { + let cb = + wasm_bindgen::closure::Closure::::new(|e: MessageEvent| { + web_sys::console::log_1(&e.data()); + }); + ch.set_onmessage(Some(cb.as_ref().unchecked_ref())); + cb.forget(); + std::mem::forget(ch); + } +} + +fn play_the_tab() { + use std::sync::atomic::{AtomicBool, Ordering}; + static INSTALLED: AtomicBool = AtomicBool::new(false); + if INSTALLED.swap(true, Ordering::Relaxed) { + return; + } + let ch = BroadcastChannel::new(LOGIN_CHANNEL).expect("login channel"); + let cb = wasm_bindgen::closure::Closure::::new(|e: MessageEvent| { + let Some(text) = e.data().as_string() else { + return; + }; + let Ok(LoginMessage::Request { url }) = serde_json::from_str::(&text) else { + return; + }; + spawn_local(async move { + let (code, state) = walk_login(&url).await; + deliver_login_code(&code, &state).expect("deliver login code"); + }); + }); + ch.set_onmessage(Some(cb.as_ref().unchecked_ref())); + cb.forget(); + std::mem::forget(ch); +} + +async fn global_fetch_str(url: &str) -> web_sys::Response { + let promise = js_sys::global() + .dyn_into::() + .map(|w| w.fetch_with_str(url)) + .unwrap_or_else(|_| web_sys::window().expect("window").fetch_with_str(url)); + JsFuture::from(promise) + .await + .expect("fetch") + .dyn_into() + .expect("Response") +} + +async fn global_fetch_req(req: &Request) -> web_sys::Response { + let promise = js_sys::global() + .dyn_into::() + .map(|w| w.fetch_with_request(req)) + .unwrap_or_else(|_| web_sys::window().expect("window").fetch_with_request(req)); + JsFuture::from(promise) + .await + .expect("fetch req") + .dyn_into() + .expect("Response") +} + +async fn walk_login(login_url: &str) -> (String, String) { + let resp = global_fetch_str(login_url).await; + let form_url = resp.url(); + let init = RequestInit::new(); + init.set_method("POST"); + init.set_body(&JsValue::from_str(&format!("username={AUTH_USERNAME}"))); + let req = Request::new_with_str_and_init(&form_url, &init).expect("login request"); + req.headers() + .set("content-type", "application/x-www-form-urlencoded") + .expect("content-type header"); + let resp = global_fetch_req(&req).await; + assert!( + resp.ok(), + "login chain ended at {} with status {}", + resp.url(), + resp.status() + ); + let final_url = resp.url(); + let parsed = web_sys::Url::new(&final_url).expect("parse final url"); + let params = parsed.search_params(); + ( + params + .get("code") + .unwrap_or_else(|| panic!("no code in {final_url}")), + params + .get("state") + .unwrap_or_else(|| panic!("no state in {final_url}")), + ) +} + +async fn mint_session() -> (String, String) { + let storage = ReplicaStorage::install().await; + let keys = IdbKeyStore::open().await.expect("open key store"); + let device = device_key(&keys).await.expect("device key"); + static N: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + let n = N.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let db = format!("yew-photo-mint-{n}.sqlite"); + let store = RefreshStore::open(&storage.db_url(&db), &device).expect("refresh store"); + let auth = BrowserAuthenticator::new( + WorkerAuthConfig::new(AUTH_BASE, AUTH_PROVIDER, AUTH_LANDING), + None, + ); + let pending = match auth + .acquire::(&store) + .await + .expect("acquire session") + { + Acquired::NeedLogin(p) => p, + Acquired::Access(_) => panic!("fresh store cannot refresh silently"), + }; + let (code, state) = walk_login(&pending.login_url).await; + let session = auth + .complete::(&pending, &code, &state, &store) + .await + .expect("complete login"); + drop(store); + storage.delete_db(&db).ok(); + (session.access_token, session.user_id) +} + +fn glue_url() -> String { + let found = js_sys::eval( + r#"performance.getEntriesByType("resource").map(e=>e.name).find(n=>n.endsWith("_bg.wasm"))"#, + ) + .expect("resource entries") + .as_string() + .expect("wasm resource entry"); + let base = found.strip_suffix("_bg.wasm").expect("wasm suffix"); + format!("{base}.js") +} + +fn spawn_photo_worker(glue_url: &str) -> Worker { + let wasm_url = glue_url + .strip_suffix(".js") + .map_or_else(|| format!("{glue_url}_bg.wasm"), |b| format!("{b}_bg.wasm")); + let src = format!( + "const ch=new BroadcastChannel('connetto-debug');\n\ + try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n ch.postMessage('yew photo worker ready');\n await mod.db_worker_photo_boot();\n}}catch(e){{\n ch.postMessage('yew photo worker FAILED: '+e);\n throw e;\n}}", + g = glue_url, + w = wasm_url, + ); + let parts = Array::of1(&JsValue::from_str(&src)); + let opts = web_sys::BlobPropertyBag::new(); + opts.set_type("text/javascript"); + let blob = + web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("bootstrap blob"); + let url = web_sys::Url::create_object_url_with_blob(&blob).expect("object url"); + let worker_opts = web_sys::WorkerOptions::new(); + worker_opts.set_type(web_sys::WorkerType::Module); + worker_opts.set_name("connetto-yew-photo"); + let w = Worker::new_with_options(&url, &worker_opts).expect("spawn worker"); + web_sys::Url::revoke_object_url(&url).ok(); + w +} + +fn photo_bytes() -> Vec { + const SEED: [u8; 16] = [ + 0x10, 0x32, 0x54, 0x76, 0x98, 0xba, 0xdc, 0xfe, 0x01, 0x23, 0x45, 0x67, 0x89, 0xab, 0xcd, + 0xef, + ]; + let mut v = Vec::with_capacity(4096); + while v.len() < 4096 { + v.extend_from_slice(&SEED); + } + v.truncate(4096); + v +} + +fn blob_of(bytes: &[u8]) -> web_sys::Blob { + let arr = Uint8Array::from(bytes); + web_sys::Blob::new_with_u8_array_sequence(&Array::of1(&arr)).expect("blob from bytes") +} + +async fn fetch_bytes(url: &str) -> Vec { + let scope = js_sys::global() + .dyn_into::() + .expect("dedicated worker"); + let resp: web_sys::Response = JsFuture::from(scope.fetch_with_str(url)) + .await + .expect("fetch content") + .dyn_into() + .expect("Response"); + assert!( + resp.ok(), + "content fetch at {} returned {}", + resp.url(), + resp.status() + ); + let buf = JsFuture::from(resp.array_buffer().expect("array_buffer promise")) + .await + .expect("array buffer"); + Uint8Array::new(&buf).to_vec() +} + +async fn connect_tab( + client_id: &str, + token: String, + identity: &str, +) -> ( + TabContent, + ConnettoConnection>, +) { + let wire = format!("connetto-wire-{client_id}"); + workers::announce_tab(&wire).await.expect("announce tab"); + let mut transport = MessageTransport::::new(&wire).expect("transport"); + let content = TabContent::new(&mut transport); + let config = ClientConfig::new(client_id.to_owned()) + .with_login(Some(Grant::new(token))) + .with_schema_version(Some(demo_schema_version())) + .with_sql_functions(uuidv4_functions()) + .with_policy_tables(demo_policy_tables()) + .with_caller(CALLER_FUNCTION, identity); + let conn = ConnettoConnection::connect( + transport, + &Replica::in_memory(), + DEMO_TAB_DDL, + &config, + None, + ) + .await + .expect("tab connect"); + (content, conn) +} + +// --- Test 1: alignment proof — version agreement via an orders round trip --- + +#[wasm_bindgen_test] +async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { + relay_worker_breadcrumbs(); + play_the_tab(); + let _serial = locks::hold_lock(SUITE_LOCK).await; + let worker = spawn_photo_worker(&glue_url()); + workers::await_db_worker_ready(&[]) + .await + .expect("db worker ready"); + stage("yew demo worker booted (align)"); + + let (token, identity) = mint_session().await; + let client_id = rosetta_uuid::Uuid::new_v4().to_string(); + let _tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; + let (_content, mut conn) = connect_tab(&client_id, token, &identity).await; + conn.subscribe("align-orders", "SELECT * FROM orders") + .await + .expect("orders subscribe"); + loop { + let event = conn.pump_one().await.expect("pump"); + assert_ne!(event, ClientEvent::Closed, "connection closed early"); + if matches!(event, ClientEvent::SnapshotEnd { .. }) { + break; + } + } + stage("orders subscription ready (align)"); + + let (client, pump) = ConnettoClient::with_pump(conn); + let (done_tx, done_rx) = oneshot::channel::<()>(); + spawn_local(async move { + pump.await; + let _ = done_tx.send(()); + }); + + let mut live: LiveQuery = orders::table + .order(orders::id) + .select(Order::as_select()) + .live(&client) + .await + .expect("orders live query"); + + let qty = 77_i64; + let id_str = identity.clone(); + client + .with_conn(move |conn| { + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(id_str.as_str()), + orders::quantity.eq(qty), + )) + .execute(conn.conn()) + }) + .await + .expect("order insert"); + client.replay_pending().await.expect("replay pending"); + stage("order written (align)"); + + loop { + if live.rows().iter().any(|o| o.quantity == qty) { + break; + } + live.changed().await.expect("live refresh"); + } + stage("order arrived at replica — version agreement confirmed"); + + drop(live); + drop(client); + done_rx.await.expect("pump exited"); + worker.terminate(); +} + +// --- Test 2: photo round trip through stage, commit, resolve, and HTTP fetch --- + +#[wasm_bindgen_test] +async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { + relay_worker_breadcrumbs(); + play_the_tab(); + let _serial = locks::hold_lock(SUITE_LOCK).await; + let worker = spawn_photo_worker(&glue_url()); + workers::await_db_worker_ready(&[]) + .await + .expect("db worker ready"); + stage("yew demo worker booted (photo)"); + + let (token, identity) = mint_session().await; + let client_id = rosetta_uuid::Uuid::new_v4().to_string(); + let _tab_lock = locks::hold_lock(&locks::tab_lock_name(&client_id)).await; + let (content, mut conn) = connect_tab(&client_id, token, &identity).await; + conn.subscribe("photo-test-photos", "SELECT * FROM photos") + .await + .expect("photos subscribe"); + loop { + let event = conn.pump_one().await.expect("pump"); + assert_ne!(event, ClientEvent::Closed, "connection closed early"); + if matches!(event, ClientEvent::SnapshotEnd { .. }) { + break; + } + } + stage("photos subscription ready"); + + let (client, pump) = ConnettoClient::with_pump(conn); + let (done_tx, done_rx) = oneshot::channel::<()>(); + spawn_local(async move { + pump.await; + let _ = done_tx.send(()); + }); + + let mut live: LiveQuery = photos::table + .order(photos::id) + .select(Photo::as_select()) + .live(&client) + .await + .expect("photo live query"); + + let bytes = photo_bytes(); + let blob = blob_of(&bytes); + let (file_id, photo_id) = content + .stage(&blob, MimeClass::Jpeg, &client, |conn, file_id| { + conn.transaction(|conn| { + let before_orders: std::collections::HashSet = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(orders::table) + .values(( + orders::owner_id.eq(identity.as_str()), + orders::quantity.eq(1_i64), + )) + .execute(conn)?; + let order_id = orders::table + .select(orders::id) + .load::(conn)? + .into_iter() + .find(|id| !before_orders.contains(id)) + .expect("minted order id"); + let before_photos: std::collections::HashSet = photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .collect(); + diesel::insert_into(photos::table) + .values(( + photos::order_id.eq(order_id), + photos::owner_id.eq(identity.as_str()), + photos::content_id.eq(file_id.as_bytes().to_vec()), + photos::content_state.eq::>(None), + )) + .execute(conn)?; + Ok::( + photos::table + .select(photos::id) + .load::(conn)? + .into_iter() + .find(|id| !before_photos.contains(id)) + .expect("minted photo id"), + ) + }) + }) + .await + .expect("stage photo and insert rows"); + client.replay_pending().await.expect("send mutation"); + assert_eq!(file_id, FileId::from_chunks([bytes.as_slice()])); + stage("photo staged"); + + let available = loop { + if let Some(p) = live + .rows() + .iter() + .find(|p| p.id == photo_id && p.content_state.as_deref() == Some("available")) + .cloned() + { + break p; + } + live.changed().await.expect("live refresh"); + }; + assert_eq!(available.owner_id, identity); + assert_eq!(available.content_id, file_id.as_bytes().to_vec()); + stage("photo available"); + + let url = match content.resolve(file_id).await { + TabResolved::Remote { url } => url, + TabResolved::Local { .. } => panic!("uploaded photo must resolve to server"), + TabResolved::Unavailable => panic!("available photo must resolve"), + }; + assert_eq!(fetch_bytes(&url).await, bytes); + stage("photo bytes verified via HTTP fetch"); + + drop(live); + drop(content); + drop(client); + done_rx.await.expect("pump exited"); + worker.terminate(); +} From 41e0089a30e8812bad07c6b84ceea1b6724b9891 Mon Sep 17 00:00:00 2001 From: LucaCappelletti94 Date: Fri, 18 Sep 2026 12:48:53 +0200 Subject: [PATCH 3/3] Give each demo suite test its own worker databases so OPFS handles cannot strand the next boot --- examples/dioxus-web-demo/src/lib.rs | 62 +++++++------ examples/dioxus-web-demo/tests/photo_flow.rs | 94 ++++++++++--------- examples/yew-web-demo/src/lib.rs | 62 +++++++------ examples/yew-web-demo/tests/photo_flow.rs | 95 ++++++++++---------- 4 files changed, 161 insertions(+), 152 deletions(-) diff --git a/examples/dioxus-web-demo/src/lib.rs b/examples/dioxus-web-demo/src/lib.rs index 4dd2507e..07badd79 100644 --- a/examples/dioxus-web-demo/src/lib.rs +++ b/examples/dioxus-web-demo/src/lib.rs @@ -1,44 +1,36 @@ -//! Browser photo worker infrastructure for the dioxus web demo. -//! -//! This lib target exposes the `db_worker_photo_boot` wasm-bindgen entry point -//! for browser tests that drive the photo flow through the real content routes. -//! The main binary (`src/main.rs`) is the actual Dioxus application. +//! Browser test worker entry points for the dioxus web demo photo surface. use wasm_bindgen::JsValue; use wasm_bindgen::prelude::wasm_bindgen; include!(concat!(env!("OUT_DIR"), "/replica-tables.rs")); -/// Schema SQL this build was compiled against (matches the browser-stack server). +// SchemaVersion hashes this string; it must be byte-identical to +// examples/wasm-smoke/schema.sql, which the browser-stack server uses. pub const SCHEMA_SQL: &str = include_str!("../schema.sql"); -/// The synced replica schema (worker replica, policy-split by build.rs). pub const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); - -/// The local tier schema (device-private, attached, never synced). pub const DEMO_FRONTEND_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")); - -/// The tab mirror schema: both tiers in the tab's main schema. pub const DEMO_TAB_DDL: &str = concat!( include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")), "\n", include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")), ); -/// The demo server the DB worker connects upstream to. pub const DEMO_WS_URL: &str = "ws://127.0.0.1:7777/"; - -/// The upstream subscription the DB worker registers. pub const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; - -/// The extra upstream subscription for photos. pub const PHOTO_QUERY: &str = "SELECT * FROM photos"; +pub const CALLER_FUNCTION: &str = "current_app_user"; -/// The OPFS file base for the worker's durable synced replica. -pub const DB_NAME: &str = "connetto-photo-dioxus.sqlite"; +// Each test gets unique OPFS filenames so Chrome's delayed handle release after +// Worker.terminate() never blocks the next worker's file open. +const ALIGN_DB_PREFIX: &str = "connetto-dioxus-align"; +const ALIGN_HUB_META: &str = "connetto-dioxus-align-hub-meta.sqlite"; +const ALIGN_AUTH_DB: &str = "connetto-dioxus-align-auth.sqlite"; -/// The registered caller-identity function connetto installs on every connection. -pub const CALLER_FUNCTION: &str = "current_app_user"; +const PHOTO_DB_PREFIX: &str = "connetto-dioxus-photo"; +const PHOTO_HUB_META: &str = "connetto-dioxus-photo-hub-meta.sqlite"; +const PHOTO_AUTH_DB: &str = "connetto-dioxus-photo-auth.sqlite"; /// The schema version this build was compiled against. #[must_use] @@ -46,8 +38,6 @@ pub fn demo_schema_version() -> connetto_core::SchemaVersion { connetto_core::SchemaVersion::from_source(SCHEMA_SQL) } -// The uuidv4 SQL function registered on every connection so the orders -// and photos DEFAULT (uuidv4()) mints a UUID on local writes. #[diesel::declare_sql_function] extern "SQL" { /// Client-authored primary key: a 16-byte UUID v4, stored as a BLOB. @@ -68,7 +58,7 @@ pub fn uuidv4_functions() -> connetto_client::SqlFunctions { )) } -/// The policy table map for this build, for `ClientConfig::with_policy_tables`. +/// The policy table map for this build. #[must_use] pub fn demo_policy_tables() -> connetto_client::PolicyTables { connetto_client::PolicyTables::from_translation( @@ -77,26 +67,42 @@ pub fn demo_policy_tables() -> connetto_client::PolicyTables { ) } -/// DB worker entry point: boot the connetto DB tier with the photo config. +/// Worker entry point for the alignment test; uses `connetto-dioxus-align*` OPFS files. +/// +/// # Errors /// -/// The test's blob worker bootstrap imports this crate's wasm module and awaits this. +/// A string describing the VFS, upstream connect, or subscribe failure. +#[wasm_bindgen] +pub async fn db_worker_boot_align() -> Result<(), JsValue> { + boot_with(ALIGN_DB_PREFIX, ALIGN_HUB_META, ALIGN_AUTH_DB).await +} + +/// Worker entry point for the photo test; uses `connetto-dioxus-photo*` OPFS files. /// /// # Errors /// /// A string describing the VFS, upstream connect, or subscribe failure. #[wasm_bindgen] pub async fn db_worker_photo_boot() -> Result<(), JsValue> { + boot_with(PHOTO_DB_PREFIX, PHOTO_HUB_META, PHOTO_AUTH_DB).await +} + +async fn boot_with( + db_prefix: &'static str, + hub_meta: &'static str, + auth_db: &'static str, +) -> Result<(), JsValue> { connetto_web::logging::init_console(); connetto_web::workers::boot_db_worker::( &connetto_web::workers::DbWorkerConfig::new(demo_schema_version()) .with_ws_url(DEMO_WS_URL) - .with_replica_db_prefix(DB_NAME) + .with_replica_db_prefix(db_prefix) .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-photo-dioxus-hub-meta.sqlite") + .with_hub_meta_name(hub_meta) .with_content_namespace("connetto-photo-content") .with_sql_functions(uuidv4_functions()) .with_policy_tables(demo_policy_tables()) @@ -106,7 +112,7 @@ pub async fn db_worker_photo_boot() -> Result<(), JsValue> { "dev-idp", "http://127.0.0.1:18099/dev/landing", ))) - .with_auth_db_name("connetto-photo-dioxus-auth.sqlite"), + .with_auth_db_name(auth_db), ) .await .map(drop) diff --git a/examples/dioxus-web-demo/tests/photo_flow.rs b/examples/dioxus-web-demo/tests/photo_flow.rs index 90097ef5..da309f81 100644 --- a/examples/dioxus-web-demo/tests/photo_flow.rs +++ b/examples/dioxus-web-demo/tests/photo_flow.rs @@ -1,17 +1,19 @@ //! Photo surface browser tests for the dioxus web demo. //! -//! Drives the real demo boot path (`db_worker_photo_boot`) through the browser -//! stack. Two tests run in this binary and are serialized by a shared Web Lock -//! to avoid OPFS conflicts between concurrent workers. +//! COUPLING: `SchemaVersion` is a hash of `schema.sql`. This file must be +//! byte-identical to `examples/wasm-smoke/schema.sql`, which the browser-stack +//! server uses. The first test is the permanent drift detector: it fails at +//! the WebSocket handshake the instant the two schemas diverge. Maintain +//! identity with `cp examples/wasm-smoke/schema.sql +//! examples/dioxus-web-demo/schema.sql`. //! -//! COUPLING: `SchemaVersion` is a hash of the schema SQL source string. -//! `schema.sql` in this workspace MUST be byte-identical to -//! `examples/wasm-smoke/schema.sql`, which is what the browser stack server -//! is compiled against. The first test (`a_tab_order_reaches_the_replica…`) -//! acts as the permanent drift detector: it fails at the WebSocket handshake -//! the moment the two schemas diverge, long before any photo assertion runs. -//! Maintain the identity with `cp examples/wasm-smoke/schema.sql -//! examples/dioxus-web-demo/schema.sql` whenever wasm-smoke's schema changes. +//! ISOLATION: each test uses a distinct `db_worker_boot_*` function whose OPFS +//! filenames are unique, so Chrome's delayed `FileSystemSyncAccessHandle` +//! release after `Worker.terminate()` never blocks the next worker's file open. +//! `DB_ALIVE_LOCK` is shared across all connetto workers on the same origin; +//! each test holds its own `SUITE_LOCK_*` that serialises that test's lifecycle +//! and a 200 ms sleep before releasing the lock gives Chrome time to finish the +//! OPFS cleanup before the next worker starts. #![cfg(target_arch = "wasm32")] @@ -42,16 +44,13 @@ use web_sys::{ wasm_bindgen_test_configure!(run_in_dedicated_worker); -// Auth server coordinates matching the browser stack setup. const AUTH_BASE: &str = "http://127.0.0.1:18099"; const AUTH_LANDING: &str = "http://127.0.0.1:18099/dev/landing"; const AUTH_PROVIDER: &str = "dev-idp"; const AUTH_USERNAME: &str = "startup"; -// Serializes tests in this binary to prevent OPFS worker conflicts. -const SUITE_LOCK: &str = "connetto-dioxus-photo-suite"; - -// --- Diesel table schema matching DEMO_TAB_DDL (policy-split, uses logical names) --- +const SUITE_LOCK_ALIGN: &str = "connetto-dioxus-photo-suite-align"; +const SUITE_LOCK_PHOTO: &str = "connetto-dioxus-photo-suite-photo"; diesel::table! { orders (id) { @@ -91,8 +90,6 @@ struct Photo { content_state: Option, } -// --- Helpers --- - fn stage(msg: &str) { web_sys::console::log_1(&msg.into()); } @@ -171,7 +168,7 @@ async fn walk_login(login_url: &str) -> (String, String) { let req = Request::new_with_str_and_init(&form_url, &init).expect("login request"); req.headers() .set("content-type", "application/x-www-form-urlencoded") - .expect("content-type header"); + .expect("header"); let resp = global_fetch_req(&req).await; assert!( resp.ok(), @@ -180,7 +177,7 @@ async fn walk_login(login_url: &str) -> (String, String) { resp.status() ); let final_url = resp.url(); - let parsed = web_sys::Url::new(&final_url).expect("parse final url"); + let parsed = web_sys::Url::new(&final_url).expect("parse url"); let params = parsed.search_params(); ( params @@ -204,11 +201,7 @@ async fn mint_session() -> (String, String) { WorkerAuthConfig::new(AUTH_BASE, AUTH_PROVIDER, AUTH_LANDING), None, ); - let pending = match auth - .acquire::(&store) - .await - .expect("acquire session") - { + let pending = match auth.acquire::(&store).await.expect("acquire") { Acquired::NeedLogin(p) => p, Acquired::Access(_) => panic!("fresh store cannot refresh silently"), }; @@ -233,25 +226,24 @@ fn glue_url() -> String { format!("{base}.js") } -fn spawn_photo_worker(glue_url: &str) -> Worker { +fn spawn_worker(glue_url: &str, boot_fn: &str) -> Worker { let wasm_url = glue_url .strip_suffix(".js") .map_or_else(|| format!("{glue_url}_bg.wasm"), |b| format!("{b}_bg.wasm")); let src = format!( "const ch=new BroadcastChannel('connetto-debug');\n\ - try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n ch.postMessage('dioxus photo worker ready');\n await mod.db_worker_photo_boot();\n}}catch(e){{\n ch.postMessage('dioxus photo worker FAILED: '+e);\n throw e;\n}}", + try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n await mod.{f}();\n}}catch(e){{\n ch.postMessage('dioxus worker FAILED: '+e);\n throw e;\n}}", g = glue_url, w = wasm_url, + f = boot_fn, ); let parts = Array::of1(&JsValue::from_str(&src)); let opts = web_sys::BlobPropertyBag::new(); opts.set_type("text/javascript"); - let blob = - web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("bootstrap blob"); + let blob = web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("blob"); let url = web_sys::Url::create_object_url_with_blob(&blob).expect("object url"); let worker_opts = web_sys::WorkerOptions::new(); worker_opts.set_type(web_sys::WorkerType::Module); - worker_opts.set_name("connetto-dioxus-photo"); let w = Worker::new_with_options(&url, &worker_opts).expect("spawn worker"); web_sys::Url::revoke_object_url(&url).ok(); w @@ -278,10 +270,10 @@ fn blob_of(bytes: &[u8]) -> web_sys::Blob { async fn fetch_bytes(url: &str) -> Vec { let scope = js_sys::global() .dyn_into::() - .expect("dedicated worker"); + .expect("worker"); let resp: web_sys::Response = JsFuture::from(scope.fetch_with_str(url)) .await - .expect("fetch content") + .expect("fetch") .dyn_into() .expect("Response"); assert!( @@ -326,18 +318,18 @@ async fn connect_tab( (content, conn) } -// --- Test 1: alignment proof — version agreement via an orders round trip --- +// --- Test 1: alignment proof --- #[wasm_bindgen_test] async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { relay_worker_breadcrumbs(); play_the_tab(); - let _serial = locks::hold_lock(SUITE_LOCK).await; - let worker = spawn_photo_worker(&glue_url()); + let _serial = locks::hold_lock(SUITE_LOCK_ALIGN).await; + let worker = spawn_worker(&glue_url(), "db_worker_boot_align"); workers::await_db_worker_ready(&[]) .await .expect("db worker ready"); - stage("dioxus demo worker booted (align)"); + stage("dioxus align worker booted"); let (token, identity) = mint_session().await; let client_id = rosetta_uuid::Uuid::new_v4().to_string(); @@ -345,7 +337,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { let (_content, mut conn) = connect_tab(&client_id, token, &identity).await; conn.subscribe("align-orders", "SELECT * FROM orders") .await - .expect("orders subscribe"); + .expect("subscribe"); loop { let event = conn.pump_one().await.expect("pump"); assert_ne!(event, ClientEvent::Closed, "connection closed early"); @@ -353,7 +345,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { break; } } - stage("orders subscription ready (align)"); + stage("orders subscription ready"); let (client, pump) = ConnettoClient::with_pump(conn); let (done_tx, done_rx) = oneshot::channel::<()>(); @@ -367,7 +359,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { .select(Order::as_select()) .live(&client) .await - .expect("orders live query"); + .expect("live query"); let qty = 77_i64; let id_str = identity.clone(); @@ -383,7 +375,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { .await .expect("order insert"); client.replay_pending().await.expect("replay pending"); - stage("order written (align)"); + stage("order written"); loop { if live.rows().iter().any(|o| o.quantity == qty) { @@ -391,26 +383,29 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { } live.changed().await.expect("live refresh"); } - stage("order arrived at replica — version agreement confirmed"); + stage("order arrived — version agreement confirmed"); drop(live); drop(client); done_rx.await.expect("pump exited"); worker.terminate(); + // Chrome headless does not release the OPFS access handle synchronously on + // terminate; per-test file names are the isolation guarantee, this is margin. + workers::sleep(core::time::Duration::from_millis(200)).await; } -// --- Test 2: photo round trip through stage, commit, resolve, and HTTP fetch --- +// --- Test 2: photo round trip --- #[wasm_bindgen_test] async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { relay_worker_breadcrumbs(); play_the_tab(); - let _serial = locks::hold_lock(SUITE_LOCK).await; - let worker = spawn_photo_worker(&glue_url()); + let _serial = locks::hold_lock(SUITE_LOCK_PHOTO).await; + let worker = spawn_worker(&glue_url(), "db_worker_photo_boot"); workers::await_db_worker_ready(&[]) .await .expect("db worker ready"); - stage("dioxus demo worker booted (photo)"); + stage("dioxus photo worker booted"); let (token, identity) = mint_session().await; let client_id = rosetta_uuid::Uuid::new_v4().to_string(); @@ -418,7 +413,7 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { let (content, mut conn) = connect_tab(&client_id, token, &identity).await; conn.subscribe("photo-test-photos", "SELECT * FROM photos") .await - .expect("photos subscribe"); + .expect("subscribe"); loop { let event = conn.pump_one().await.expect("pump"); assert_ne!(event, ClientEvent::Closed, "connection closed early"); @@ -440,7 +435,7 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { .select(Photo::as_select()) .live(&client) .await - .expect("photo live query"); + .expect("live query"); let bytes = photo_bytes(); let blob = blob_of(&bytes); @@ -514,11 +509,14 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { TabResolved::Unavailable => panic!("available photo must resolve"), }; assert_eq!(fetch_bytes(&url).await, bytes); - stage("photo bytes verified via HTTP fetch"); + stage("photo bytes verified"); drop(live); drop(content); drop(client); done_rx.await.expect("pump exited"); worker.terminate(); + // Chrome headless does not release the OPFS access handle synchronously on + // terminate; per-test file names are the isolation guarantee, this is margin. + workers::sleep(core::time::Duration::from_millis(200)).await; } diff --git a/examples/yew-web-demo/src/lib.rs b/examples/yew-web-demo/src/lib.rs index 6434255b..29ef9d9e 100644 --- a/examples/yew-web-demo/src/lib.rs +++ b/examples/yew-web-demo/src/lib.rs @@ -1,44 +1,36 @@ -//! Browser photo worker infrastructure for the yew web demo. -//! -//! This lib target exposes the `db_worker_photo_boot` wasm-bindgen entry point -//! for browser tests that drive the photo flow through the real content routes. -//! The main binary (`src/main.rs`) is the actual Yew application. +//! Browser test worker entry points for the yew web demo photo surface. use wasm_bindgen::JsValue; use wasm_bindgen::prelude::wasm_bindgen; include!(concat!(env!("OUT_DIR"), "/replica-tables.rs")); -/// Schema SQL this build was compiled against (matches the browser-stack server). +// SchemaVersion hashes this string; it must be byte-identical to +// examples/wasm-smoke/schema.sql, which the browser-stack server uses. pub const SCHEMA_SQL: &str = include_str!("../schema.sql"); -/// The synced replica schema (worker replica, policy-split by build.rs). pub const DEMO_SQLITE_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")); - -/// The local tier schema (device-private, attached, never synced). pub const DEMO_FRONTEND_DDL: &str = include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")); - -/// The tab mirror schema: both tiers in the tab's main schema. pub const DEMO_TAB_DDL: &str = concat!( include_str!(concat!(env!("OUT_DIR"), "/replica-ddl.sql")), "\n", include_str!(concat!(env!("OUT_DIR"), "/frontend-ddl.sql")), ); -/// The demo server the DB worker connects upstream to. pub const DEMO_WS_URL: &str = "ws://127.0.0.1:7777/"; - -/// The upstream subscription the DB worker registers. pub const DEMO_QUERY: &str = "SELECT * FROM orders WHERE quantity > 0"; - -/// The extra upstream subscription for photos. pub const PHOTO_QUERY: &str = "SELECT * FROM photos"; +pub const CALLER_FUNCTION: &str = "current_app_user"; -/// The OPFS file base for the worker's durable synced replica. -pub const DB_NAME: &str = "connetto-photo-yew.sqlite"; +// Each test gets unique OPFS filenames so Chrome's delayed handle release after +// Worker.terminate() never blocks the next worker's file open. +const ALIGN_DB_PREFIX: &str = "connetto-yew-align"; +const ALIGN_HUB_META: &str = "connetto-yew-align-hub-meta.sqlite"; +const ALIGN_AUTH_DB: &str = "connetto-yew-align-auth.sqlite"; -/// The registered caller-identity function connetto installs on every connection. -pub const CALLER_FUNCTION: &str = "current_app_user"; +const PHOTO_DB_PREFIX: &str = "connetto-yew-photo"; +const PHOTO_HUB_META: &str = "connetto-yew-photo-hub-meta.sqlite"; +const PHOTO_AUTH_DB: &str = "connetto-yew-photo-auth.sqlite"; /// The schema version this build was compiled against. #[must_use] @@ -46,8 +38,6 @@ pub fn demo_schema_version() -> connetto_core::SchemaVersion { connetto_core::SchemaVersion::from_source(SCHEMA_SQL) } -// The uuidv4 SQL function registered on every connection so the orders -// and photos DEFAULT (uuidv4()) mints a UUID on local writes. #[diesel::declare_sql_function] extern "SQL" { /// Client-authored primary key: a 16-byte UUID v4, stored as a BLOB. @@ -68,7 +58,7 @@ pub fn uuidv4_functions() -> connetto_client::SqlFunctions { )) } -/// The policy table map for this build, for `ClientConfig::with_policy_tables`. +/// The policy table map for this build. #[must_use] pub fn demo_policy_tables() -> connetto_client::PolicyTables { connetto_client::PolicyTables::from_translation( @@ -77,26 +67,42 @@ pub fn demo_policy_tables() -> connetto_client::PolicyTables { ) } -/// DB worker entry point: boot the connetto DB tier with the photo config. +/// Worker entry point for the alignment test; uses `connetto-yew-align*` OPFS files. +/// +/// # Errors /// -/// The test's blob worker bootstrap imports this crate's wasm module and awaits this. +/// A string describing the VFS, upstream connect, or subscribe failure. +#[wasm_bindgen] +pub async fn db_worker_boot_align() -> Result<(), JsValue> { + boot_with(ALIGN_DB_PREFIX, ALIGN_HUB_META, ALIGN_AUTH_DB).await +} + +/// Worker entry point for the photo test; uses `connetto-yew-photo*` OPFS files. /// /// # Errors /// /// A string describing the VFS, upstream connect, or subscribe failure. #[wasm_bindgen] pub async fn db_worker_photo_boot() -> Result<(), JsValue> { + boot_with(PHOTO_DB_PREFIX, PHOTO_HUB_META, PHOTO_AUTH_DB).await +} + +async fn boot_with( + db_prefix: &'static str, + hub_meta: &'static str, + auth_db: &'static str, +) -> Result<(), JsValue> { connetto_web::logging::init_console(); connetto_web::workers::boot_db_worker::( &connetto_web::workers::DbWorkerConfig::new(demo_schema_version()) .with_ws_url(DEMO_WS_URL) - .with_replica_db_prefix(DB_NAME) + .with_replica_db_prefix(db_prefix) .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-photo-yew-hub-meta.sqlite") + .with_hub_meta_name(hub_meta) .with_content_namespace("connetto-photo-content") .with_sql_functions(uuidv4_functions()) .with_policy_tables(demo_policy_tables()) @@ -106,7 +112,7 @@ pub async fn db_worker_photo_boot() -> Result<(), JsValue> { "dev-idp", "http://127.0.0.1:18099/dev/landing", ))) - .with_auth_db_name("connetto-photo-yew-auth.sqlite"), + .with_auth_db_name(auth_db), ) .await .map(drop) diff --git a/examples/yew-web-demo/tests/photo_flow.rs b/examples/yew-web-demo/tests/photo_flow.rs index 4f8dea38..38e1e362 100644 --- a/examples/yew-web-demo/tests/photo_flow.rs +++ b/examples/yew-web-demo/tests/photo_flow.rs @@ -1,17 +1,18 @@ //! Photo surface browser tests for the yew web demo. //! -//! Drives the real demo boot path (`db_worker_photo_boot`) through the browser -//! stack. Two tests run in this binary and are serialized by a shared Web Lock -//! to avoid OPFS conflicts between concurrent workers. +//! COUPLING: `SchemaVersion` is a hash of `schema.sql`. This file must be +//! byte-identical to `examples/wasm-smoke/schema.sql`. The first test is the +//! permanent drift detector: it fails at the WebSocket handshake the instant +//! the two schemas diverge. Maintain identity with +//! `cp examples/wasm-smoke/schema.sql examples/yew-web-demo/schema.sql`. //! -//! COUPLING: `SchemaVersion` is a hash of the schema SQL source string. -//! `schema.sql` in this workspace MUST be byte-identical to -//! `examples/wasm-smoke/schema.sql`, which is what the browser stack server -//! is compiled against. The first test (`a_tab_order_reaches_the_replica…`) -//! acts as the permanent drift detector: it fails at the WebSocket handshake -//! the moment the two schemas diverge, long before any photo assertion runs. -//! Maintain the identity with `cp examples/wasm-smoke/schema.sql -//! examples/yew-web-demo/schema.sql` whenever wasm-smoke's schema changes. +//! ISOLATION: each test uses a distinct `db_worker_boot_*` function whose OPFS +//! filenames are unique, so Chrome's delayed `FileSystemSyncAccessHandle` +//! release after `Worker.terminate()` never blocks the next worker's file open. +//! `DB_ALIVE_LOCK` is shared across all connetto workers on the same origin; +//! each test holds its own `SUITE_LOCK_*` that serialises that test's lifecycle +//! and a 200 ms sleep before releasing the lock gives Chrome time to finish the +//! OPFS cleanup before the next worker starts. #![cfg(target_arch = "wasm32")] @@ -42,16 +43,15 @@ use web_sys::{ wasm_bindgen_test_configure!(run_in_dedicated_worker); -// Auth server coordinates matching the browser stack setup. const AUTH_BASE: &str = "http://127.0.0.1:18099"; const AUTH_LANDING: &str = "http://127.0.0.1:18099/dev/landing"; const AUTH_PROVIDER: &str = "dev-idp"; const AUTH_USERNAME: &str = "startup"; -// Serializes tests in this binary to prevent OPFS worker conflicts. -const SUITE_LOCK: &str = "connetto-yew-photo-suite"; - -// --- Diesel table schema matching DEMO_TAB_DDL (policy-split, uses logical names) --- +// Separate lock names: each test holds its own lock for its full lifecycle so +// DB_ALIVE_LOCK is never contended between a not-yet-GC'd worker and the next. +const SUITE_LOCK_ALIGN: &str = "connetto-yew-photo-suite-align"; +const SUITE_LOCK_PHOTO: &str = "connetto-yew-photo-suite-photo"; diesel::table! { orders (id) { @@ -91,8 +91,6 @@ struct Photo { content_state: Option, } -// --- Helpers --- - fn stage(msg: &str) { web_sys::console::log_1(&msg.into()); } @@ -171,7 +169,7 @@ async fn walk_login(login_url: &str) -> (String, String) { let req = Request::new_with_str_and_init(&form_url, &init).expect("login request"); req.headers() .set("content-type", "application/x-www-form-urlencoded") - .expect("content-type header"); + .expect("header"); let resp = global_fetch_req(&req).await; assert!( resp.ok(), @@ -180,7 +178,7 @@ async fn walk_login(login_url: &str) -> (String, String) { resp.status() ); let final_url = resp.url(); - let parsed = web_sys::Url::new(&final_url).expect("parse final url"); + let parsed = web_sys::Url::new(&final_url).expect("parse url"); let params = parsed.search_params(); ( params @@ -204,11 +202,7 @@ async fn mint_session() -> (String, String) { WorkerAuthConfig::new(AUTH_BASE, AUTH_PROVIDER, AUTH_LANDING), None, ); - let pending = match auth - .acquire::(&store) - .await - .expect("acquire session") - { + let pending = match auth.acquire::(&store).await.expect("acquire") { Acquired::NeedLogin(p) => p, Acquired::Access(_) => panic!("fresh store cannot refresh silently"), }; @@ -233,25 +227,24 @@ fn glue_url() -> String { format!("{base}.js") } -fn spawn_photo_worker(glue_url: &str) -> Worker { +fn spawn_worker(glue_url: &str, boot_fn: &str) -> Worker { let wasm_url = glue_url .strip_suffix(".js") .map_or_else(|| format!("{glue_url}_bg.wasm"), |b| format!("{b}_bg.wasm")); let src = format!( "const ch=new BroadcastChannel('connetto-debug');\n\ - try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n ch.postMessage('yew photo worker ready');\n await mod.db_worker_photo_boot();\n}}catch(e){{\n ch.postMessage('yew photo worker FAILED: '+e);\n throw e;\n}}", + try{{\n const mod=await import({g:?});\n await mod.default({{module_or_path:{w:?}}});\n await mod.{f}();\n}}catch(e){{\n ch.postMessage('yew worker FAILED: '+e);\n throw e;\n}}", g = glue_url, w = wasm_url, + f = boot_fn, ); let parts = Array::of1(&JsValue::from_str(&src)); let opts = web_sys::BlobPropertyBag::new(); opts.set_type("text/javascript"); - let blob = - web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("bootstrap blob"); + let blob = web_sys::Blob::new_with_str_sequence_and_options(&parts, &opts).expect("blob"); let url = web_sys::Url::create_object_url_with_blob(&blob).expect("object url"); let worker_opts = web_sys::WorkerOptions::new(); worker_opts.set_type(web_sys::WorkerType::Module); - worker_opts.set_name("connetto-yew-photo"); let w = Worker::new_with_options(&url, &worker_opts).expect("spawn worker"); web_sys::Url::revoke_object_url(&url).ok(); w @@ -278,10 +271,10 @@ fn blob_of(bytes: &[u8]) -> web_sys::Blob { async fn fetch_bytes(url: &str) -> Vec { let scope = js_sys::global() .dyn_into::() - .expect("dedicated worker"); + .expect("worker"); let resp: web_sys::Response = JsFuture::from(scope.fetch_with_str(url)) .await - .expect("fetch content") + .expect("fetch") .dyn_into() .expect("Response"); assert!( @@ -326,18 +319,18 @@ async fn connect_tab( (content, conn) } -// --- Test 1: alignment proof — version agreement via an orders round trip --- +// --- Test 1: alignment proof — proves version agreement before the photo test --- #[wasm_bindgen_test] async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { relay_worker_breadcrumbs(); play_the_tab(); - let _serial = locks::hold_lock(SUITE_LOCK).await; - let worker = spawn_photo_worker(&glue_url()); + let _serial = locks::hold_lock(SUITE_LOCK_ALIGN).await; + let worker = spawn_worker(&glue_url(), "db_worker_boot_align"); workers::await_db_worker_ready(&[]) .await .expect("db worker ready"); - stage("yew demo worker booted (align)"); + stage("yew align worker booted"); let (token, identity) = mint_session().await; let client_id = rosetta_uuid::Uuid::new_v4().to_string(); @@ -345,7 +338,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { let (_content, mut conn) = connect_tab(&client_id, token, &identity).await; conn.subscribe("align-orders", "SELECT * FROM orders") .await - .expect("orders subscribe"); + .expect("subscribe"); loop { let event = conn.pump_one().await.expect("pump"); assert_ne!(event, ClientEvent::Closed, "connection closed early"); @@ -353,7 +346,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { break; } } - stage("orders subscription ready (align)"); + stage("orders subscription ready"); let (client, pump) = ConnettoClient::with_pump(conn); let (done_tx, done_rx) = oneshot::channel::<()>(); @@ -367,7 +360,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { .select(Order::as_select()) .live(&client) .await - .expect("orders live query"); + .expect("live query"); let qty = 77_i64; let id_str = identity.clone(); @@ -383,7 +376,7 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { .await .expect("order insert"); client.replay_pending().await.expect("replay pending"); - stage("order written (align)"); + stage("order written"); loop { if live.rows().iter().any(|o| o.quantity == qty) { @@ -391,26 +384,29 @@ async fn a_tab_order_reaches_the_replica_through_the_demo_boot_path() { } live.changed().await.expect("live refresh"); } - stage("order arrived at replica — version agreement confirmed"); + stage("order arrived — version agreement confirmed"); drop(live); drop(client); done_rx.await.expect("pump exited"); worker.terminate(); + // Chrome headless does not release the OPFS access handle synchronously on + // terminate; per-test file names are the isolation guarantee, this is margin. + workers::sleep(core::time::Duration::from_millis(200)).await; } -// --- Test 2: photo round trip through stage, commit, resolve, and HTTP fetch --- +// --- Test 2: photo round trip --- #[wasm_bindgen_test] async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { relay_worker_breadcrumbs(); play_the_tab(); - let _serial = locks::hold_lock(SUITE_LOCK).await; - let worker = spawn_photo_worker(&glue_url()); + let _serial = locks::hold_lock(SUITE_LOCK_PHOTO).await; + let worker = spawn_worker(&glue_url(), "db_worker_photo_boot"); workers::await_db_worker_ready(&[]) .await .expect("db worker ready"); - stage("yew demo worker booted (photo)"); + stage("yew photo worker booted"); let (token, identity) = mint_session().await; let client_id = rosetta_uuid::Uuid::new_v4().to_string(); @@ -418,7 +414,7 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { let (content, mut conn) = connect_tab(&client_id, token, &identity).await; conn.subscribe("photo-test-photos", "SELECT * FROM photos") .await - .expect("photos subscribe"); + .expect("subscribe"); loop { let event = conn.pump_one().await.expect("pump"); assert_ne!(event, ClientEvent::Closed, "connection closed early"); @@ -440,7 +436,7 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { .select(Photo::as_select()) .live(&client) .await - .expect("photo live query"); + .expect("live query"); let bytes = photo_bytes(); let blob = blob_of(&bytes); @@ -514,11 +510,14 @@ async fn a_photo_round_trips_through_stage_commit_resolve_and_http_fetch() { TabResolved::Unavailable => panic!("available photo must resolve"), }; assert_eq!(fetch_bytes(&url).await, bytes); - stage("photo bytes verified via HTTP fetch"); + stage("photo bytes verified"); drop(live); drop(content); drop(client); done_rx.await.expect("pump exited"); worker.terminate(); + // Chrome headless does not release the OPFS access handle synchronously on + // terminate; per-test file names are the isolation guarantee, this is margin. + workers::sleep(core::time::Duration::from_millis(200)).await; }