diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 04e1b23f..d3e32551 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -404,6 +404,7 @@ jobs: -p connetto-yew -p connetto-file-core -p connetto-file-server + -p connetto-file-client bins: '' steps: - uses: actions/checkout@v7 @@ -509,7 +510,6 @@ jobs: echo "${CHROMEWEBDRIVER}" >> "${GITHUB_PATH}" echo "CHROMEDRIVER=${CHROMEWEBDRIVER}/chromedriver" >> "${GITHUB_ENV}" - # The identity-provider image is pulled inside each shard's startup # wait, which a congested ghcr.io can overrun. Pulling it once here # takes the registry out of the container startup budget. diff --git a/crates/connetto-file-client/src/error.rs b/crates/connetto-file-client/src/error.rs index 94bdb7ee..167c6bcf 100644 --- a/crates/connetto-file-client/src/error.rs +++ b/crates/connetto-file-client/src/error.rs @@ -112,6 +112,24 @@ pub enum ContentError { }, } +/// Why a staged commit did not land, keeping the row's refusal distinguishable +/// from a bookkeeping failure so the caller can word the answer. +#[derive(Debug, Error)] +pub enum StageCommitError { + /// Recording the manifest or queueing the upload failed. + #[error("staged content bookkeeping: {0}")] + Bookkeeping(#[from] ContentError), + /// The row closure refused, carrying the detail it gave. + #[error("staged row refused: {0}")] + Row(String), +} + +impl From for StageCommitError { + fn from(err: diesel::result::Error) -> Self { + Self::Bookkeeping(err.into()) + } +} + /// What the outbox walk does with a failed upload attempt, decided by where the fact came from. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum AttemptOutcome { diff --git a/crates/connetto-file-client/src/lib.rs b/crates/connetto-file-client/src/lib.rs index e0b4130b..5d397685 100644 --- a/crates/connetto-file-client/src/lib.rs +++ b/crates/connetto-file-client/src/lib.rs @@ -19,15 +19,15 @@ mod worker; pub use browser_store::{BrowserStore, BrowserStoreError}; pub use client::{ContentClient, ContentEvent}; -pub use connetto_file_core::{ChunkHash, FileId}; -pub use error::{AttemptOutcome, ContentError}; +pub use connetto_file_core::{ChunkHash, FileId, FileIdHasher, MimeClass}; +pub use error::{AttemptOutcome, ContentError, StageCommitError}; pub use import::ContentImportPlan; pub use resolve::{BoxedSource, ChunkStoreSource, LocalContentSource, Resolved, SourceFuture}; #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] pub use store::{FsStore, FsStoreError}; pub use worker::{ ChunkScan, ContentArchive, ContentFlush, ContentFlushStart, ContentFlushState, ContentUpload, - FlushCursor, ScanStep, + FlushCursor, PendingConnectionResolve, ResolveRoute, ResolveStart, ScanStep, }; #[cfg(all(target_family = "wasm", target_os = "unknown"))] diff --git a/crates/connetto-file-client/src/ticket.rs b/crates/connetto-file-client/src/ticket.rs index a87ca3ff..cb8d1ec9 100644 --- a/crates/connetto-file-client/src/ticket.rs +++ b/crates/connetto-file-client/src/ticket.rs @@ -88,7 +88,7 @@ where } } -async fn ensure_pending_request( +pub(crate) async fn ensure_pending_request( connection: &mut ConnettoConnection, file_id: FileId, verb: ContentVerb, @@ -170,6 +170,26 @@ fn resolve_connection_event( Some(result) } +/// What one pumped event has to say about one in-flight ticket request, +/// without consuming the event or the waiting slot. +pub(crate) enum RouteAnswer { + /// The request's answer arrived, or the link carrying it is gone. + Settled(Result), + /// The event concerns something else. + Other, +} + +pub(crate) fn route_connection_event(pending: &PendingTicket, event: &ClientEvent) -> RouteAnswer { + match classify(event.clone(), &pending.request_id, pending.file_id) { + TicketAnswer::Granted(url) => RouteAnswer::Settled(Ok(url)), + TicketAnswer::Failed(error) => RouteAnswer::Settled(Err(error)), + TicketAnswer::Unrelated(ClientEvent::Closed | ClientEvent::ServerClosed { .. }) => { + RouteAnswer::Settled(Err(ContentError::TicketAbandoned)) + } + TicketAnswer::Unrelated(_) => RouteAnswer::Other, + } +} + fn next_request_id() -> String { format!("content-{}", NEXT_REQUEST.fetch_add(1, Ordering::Relaxed)) } @@ -185,3 +205,56 @@ fn refusal(detail: &str, file_id: FileId) -> ContentError { ); ContentError::TicketRefused { file_id } } + +#[cfg(test)] +mod tests { + use super::*; + + fn waiting() -> PendingTicket { + PendingTicket { + file_id: FileId::from_bytes([1; 32]), + request_id: "content-7".to_owned(), + } + } + + #[test] + fn a_routed_event_settles_only_its_own_request() { + let pending = waiting(); + let grant = ClientEvent::ContentTicket { + request_id: "content-7".to_owned(), + url: "https://content.invalid/f".to_owned(), + }; + assert!(matches!( + route_connection_event(&pending, &grant), + RouteAnswer::Settled(Ok(url)) if url == "https://content.invalid/f" + )); + let someone_elses_grant = ClientEvent::ContentTicket { + request_id: "content-8".to_owned(), + url: "https://content.invalid/f".to_owned(), + }; + assert!(matches!( + route_connection_event(&pending, &someone_elses_grant), + RouteAnswer::Other + )); + let refusal = ClientEvent::NonFatal { + related_to: Some("content-7".to_owned()), + detail: CONTENT_TICKET_REFUSED.to_owned(), + }; + assert!(matches!( + route_connection_event(&pending, &refusal), + RouteAnswer::Settled(Err(_)) + )); + let unrelated = ClientEvent::NonFatal { + related_to: Some("subscription-1".to_owned()), + detail: "something else failed".to_owned(), + }; + assert!(matches!( + route_connection_event(&pending, &unrelated), + RouteAnswer::Other + )); + assert!(matches!( + route_connection_event(&pending, &ClientEvent::Closed), + RouteAnswer::Settled(Err(ContentError::TicketAbandoned)) + )); + } +} diff --git a/crates/connetto-file-client/src/worker.rs b/crates/connetto-file-client/src/worker.rs index 043e7edf..c7f05317 100644 --- a/crates/connetto-file-client/src/worker.rs +++ b/crates/connetto-file-client/src/worker.rs @@ -2,20 +2,23 @@ use core::fmt::Display; use std::collections::HashSet; +use std::io::Read; use connetto_client::{ClientEvent, ConnettoConnection, ExportScope, ImportChoices, ImportOutcome}; use connetto_core::messages::ContentVerb; use connetto_core::traits::Transport; use connetto_file_core::{ ChunkHash, ChunkStore, EncryptStoreError, EncryptingStore, FileId, Manifest, MaybeSend, + MimeClass, process_file_from_reader, }; use diesel::connection::SimpleConnection; use diesel::prelude::*; use crate::db; -use crate::error::{AttemptOutcome, ContentError}; +use crate::error::{AttemptOutcome, ContentError, StageCommitError}; use crate::http::ContentHttp; use crate::import::{apply_content_import, prepare_content_import, write_import_chunks}; +use crate::resolve::{ChunkStoreSource, LocalContentSource, Resolved}; use crate::{ticket, upload}; /// Result of one worker-owned outbox attempt. @@ -152,6 +155,49 @@ where } } +/// The decision of a [`ContentArchive::start_resolve_connection`] at the +/// moment it was asked. +pub enum ResolveStart { + /// The answer is already in hand: the local store had the bytes, or the + /// connection cannot ask the server. + Answered(Resolved), + /// The read-ticket request is in flight; pumped events settle it through + /// [`PendingConnectionResolve::route`]. + Waiting(PendingConnectionResolve), +} + +/// A connection resolve whose ticket request is in flight. +/// +/// Opaque: the caller keeps it beside its own request bookkeeping until +/// routing an event settles the wait. +pub struct PendingConnectionResolve { + ticket: ticket::PendingTicket, +} + +impl PendingConnectionResolve { + /// What one pumped upstream event has to say about this waiting resolve. + /// + /// `Settled` ends the wait; `Err` means answer + /// [`Resolved::Unavailable`]. `Other` leaves the event to the caller's + /// own handling, the same events [`resolve_connection`](ContentArchive::resolve_connection) + /// returns as observed. + pub fn route(&self, event: &ClientEvent) -> ResolveRoute { + match ticket::route_connection_event(&self.ticket, event) { + ticket::RouteAnswer::Settled(result) => ResolveRoute::Settled(result), + ticket::RouteAnswer::Other => ResolveRoute::Other, + } + } +} + +/// What routing one pumped event did to a waiting resolve. +pub enum ResolveRoute { + /// The wait ends. `Ok(url)` is the granted read; `Err` means answer + /// [`Resolved::Unavailable`] and log. + Settled(Result), + /// The event belongs to something else; the caller still owns it. + Other, +} + /// Content archive policy for an owner of a raw sync connection. pub struct ContentArchive { store: B, @@ -182,6 +228,220 @@ where db::add_refused_column(conn).map_err(ContentError::from) } + /// Chunks one file into the encrypted store without touching the replica. + /// + /// The slow half of staging a file, so a hub can keep serving local work + /// while a large file is split. [`commit_staged`](Self::commit_staged) + /// closes the pair in one transaction. The identity inside the returned + /// manifest is what the bytes hash to, and it is the only identity the + /// chunks can ever be served under. + /// + /// # Errors + /// + /// [`ContentError::Store`] when reading the bytes or writing a chunk fails. + pub async fn chunk_file(&self, reader: R, mime: MimeClass) -> Result + where + R: Read + MaybeSend, + { + let store = EncryptingStore::new_with( + self.store.clone(), + &self.root_key, + mime.params().skip_compression, + ); + process_file_from_reader(reader, mime, &store) + .await + .map_err(|err| ContentError::Store(err.to_string())) + } + + /// Commits staged content and the row that names it as one transaction. + /// + /// The raw-connection shape of [`ContentClient::stage`](crate::ContentClient::stage): + /// the manifest and its outbox entry are recorded with capture suspended, + /// and `row` then runs with capture live, so the row syncs as part of the + /// same mutation and the upload leg sees the queued entry. `row` receives + /// the identity the chunked bytes hash to; a caller told an identity + /// elsewhere compares it here, because bytes that hash to something else + /// can never be served under the declared one. + /// + /// # Errors + /// + /// Anything `row` returns, or [`StageCommitError::Bookkeeping`] when the + /// manifest or outbox write fails. Either way nothing committed, and the + /// chunks stand until the next orphan sweep. + pub fn commit_staged( + &self, + connection: &mut ConnettoConnection, + manifest: &Manifest, + row: F, + ) -> Result + where + T: Transport, + F: FnOnce(&mut SqliteConnection, FileId) -> Result, + { + let file_id = manifest.file_id(); + connection + .transact_with_bookkeeping( + |conn| { + db::put_manifest(conn, manifest)?; + db::enqueue(conn, file_id)?; + Ok(()) + }, + |conn| row(conn, file_id), + ) + .map(|((), outcome)| outcome) + } + + /// The answer about this file's bytes this device can give on its own. + /// + /// The shape of [`ContentClient`](crate::ContentClient)'s own local + /// answer: unsent content answers from the chunk store or nowhere, pinned + /// content prefers what the pin paid to keep, and anything else answers + /// `None`, which means only a server knows. + /// + /// # Errors + /// + /// [`ContentError::Replica`] on a bookkeeping read failure. + pub async fn resolve_local( + &self, + connection: &mut ConnettoConnection, + file_id: FileId, + ) -> Result, ContentError> { + let Some(manifest) = db::load_manifest(connection.conn(), file_id)? else { + return Ok(None); + }; + let source = + ChunkStoreSource::new(EncryptingStore::new(self.store.clone(), &self.root_key)); + if db::is_unsent(connection.conn(), file_id)? { + return Ok(Some(match source.bytes(&manifest).await? { + Some(bytes) => Resolved::Local { + source: source.name(), + bytes, + }, + None => Resolved::Unavailable, + })); + } + if !crate::retain::pinned_ids(connection.conn())?.contains(&file_id) { + return Ok(None); + } + Ok(source.bytes(&manifest).await?.map(|bytes| Resolved::Local { + source: source.name(), + bytes, + })) + } + + /// Where this file's bytes are to be had, answering the way the owning + /// client would, for a hub that serves resolution questions for others. + /// + /// [`resolve_local`](Self::resolve_local) first, then the server for + /// anything local sources cannot answer, exactly as + /// [`ContentClient::resolve`](crate::ContentClient::resolve) does. A + /// `cancel` that fires, a server refusal and an offline connection all + /// answer [`Resolved::Unavailable`], because the question deserves an + /// answer rather than a wait. Events pumped while the ticket round trip + /// runs come back in the returned vector for the caller to re-apply; a + /// cancelled wait re-applies them without an answer. + /// + /// There is no error to handle. A bookkeeping failure, a server refusal, + /// a transport failure and a cancelled wait all answer + /// [`Resolved::Unavailable`], the only answer better than no answer. + pub async fn resolve_connection( + &self, + connection: &mut ConnettoConnection, + file_id: FileId, + cancel: C, + ) -> (Resolved, Vec) + where + T: Transport, + T::Error: Display, + C: core::future::Future, + { + let mut observed = Vec::new(); + match self + .answer_connection(connection, file_id, cancel, &mut observed) + .await + { + Ok(answer) => (answer, observed), + Err(_) => (Resolved::Unavailable, observed), + } + } + + /// Starts a connection resolve without waiting for the server. + /// + /// [`resolve_local`](Self::resolve_local) answers first and + /// `Answered(Resolved::Unavailable)` follows when the connection cannot + /// ask; otherwise the read-ticket request goes out and the handle comes + /// back in `Waiting`, to be settled event by event with + /// [`route`](PendingConnectionResolve::route). This is the shape + /// a hub that must keep serving its queue uses instead of + /// [`resolve_connection`](Self::resolve_connection), which holds a task + /// parked until the answer or the cancel. + /// + /// # Errors + /// + /// A replica read failure or a transport failure on the way to sending + /// the ticket request is returned. Nothing else is an error: a refusal + /// and a closed link both answer through the routed events. + /// + /// # Panics + /// + /// Never. A request that reports success always leaves its handle + /// behind, and the handle is what the returned waiting state carries. + pub async fn start_resolve_connection( + &self, + connection: &mut ConnettoConnection, + file_id: FileId, + ) -> Result + where + T: Transport, + T::Error: Display, + { + if let Some(answer) = self.resolve_local(connection, file_id).await? { + return Ok(ResolveStart::Answered(answer)); + } + if !connection.is_connected() { + return Ok(ResolveStart::Answered(Resolved::Unavailable)); + } + let mut pending = None; + ticket::ensure_pending_request(connection, file_id, ContentVerb::Read, &mut pending) + .await?; + let ticket = pending.expect("a sent request leaves its handle behind"); + Ok(ResolveStart::Waiting(PendingConnectionResolve { ticket })) + } + + async fn answer_connection( + &self, + connection: &mut ConnettoConnection, + file_id: FileId, + cancel: C, + observed: &mut Vec, + ) -> Result + where + T: Transport, + T::Error: Display, + C: core::future::Future, + { + if let Some(answer) = self.resolve_local(connection, file_id).await? { + return Ok(answer); + } + if !connection.is_connected() { + return Ok(Resolved::Unavailable); + } + let mut pending = None; + let url = ticket::request_connection_or( + connection, + file_id, + ContentVerb::Read, + observed, + cancel, + &mut pending, + ) + .await?; + Ok(match url { + Some(url) => Resolved::Remote { url }, + None => Resolved::Unavailable, + }) + } + /// Counts content files that have not reached the server. /// /// # Errors @@ -2079,4 +2339,218 @@ mod tests { "sendable_files counts only unmarked entries" ); } + + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + #[derive(diesel::QueryableByName)] + struct Count { + #[diesel(sql_type = diesel::sql_types::BigInt)] + n: i64, + } + + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + fn count(connection: &mut diesel::SqliteConnection, sql: &str) -> i64 { + use diesel::RunQueryDsl; + diesel::sql_query(sql) + .get_result::(connection) + .expect("the count reads") + .n + } + + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + fn staged_fixture( + name: &str, + ) -> ( + tempfile::TempDir, + super::ContentArchive, + connetto_client::ConnettoConnection, + ) { + use connetto_client::{ClientConfig, ConnettoConnection, Replica}; + use connetto_core::test_support::FakeTransport; + + let dir = tempfile::tempdir().expect("a temporary directory"); + let store = crate::store::FsStore::new(dir.path().join("chunks")); + let mut connection = ConnettoConnection::::open( + &Replica::in_memory(), + "CREATE TABLE photos (id INTEGER PRIMARY KEY, content_id BLOB NOT NULL)", + &ClientConfig::new(name), + None, + ) + .expect("the replica opens offline"); + let archive = super::ContentArchive::new(store, [1; 32]); + archive.install(&mut connection).expect("content tables"); + (dir, archive, connection) + } + + /// The manifest, the outbox entry and the naming row commit together, and + /// the identity they carry is the one the bytes hash to. + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + #[tokio::test] + async fn a_staged_file_and_its_row_commit_together() { + use crate::db; + use connetto_file_core::MimeClass; + use diesel::RunQueryDsl; + + let (_dir, archive, mut connection) = staged_fixture("stage-commit"); + let bytes = vec![7u8; 1024]; + let manifest = archive + .chunk_file(&bytes[..], MimeClass::Jpeg) + .await + .expect("the bytes chunk"); + let expected = manifest.file_id(); + + let written = archive + .commit_staged(&mut connection, &manifest, |conn, id| { + diesel::sql_query("INSERT INTO photos (id, content_id) VALUES (1, ?)") + .bind::(id.as_bytes().to_vec()) + .execute(conn)?; + Ok(id) + }) + .expect("the staged commit lands"); + + assert_eq!(written, expected, "the row is told the computed identity"); + assert!( + db::load_manifest(connection.conn(), expected) + .expect("the manifest reads") + .is_some(), + "the manifest committed with the row" + ); + assert_eq!( + db::outbox(connection.conn()).expect("the outbox reads"), + vec![expected], + "the upload was queued with the row" + ); + assert_eq!( + count(connection.conn(), "SELECT COUNT(*) AS n FROM photos"), + 1, + "the row committed" + ); + } + + /// A row closure that refuses the computed identity rolls back the + /// manifest and the outbox entry with it, leaving nothing staged behind. + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + #[tokio::test] + async fn a_row_that_refuses_the_computed_id_rolls_the_manifest_back() { + use crate::db; + use crate::error::StageCommitError; + use connetto_file_core::MimeClass; + use diesel::RunQueryDsl; + + let (_dir, archive, mut connection) = staged_fixture("stage-mismatch"); + let bytes = vec![8u8; 1024]; + let manifest = archive + .chunk_file(&bytes[..], MimeClass::Jpeg) + .await + .expect("the bytes chunk"); + let declared = FileId::from_bytes([9; 32]); + + let error = archive + .commit_staged(&mut connection, &manifest, |conn, id| { + if id != declared { + return Err(StageCommitError::Row(format!( + "declared {declared}, the bytes hash to {id}" + ))); + } + diesel::sql_query("INSERT INTO photos (id, content_id) VALUES (1, ?)") + .bind::(id.as_bytes().to_vec()) + .execute(conn)?; + Ok(id) + }) + .expect_err("a mismatched identity must refuse"); + assert!(matches!(error, StageCommitError::Row(_))); + + assert!( + db::load_manifest(connection.conn(), manifest.file_id()) + .expect("the manifest reads") + .is_none(), + "a refused row leaves no manifest behind" + ); + assert!( + db::outbox(connection.conn()) + .expect("the outbox reads") + .is_empty(), + "a refused row queues no upload" + ); + assert_eq!( + count(connection.conn(), "SELECT COUNT(*) AS n FROM photos"), + 0, + "nothing wrote to the application table" + ); + } + + /// Content the worker just staged answers locally before any upload, and + /// answers with the exact bytes. + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + #[tokio::test] + async fn unsent_staged_content_resolves_from_the_chunk_store() { + use crate::resolve::Resolved; + use connetto_file_core::MimeClass; + + let (_dir, archive, mut connection) = staged_fixture("stage-resolve"); + let bytes = vec![5u8; 2048]; + let manifest = archive + .chunk_file(&bytes[..], MimeClass::Jpeg) + .await + .expect("the bytes chunk"); + archive + .commit_staged(&mut connection, &manifest, |_conn, id| Ok(id)) + .expect("the staged commit lands"); + + let (answer, observed) = archive + .resolve_connection( + &mut connection, + manifest.file_id(), + core::future::pending::<()>(), + ) + .await; + match answer { + Resolved::Local { bytes: served, .. } => assert_eq!(served, bytes), + other => panic!("unsent content must resolve locally, got {other:?}"), + } + assert!(observed.is_empty(), "no round trip ran"); + } + + /// Content the worker cannot produce answers `Unavailable`: a file with no + /// manifest, and an uploaded file no pin covers, both with no server + /// online. + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + #[tokio::test] + async fn content_the_worker_cannot_produce_resolves_unavailable_offline() { + use crate::db; + use crate::resolve::Resolved; + use connetto_file_core::{EncryptingStore, MimeClass, process_file}; + + let (dir, archive, mut connection) = staged_fixture("stage-unavailable"); + let (answer, _) = archive + .resolve_connection( + &mut connection, + FileId::from_bytes([3; 32]), + core::future::pending::<()>(), + ) + .await; + assert!( + matches!(answer, Resolved::Unavailable), + "an unknown file must resolve Unavailable, got {answer:?}" + ); + + let encrypted = EncryptingStore::new( + crate::store::FsStore::new(dir.path().join("chunks")), + &[1; 32], + ); + let uploaded = process_file(&vec![4u8; 1024], MimeClass::Jpeg, &encrypted) + .await + .expect("the file chunks"); + db::put_manifest(connection.conn(), &uploaded).expect("record the manifest"); + let (answer, _) = archive + .resolve_connection( + &mut connection, + uploaded.file_id(), + core::future::pending::<()>(), + ) + .await; + assert!( + matches!(answer, Resolved::Unavailable), + "an uploaded unpinned file offline must resolve Unavailable, got {answer:?}" + ); + } } diff --git a/crates/connetto-file-client/tests/it/resolving.rs b/crates/connetto-file-client/tests/it/resolving.rs index 76f46f62..018bbccf 100644 --- a/crates/connetto-file-client/tests/it/resolving.rs +++ b/crates/connetto-file-client/tests/it/resolving.rs @@ -5,8 +5,9 @@ use connetto_file_core::MimeClass; use tempfile::tempdir; use crate::support::{ - RecordingHttp, Scripted, TicketAnswer, attach_content, connected_client, connected_content, - insert_row_and_pin_album, learn_file_id, offline_client, stage_photo, + ROOT_KEY, RecordingHttp, Scripted, TicketAnswer, attach_content, cold_connection, + connected_client, connected_connection, connected_content, insert_row_and_pin_album, + learn_file_id, offline_client, stage_photo, }; /// A granted read address, the shape the file server mints for a download. @@ -289,3 +290,94 @@ async fn a_refused_read_ticket_is_reported() { "a caller learns it may not read, and learns nothing else, got {refused:?}" ); } + +/// The split resolve API a relay hub drives: the request leaves without the +/// caller waiting, foreign events pass by untouched, and the grant this +/// request asked for settles it. +#[tokio::test] +async fn a_started_resolve_waits_and_settles_on_its_own_grant() { + use connetto_client::ClientEvent; + use connetto_file_client::{ContentArchive, FsStore, ResolveRoute, ResolveStart}; + + let dir = tempdir().expect("temp dir"); + let mut conn = connected_connection( + &dir.path().join("replica.sqlite"), + Scripted::granting(READ_URL), + ) + .await; + let archive = ContentArchive::new(FsStore::new(dir.path().join("chunks")), ROOT_KEY); + archive + .install(&mut conn) + .expect("install content bookkeeping"); + let unknown = connetto_file_core::FileId::from_bytes([7; 32]); + + let start = archive + .start_resolve_connection(&mut conn, unknown) + .await + .expect("start a resolve on the live connection"); + let ResolveStart::Waiting(pending) = start else { + panic!("an unknown file on a live connection waits for its ticket"); + }; + assert!( + matches!( + pending.route(&ClientEvent::ContentTicket { + request_id: "content-someone-else".to_owned(), + url: READ_URL.to_owned(), + }), + ResolveRoute::Other + ), + "a grant for another request is not this one's answer" + ); + + let mut settled = None; + for _ in 0..8 { + let event = conn.pump_one().await.expect("a pumped event"); + if matches!(event, ClientEvent::ContentTicket { .. }) + && let ResolveRoute::Settled(result) = pending.route(&event) + { + settled = Some(result); + break; + } + } + assert!( + matches!(&settled, Some(Ok(url)) if url == READ_URL), + "the grant this request asked for settles the wait, got {settled:?}" + ); + + let second = archive + .start_resolve_connection(&mut conn, connetto_file_core::FileId::from_bytes([8; 32])) + .await + .expect("start the second resolve"); + let ResolveStart::Waiting(pending) = second else { + panic!("the second request waits too"); + }; + assert!( + matches!( + pending.route(&ClientEvent::Closed), + ResolveRoute::Settled(Err(_)) + ), + "a closed link abandons the wait" + ); +} + +/// On a connection with no transport there is nobody to ask, and the split +/// API says so at once instead of waiting. +#[tokio::test] +async fn a_started_resolve_on_a_cold_connection_answers_unavailable() { + use connetto_file_client::{ContentArchive, FsStore, ResolveStart}; + + let dir = tempdir().expect("temp dir"); + let mut conn = cold_connection(&dir.path().join("replica.sqlite")); + let archive = ContentArchive::new(FsStore::new(dir.path().join("chunks")), ROOT_KEY); + archive + .install(&mut conn) + .expect("install content bookkeeping"); + let start = archive + .start_resolve_connection(&mut conn, connetto_file_core::FileId::from_bytes([7; 32])) + .await + .expect("a cold start is an answer, not an error"); + assert!( + matches!(start, ResolveStart::Answered(Resolved::Unavailable)), + "no transport means Unavailable at once" + ); +} diff --git a/crates/connetto-file-client/tests/it/support.rs b/crates/connetto-file-client/tests/it/support.rs index fe18eba2..4aacd275 100644 --- a/crates/connetto-file-client/tests/it/support.rs +++ b/crates/connetto-file-client/tests/it/support.rs @@ -446,3 +446,34 @@ pub async fn connected_content( let content = attach_content(client.clone(), &root.join("chunks"), http).await; (client, content) } + +/// Opens a replica at `path` and hands back the raw attached connection. +/// +/// The test owns the pump loop this way, the same shape a relay hub drives +/// the archive's waiting resolve API against. +pub async fn connected_connection( + path: &std::path::Path, + transport: Scripted, +) -> ConnettoConnection { + let replica = Replica::encrypted_file( + path.to_str().expect("utf-8 path"), + Some(connetto_core::test_support::replica_key()), + ) + .expect("a resolved key"); + let mut conn = ConnettoConnection::::open(&replica, DDL, &config(), None) + .expect("open with no server"); + conn.attach(transport).await.expect("attach the transport"); + conn +} + +/// Opens a replica at `path` and hands back the raw connection with no +/// transport, the cold state where asking the server is impossible. +pub fn cold_connection(path: &std::path::Path) -> ConnettoConnection { + let replica = Replica::encrypted_file( + path.to_str().expect("utf-8 path"), + Some(connetto_core::test_support::replica_key()), + ) + .expect("a resolved key"); + ConnettoConnection::::open(&replica, DDL, &config(), None) + .expect("open with no server") +} diff --git a/crates/connetto-web/Cargo.toml b/crates/connetto-web/Cargo.toml index 9f6a38c2..a989fc83 100644 --- a/crates/connetto-web/Cargo.toml +++ b/crates/connetto-web/Cargo.toml @@ -41,6 +41,7 @@ web-sys = { version = "0.3", features = [ "BlobPropertyBag", "BroadcastChannel", "CloseEvent", + "console", "Credential", "CredentialCreationOptions", "CredentialRequestOptions", @@ -103,6 +104,7 @@ diesel-sqlite-session = { git = "https://github.com/LucaCappelletti94/diesel-sql # Routing live patches by table needs the patchset parser, and snapshot # payloads ride the wire zstd compressed like the server's snapshot leg. sqlite-diff-rs = { version = "0.13", default-features = false } +connetto-file-core = { path = "../connetto-file-core" } zstd = "0.13" # tokio::select! races the hub's channels against the upstream pump, and the # tokio mpsc channels carry frames between the hub core and its tab shovels. diff --git a/crates/connetto-web/src/content.rs b/crates/connetto-web/src/content.rs index 9e444a9c..87fe3b4c 100644 --- a/crates/connetto-web/src/content.rs +++ b/crates/connetto-web/src/content.rs @@ -1,16 +1,27 @@ //! Browser display handles for resolved content. use core::fmt::Display; +use std::cell::{Cell, RefCell}; +use std::collections::HashMap; +use std::rc::Rc; use std::sync::Arc; +use crate::content_wire::{ContentFrame, WireResolve, mime_code}; +use crate::frames::{InternalLane, MessageSink, MessageTransport}; +use crate::workers::helpers::sleep_ms; use connetto_client::live::ConnettoClient; use connetto_core::traits::{MaybeSend, Transport}; use connetto_file_client::{ - BrowserHttp, BrowserStore, BrowserStoreError, ContentClient, ContentError, Resolved, + BrowserHttp, BrowserStore, BrowserStoreError, ContentClient, ContentError, FileId, + FileIdHasher, MimeClass, Resolved, }; +use diesel::SqliteConnection; +use futures_channel::oneshot; +use futures_util::StreamExt; use js_sys::{Array, Uint8Array}; use thiserror::Error; use wasm_bindgen::JsCast; +use wasm_bindgen_futures::JsFuture; /// Failure to attach browser file handling. #[derive(Debug, Error)] @@ -98,6 +109,19 @@ impl ObjectUrl { pub fn as_str(&self) -> &str { &self.inner.value } + + /// Creates an object URL for a blob the browser already holds. + /// + /// # Errors + /// + /// [`ObjectUrlError`] when the browser rejects URL creation. + pub fn from_blob(blob: &Blob) -> Result { + let value = + Url::create_object_url_with_blob(blob).map_err(|value| object_url_error(&value))?; + Ok(Self { + inner: Arc::new(ObjectUrlInner { value }), + }) + } } impl AsRef for ObjectUrl { @@ -152,3 +176,202 @@ impl BrowserResolved { fn object_url_error(value: &wasm_bindgen::JsValue) -> ObjectUrlError { ObjectUrlError::Browser(value.as_string().unwrap_or_else(|| format!("{value:?}"))) } + +/// The longest a tab waits for the worker's resolve answer. The hub bounds +/// its own wait sooner, so this only fires when the hub has stopped +/// answering entirely. +const TAB_RESOLVE_WAIT_MS: i32 = 16_000; + +/// The bytes one staging hash read takes from the blob. The identity is the +/// hash of the whole file, but the file goes through one window of this size +/// at a time, so the tab's peak holds a window rather than two copies of the +/// largest file it can upload. +const STAGE_READ_BYTES: f64 = 4.0 * 1024.0 * 1024.0; + +/// Where a file's bytes are to be had, as answered by the worker hub over +/// the tab's content lane. +#[derive(Debug)] +pub enum TabResolved { + /// The server granted a read at this URL. + Remote { + /// The granted URL. + url: String, + }, + /// The worker held the bytes itself. + Local { + /// The file's bytes. + blob: Blob, + }, + /// Neither the worker nor the server can produce the bytes right now. + Unavailable, +} + +impl TabResolved { + /// Exposes a local answer as an object URL. A remote answer is already a + /// URL, and `Unavailable` has no bytes; both answer `None`. + /// + /// # Errors + /// + /// [`ObjectUrlError`] when the browser rejects URL creation. + pub fn into_object_url(self) -> Result, ObjectUrlError> { + match self { + Self::Local { blob } => ObjectUrl::from_blob(&blob).map(Some), + Self::Remote { .. } | Self::Unavailable => Ok(None), + } + } +} + +/// Why a tab-side stage did not complete. +#[derive(Debug, Error)] +pub enum TabStageError { + /// The blob's bytes could not be read for hashing. + #[error("the staged blob could not be read: {0}")] + Read(String), + /// The content lane refused the announcement. + #[error("the staged content could not be announced: {0}")] + Post(String), + /// The caller's row failed. The announced blob ages out at the worker on + /// its own; nothing needs unwinding. + #[error("the staged row failed: {0}")] + Row(E), +} + +/// One answer to a `Resolve`, with the bytes a `Local` answer carries. +type ResolveAnswer = (WireResolve, Option); + +/// A tab's content lane: stage files through the worker and ask where bytes +/// live, while the tab's client owns the transport itself. +/// +/// Staging pairs with the mutation that names the file. +/// [`stage`](Self::stage) announces the blob before the caller's transaction +/// commits, and message ports deliver in order, so the worker holds the +/// bytes by the time the mutation arrives; the worker then commits the +/// file's manifest, its upload queue entry and the rows as one mutation, +/// checking that the bytes hash to the identity the row names. +pub struct TabContent { + lane: InternalLane, + replies: Rc>>>, + next_request: Rc>, +} + +impl TabContent { + /// Split the content lane off a tab transport, before handing that + /// transport to the tab's client. + #[must_use] + pub fn new(transport: &mut MessageTransport) -> Self { + let lane = transport.internal_lane(); + let replies: Rc>>> = + Rc::new(RefCell::new(HashMap::new())); + if let Some(mut inbox) = transport.take_internal_inbox() { + let waiting = Rc::clone(&replies); + wasm_bindgen_futures::spawn_local(async move { + while let Some(inbound) = inbox.next().await { + if let Some(ContentFrame::ResolveReply { request_id, answer }) = + ContentFrame::from_json(&inbound.json) + && let Some(sender) = waiting.borrow_mut().remove(&request_id) + { + let _ = sender.send((answer, inbound.blob)); + } + } + }); + } + Self { + lane, + replies, + next_request: Rc::new(Cell::new(0)), + } + } + + /// Announce `blob` to the worker and run `row`, so the file, its upload + /// and the rows naming it reach the hub as one mutation. + /// + /// The identity handed to `row` is what the blob's bytes hash to; the + /// worker recomputes it and refuses the mutation if the rows name + /// anything else. A `row` that fails leaves the announced blob to age + /// out at the worker, which costs the bytes and nothing else. + /// + /// # Errors + /// + /// [`TabStageError::Read`] when the blob cannot be hashed, + /// [`TabStageError::Post`] when the lane refuses the announcement, and + /// whatever `row` returns as [`TabStageError::Row`]. + pub async fn stage( + &self, + blob: &Blob, + mime: MimeClass, + client: &ConnettoClient, + row: F, + ) -> Result<(FileId, O), TabStageError> + where + T: Transport + MaybeSend + 'static, + T::Error: Display, + F: FnOnce(&mut SqliteConnection, FileId) -> Result, + E: Display, + { + let mut hasher = FileIdHasher::new(); + let size = blob.size(); + let mut offset: f64 = 0.0; + while offset < size { + let end = (offset + STAGE_READ_BYTES).min(size); + let window = match blob.slice_with_f64_and_f64(offset, end) { + Ok(window) => window, + Err(err) => return Err(TabStageError::Read(format!("{err:?}"))), + }; + let bytes = match blob_bytes(&window).await { + Ok(bytes) => bytes, + Err(err) => return Err(TabStageError::Read(format!("{err:?}"))), + }; + hasher.update(&bytes); + offset = end; + } + let file_id = hasher.finalize(); + let frame = ContentFrame::Stage { + file_id: *file_id.as_bytes(), + mime: mime_code(mime), + }; + self.lane + .post_internal(&frame.to_json(), Some(blob)) + .map_err(|err| TabStageError::Post(err.to_string()))?; + client + .with_conn(|conn| row(conn.conn(), file_id)) + .await + .map(|outcome| (file_id, outcome)) + .map_err(TabStageError::Row) + } + + /// Where this file's bytes are to be had, answered by the worker. + pub async fn resolve(&self, file_id: FileId) -> TabResolved { + let request_id = { + let next = self.next_request.get() + 1; + self.next_request.set(next); + next + }; + let (sender, answer) = oneshot::channel(); + self.replies.borrow_mut().insert(request_id, sender); + let frame = ContentFrame::Resolve { + request_id, + file_id: *file_id.as_bytes(), + }; + if self.lane.post_internal(&frame.to_json(), None).is_err() { + self.replies.borrow_mut().remove(&request_id); + return TabResolved::Unavailable; + } + tokio::select! { + answer = answer => match answer.unwrap_or((WireResolve::Unavailable, None)) { + (WireResolve::Remote { url }, _) => TabResolved::Remote { url }, + (WireResolve::Local, Some(blob)) => TabResolved::Local { blob }, + (WireResolve::Local, None) | (WireResolve::Unavailable, _) => { + TabResolved::Unavailable + } + }, + () = sleep_ms(TAB_RESOLVE_WAIT_MS) => TabResolved::Unavailable, + } + } +} + +/// Read a whole blob through the browser's file reader, which works from any +/// context, page or worker. +async fn blob_bytes(blob: &Blob) -> Result, wasm_bindgen::JsValue> { + let buffer = JsFuture::from(blob.array_buffer()).await?; + Ok(Uint8Array::new(&buffer).to_vec()) +} diff --git a/crates/connetto-web/src/content_wire.rs b/crates/connetto-web/src/content_wire.rs new file mode 100644 index 00000000..f7afece7 --- /dev/null +++ b/crates/connetto-web/src/content_wire.rs @@ -0,0 +1,101 @@ +//! The content protocol that rides a tab transport's internal lane. +//! +//! Staging and resolution questions travel between a tab and the worker hub +//! as JSON frames, never as codec frames, so the sync protocol between worker +//! and server is untouched. Bytes move beside the JSON rather than inside it: +//! a `Stage` and a local `ResolveReply` attach a `Blob` to the message, which +//! structured clone hands over without a copy through Rust. + +use connetto_file_core::MimeClass; +use serde::{Deserialize, Serialize}; + +/// One content-protocol message between a tab and the worker hub. +#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)] +pub enum ContentFrame { + /// The tab announces content it is about to name, declaring the identity + /// it computed, with the blob attached to the same message. + Stage { + /// The identity the tab's bytes hash to. + file_id: [u8; 32], + /// The MIME class the tab chose, as [`mime_code`]. + mime: u8, + }, + /// The tab asks where this file's bytes are to be had. + Resolve { + /// Identifies the reply, unique per tab transport. + request_id: u64, + /// The file being asked about. + file_id: [u8; 32], + }, + /// The hub's answer to a [`ContentFrame::Resolve`]. A `Local` answer + /// carries its bytes as the blob attached to the reply message. + ResolveReply { + /// The `request_id` being answered. + request_id: u64, + /// The answer. + answer: WireResolve, + }, +} + +/// A resolution answer in its wire shape. +#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)] +pub enum WireResolve { + /// The server granted a read for these bytes. + Remote { + /// The granted URL. + url: String, + }, + /// The worker holds the bytes, attached to this message. + Local, + /// Neither the worker nor the server can produce the bytes right now. + Unavailable, +} + +impl ContentFrame { + /// Encode for the internal lane. + pub fn to_json(&self) -> Vec { + // Encoding our own value of a Serialize type cannot fail. + serde_json::to_vec(self).unwrap_or_default() + } + + /// Decode from the internal lane, rejecting anything that is not one of + /// these frames. + pub fn from_json(bytes: &[u8]) -> Option { + serde_json::from_slice(bytes).ok() + } +} + +/// The wire code for a MIME class. +#[must_use] +pub const fn mime_code(mime: MimeClass) -> u8 { + match mime { + MimeClass::Fasta => 0, + MimeClass::Mgf => 1, + MimeClass::Csv => 2, + MimeClass::Json => 3, + MimeClass::Jpeg => 4, + MimeClass::Png => 5, + MimeClass::Video => 6, + MimeClass::Gzip => 7, + MimeClass::Zip => 8, + MimeClass::Generic => 9, + } +} + +/// The class behind a wire code. An unknown code classes as `Generic`, the +/// same answer an unrecognised MIME string gets everywhere else. +#[must_use] +pub const fn mime_from_code(code: u8) -> MimeClass { + match code { + 0 => MimeClass::Fasta, + 1 => MimeClass::Mgf, + 2 => MimeClass::Csv, + 3 => MimeClass::Json, + 4 => MimeClass::Jpeg, + 5 => MimeClass::Png, + 6 => MimeClass::Video, + 7 => MimeClass::Gzip, + 8 => MimeClass::Zip, + _ => MimeClass::Generic, + } +} diff --git a/crates/connetto-web/src/frames.rs b/crates/connetto-web/src/frames.rs index d3d0824b..d0062d3c 100644 --- a/crates/connetto-web/src/frames.rs +++ b/crates/connetto-web/src/frames.rs @@ -9,6 +9,11 @@ //! wire tag byte followed by the `MessagePack` payload. Neither object reports //! the peer going away, so a private close sentinel stands in for a close //! event. +//! +//! Alongside the codec frames rides an internal lane for the content +//! protocol: a frame under the internal tag byte carries JSON rather than +//! `MessagePack`, and a message whose payload is a two-element array carries +//! a `Blob` beside that JSON. The lane never reaches the codec. use connetto_core::codec::{ TAG_BULK, TAG_CONTROL, decode_bulk, decode_control, encode_bulk, encode_control, @@ -28,6 +33,10 @@ use web_sys::MessageEvent; /// shared by every message-delimited transport in this crate. pub(crate) const TAG_CLOSE: u8 = 0xFF; +/// Private wire tag marking the internal content lane. Like the close +/// sentinel it never reaches the codec layer. +pub(crate) const TAG_INTERNAL: u8 = 0xFE; + /// A browser object that carries binary messages one frame at a time. pub trait MessageSink { /// How this sink names itself when it refuses something. @@ -71,6 +80,81 @@ impl MessageTransportError { } } +/// One message from a transport's inbound internal lane. +#[derive(Debug)] +pub struct InternalInbound { + /// The internal frame's JSON, without the tag byte. + pub json: Vec, + /// The bytes attached to the message, when it carried a blob. + pub blob: Option, +} + +/// Decode an internal message posted as `[bytes, blob]`, the shape +/// [`MessageTransport::post_internal`] writes. `None` rejects anything that +/// is not one. +fn decode_attached_internal(data: &JsValue) -> Option { + let parts = data.dyn_ref::()?; + if parts.length() != 2 { + return None; + } + let head = parts.get(0); + let bytes = match head.clone().dyn_into::() { + Ok(view) => view, + Err(_) => Uint8Array::new(head.dyn_ref::()?), + }; + let blob = parts.get(1).dyn_into::().ok()?; + let framed = bytes.to_vec(); + let (tag, json) = framed.split_first()?; + if *tag != TAG_INTERNAL { + return None; + } + Some(InternalInbound { + json: json.to_vec(), + blob: Some(blob), + }) +} + +/// Post one internal message on `sink`: the JSON under the internal tag, +/// with the blob carried beside it rather than inside it. +fn post_internal_to( + sink: &S, + json: &[u8], + blob: Option<&web_sys::Blob>, +) -> Result<(), MessageTransportError> { + let mut framed = Vec::with_capacity(1 + json.len()); + framed.push(TAG_INTERNAL); + framed.extend_from_slice(json); + let message = match blob { + None => Uint8Array::from(framed.as_slice()).into(), + Some(blob) => { + js_sys::Array::of2(&Uint8Array::from(framed.as_slice()), blob.as_ref()).into() + } + }; + sink.post(&message) + .map_err(|err| MessageTransportError::refused::(&err)) +} + +/// A handle that posts internal messages on a transport's lane without +/// owning the transport, for the side that handed the transport over. +pub struct InternalLane { + sink: S, +} + +impl InternalLane { + /// Post one internal message, shaped as [`MessageTransport::post_internal`]. + /// + /// # Errors + /// + /// [`MessageTransportError::Sink`] when the browser refuses the post. + pub fn post_internal( + &self, + json: &[u8], + blob: Option<&web_sys::Blob>, + ) -> Result<(), MessageTransportError> { + post_internal_to(&self.sink, json, blob) + } +} + /// A [`Transport`] over one end of a browser message sink. /// /// The closure stays alive as long as the transport: dropping it would @@ -78,6 +162,7 @@ impl MessageTransportError { pub struct MessageTransport { sink: S, inbound: mpsc::UnboundedReceiver>, + internal: Option>, closed: bool, _on_message: Closure, } @@ -90,11 +175,25 @@ impl MessageTransport { /// peer posted before this call arrive rather than being lost. pub(crate) fn attach(sink: S) -> (Self, mpsc::UnboundedSender>) { let (tx, inbound) = mpsc::unbounded::>(); + let (internal_tx, internal) = mpsc::unbounded::(); let on_message = { let tx = tx.clone(); Closure::::new(move |event: MessageEvent| { - if let Ok(bytes) = event.data().dyn_into::() { - let _ = tx.unbounded_send(bytes.to_vec()); + let data = event.data(); + if let Ok(bytes) = data.clone().dyn_into::() { + let raw = bytes.to_vec(); + if raw.first() == Some(&TAG_INTERNAL) { + let _ = internal_tx.unbounded_send(InternalInbound { + json: raw[1..].to_vec(), + blob: None, + }); + } else { + let _ = tx.unbounded_send(raw); + } + return; + } + if let Some(inbound) = decode_attached_internal(&data) { + let _ = internal_tx.unbounded_send(inbound); } }) }; @@ -103,6 +202,7 @@ impl MessageTransport { Self { sink, inbound, + internal: Some(internal), closed: false, _on_message: on_message, }, @@ -110,6 +210,35 @@ impl MessageTransport { ) } + /// Take the inbound half of this transport's internal lane, once. + pub fn take_internal_inbox(&mut self) -> Option> { + self.internal.take() + } + + /// Hand out a handle that posts internal messages on this lane without + /// owning the transport. + pub fn internal_lane(&self) -> InternalLane + where + S: Clone, + { + InternalLane { + sink: self.sink.clone(), + } + } + + /// Post one internal message: its JSON under the internal tag, with the + /// blob carried beside it rather than inside it. + /// + /// # Errors + /// + /// [`MessageTransportError::Sink`] when the browser refuses the post. + pub fn post_internal( + &self, + json: &[u8], + blob: Option<&web_sys::Blob>, + ) -> Result<(), MessageTransportError> { + post_internal_to(&self.sink, json, blob) + } fn send_frame(&self, tag: u8, payload: &[u8]) -> Result<(), MessageTransportError> { let mut framed = Vec::with_capacity(1 + payload.len()); framed.push(tag); diff --git a/crates/connetto-web/src/lib.rs b/crates/connetto-web/src/lib.rs index 46cbeec3..f44cb5b5 100644 --- a/crates/connetto-web/src/lib.rs +++ b/crates/connetto-web/src/lib.rs @@ -23,6 +23,7 @@ pub mod auth; pub mod broadcast; pub mod content; +pub mod content_wire; pub mod frames; pub mod leader; pub mod locks; @@ -47,9 +48,12 @@ use connetto_core::messages::{BulkMessage, ControlMessage}; use connetto_core::traits::{IncomingFrame, Transport}; pub use content::{ BrowserContentClient, BrowserContentError, BrowserResolved, ObjectUrl, ObjectUrlError, - attach_browser_content, + TabContent, TabResolved, TabStageError, attach_browser_content, +}; +pub use content_wire::{ContentFrame, WireResolve, mime_code, mime_from_code}; +pub use frames::{ + InternalInbound, InternalLane, MessageSink, MessageTransport, MessageTransportError, }; -pub use frames::{MessageSink, MessageTransport, MessageTransportError}; use futures_channel::mpsc; use futures_util::StreamExt; use js_sys::{ArrayBuffer, Uint8Array}; diff --git a/crates/connetto-web/src/relay.rs b/crates/connetto-web/src/relay.rs index 00bae2d9..275a3534 100644 --- a/crates/connetto-web/src/relay.rs +++ b/crates/connetto-web/src/relay.rs @@ -3,7 +3,7 @@ //! One worker-held [`ConnettoConnection`] owns the durable replica and the //! server session, and any number of tabs speak the ordinary connetto wire //! protocol to it over their own [`Transport`]s (an in-memory loopback, or a -//! [`MessageTransport`](crate::MessageTransport) over a `MessageChannel` +//! [`MessageTransport`] over a `MessageChannel` //! port). The hub is a single-task core fed by channels: each attached tab //! gets a shovel task that owns its transport and exchanges frames with the //! core, so the core never selects over a dynamic set of transports and sends @@ -61,6 +61,12 @@ use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use crate::auth::PendingWork; +use crate::content_wire::{ContentFrame, WireResolve, mime_from_code}; +use crate::frames::{ + InternalInbound, InternalLane, MessageSink, MessageTransport, MessageTransportError, +}; +use crate::workers::blob_io::BlobSource; +use crate::workers::helpers::sleep_ms; use connetto_client::reconnect::{ReconnectPolicy, Sleeper, TransportFactory}; use connetto_client::{ AffectedRow, ClientError, ClientEvent, ConnettoConnection, ExportScope, ImportChoices, @@ -76,7 +82,8 @@ use connetto_core::traits::MaybeSend; use connetto_core::{Cursor, IncomingFrame, Transport, quote_ident}; use connetto_file_client::{ BrowserHttp, BrowserStore, ChunkScan, ContentArchive, ContentError, ContentFlush, - ContentFlushStart, ContentFlushState, ContentUpload, FileId, ScanStep, + ContentFlushStart, ContentFlushState, ContentUpload, FileId, PendingConnectionResolve, + ResolveRoute, ResolveStart, Resolved, ScanStep, StageCommitError, }; use diesel::SqliteConnection; use diesel::connection::SimpleConnection; @@ -85,12 +92,23 @@ use diesel::query_builder::{BoxedSqlQuery, SqlQuery}; use diesel::sql_query; use diesel::sqlite::Sqlite; use diesel_sqlite_session::{ConflictAction, SqliteSessionExt}; +use futures_util::StreamExt; use sqlite_diff_rs::{ChangesetOp, ParsedDiffSet, TableSchema, Value}; use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel}; - /// Zstd level for relayed snapshot payloads, matching the client library default. const ZSTD_LEVEL: i32 = 3; +/// Staged content older than this lost its mutation at the hub; the blob is +/// dropped when the tab next stages or writes. +const STALE_CONTENT_MS: f64 = 15_000.0; + +/// Staged blobs a single tab may hold at once, oldest dropped first. +const MAX_STAGED_CONTENT: usize = 8; + +/// Longest a tab's resolve question waits on a ticket round trip before the +/// hub answers `Unavailable`. +const RESOLVE_WAIT_MS: i32 = 15_000; + /// Upstream sequence numbers retained for mapping rejections back to a tab. /// A rejection arrives well within this window, mirroring the client's own /// pending cap. @@ -261,6 +279,9 @@ enum HubEvent { FileId, futures_channel::oneshot::Sender>, ), + /// One content-protocol message from a tab's internal lane. A `Stage` + /// arrives with its blob; the other frames carry none. + Internal(TabId, ContentFrame, Option), } /// One outbound frame toward a tab. Dropping a tab's sender closes it: the @@ -268,6 +289,10 @@ enum HubEvent { enum TabOut { Control(ControlMessage), Bulk(BulkMessage), + /// One content-protocol reply toward a tab, with the blob a `Local` + /// resolution carries. Only transports attached with an internal lane + /// can receive these. + Internal(Vec, Option), } /// A fault while handling one tab's frame: either close that tab, or a @@ -310,6 +335,20 @@ struct TabState { credits: u32, /// Frames queued toward this tab, drained in FIFO order as credits return. pending: VecDeque, + /// Content blobs this tab has staged and not yet paired with a mutation, + /// oldest first. + staged: VecDeque, +} + +/// One staged blob awaiting the mutation that names it. The mutation pairs +/// by the declared identity, and content nothing paired by the time it goes +/// stale is dropped. +struct StagedContent { + file_id: FileId, + mime: connetto_file_core::MimeClass, + blob: web_sys::Blob, + /// `Date::now()` at arrival, milliseconds. + taken: f64, } /// One item waiting on a tab's outbound queue. @@ -468,6 +507,20 @@ struct HubState { /// `SnapshotEnd` triggers the tab re-snapshot, which carries that reason on /// rather than restating one cause as another. resyncing: HashMap, + + /// Resolves waiting on a server ticket answer, each with the tab and + /// request to answer and the `Date::now` instant its wait ends. + pending_resolve: Vec, +} + +/// One resolve whose ticket request is in flight, kept in the hub state so +/// the wait costs a slot rather than the event loop. +struct PendingHubResolve { + tab: TabId, + request_id: u64, + ticket: PendingConnectionResolve, + /// `Date::now` milliseconds by which the tab hears at the latest. + deadline: f64, } /// Handle for attaching tabs to a running hub. Cloneable, and every clone @@ -727,6 +780,31 @@ impl RelayHub { /// Attach one tab transport and spawn its shovel task. pub fn attach(&self, tab: D) -> TabId + where + D: Transport + 'static, + D::Error: core::fmt::Display, + { + self.spawn_tab(tab, None, None) + } + + /// Attach one message transport together with its internal lane, so the + /// tab can stage content and ask resolution questions alongside its + /// sync frames. + pub fn attach_with_content(&self, mut tab: MessageTransport) -> TabId + where + S: MessageSink + Clone + 'static, + { + let internal_rx = tab.take_internal_inbox(); + let replier: Box = Box::new(InternalLaneReplier(tab.internal_lane())); + self.spawn_tab(tab, internal_rx, Some(replier)) + } + + fn spawn_tab( + &self, + tab: D, + internal_rx: Option>, + replier: Option>, + ) -> TabId where D: Transport + 'static, D::Error: core::fmt::Display, @@ -737,7 +815,14 @@ impl RelayHub { // Queued before the shovel exists, so the core learns the tab // before its first frame can possibly arrive on the same channel. let _ = self.events.send(HubEvent::Attached(id, out_tx)); - wasm_bindgen_futures::spawn_local(shovel(id, tab, out_rx, self.events.clone())); + wasm_bindgen_futures::spawn_local(shovel( + id, + tab, + out_rx, + internal_rx, + replier, + self.events.clone(), + )); id } @@ -882,24 +967,77 @@ fn prepare_hub_worker( Ok(local_tables) } -/// The per-tab I/O task: owns the transport, feeds inbound frames to the -/// core, writes outbound frames, and closes the transport when the core -/// drops the tab. +/// The reply leg of an attached tab's internal lane, boxed shape for the +/// shovel so the lane's sink type stays private to the attach call. +trait InternalReplier { + /// Post one internal message to the tab. + /// + /// # Errors + /// + /// [`MessageTransportError::Sink`] when the browser refuses the post. + fn post(&self, json: &[u8], blob: Option<&web_sys::Blob>) -> Result<(), MessageTransportError>; +} + +struct InternalLaneReplier(InternalLane); + +impl InternalReplier for InternalLaneReplier { + fn post(&self, json: &[u8], blob: Option<&web_sys::Blob>) -> Result<(), MessageTransportError> { + self.0.post_internal(json, blob) + } +} + +/// The per-tab I/O task: owns the transport, feeds inbound frames and +/// internal messages to the core, writes outbound frames and internal +/// replies, and closes the transport when the core drops the tab. async fn shovel( id: TabId, mut tab: D, mut out_rx: UnboundedReceiver, + mut internal_rx: Option>, + replier: Option>, events: UnboundedSender, ) where D: Transport, D::Error: core::fmt::Display, { loop { - // Cancel safety: both legs park on an mpsc backed receive, which + // Cancel safety: every leg parks on an mpsc backed receive, which // loses nothing when dropped, and sends on the transports this hub // runs over (loopback and message ports) complete in one poll, so a // losing branch is only ever dropped while parked. + // + // Biased on purpose: the content lane and the codec frames share one + // message port, so a staged blob is posted before the mutation that + // names it, and poll order here is the only thing that keeps that + // order at the hub. A fair select could deliver the mutation first + // and lose the pairing. The handshake needs no such guard: a tab + // stages only after its ack round trip, which the hub already + // answered. tokio::select! { + biased; + inbound = async { + match internal_rx.as_mut() { + Some(rx) => rx.next().await, + None => std::future::pending::>().await, + } + } => match inbound { + Some(inbound) => match ContentFrame::from_json(&inbound.json) { + Some(frame) => { + if events + .send(HubEvent::Internal(id, frame, inbound.blob)) + .is_err() + { + break; + } + } + None => { + tracing::warn!(tab = %id, "relay hub dropped an undecodable internal frame"); + } + }, + // The lane ended with its transport; parking on a closed + // channel would spin until the frame leg catches up. + None => break, + }, frame = tab.recv() => match frame { Ok(Some(frame)) => { if events.send(HubEvent::Frame(id, frame)).is_err() { @@ -919,6 +1057,13 @@ async fn shovel( break; } } + Some(TabOut::Internal(json, blob)) => { + // Only core replies reach this arm, and only tabs + // attached with a lane receive them. + if let Some(replier) = &replier { + let _ = replier.post(&json, blob.as_ref()); + } + } None => { let _ = tab.close().await; break; @@ -1032,6 +1177,8 @@ enum Wake { Upstream(Result), /// The content integrity walk has a file left to check. Verify, + /// A resolve's ticket wait ran out while the hub served everything else. + Resolve, } impl HubRuntime @@ -1082,6 +1229,17 @@ where } }; tokio::pin!(content_wait); + // The resolve waits ride a timer here rather than parking a handler: + // when one fires, the sweep answers only the waits that have truly + // elapsed and the loop keeps serving everything else in between. + let resolve_ms = resolve_deadline_ms(&self.state.pending_resolve); + let resolve_wait = async { + match resolve_ms { + Some(ms) => sleep_ms(ms).await, + None => core::future::pending().await, + } + }; + tokio::pin!(resolve_wait); // Each arm only names its wake reason, so a losing branch leaves // nothing half applied: every mpsc receive loses nothing when dropped. let wake = { @@ -1098,6 +1256,7 @@ where // Always ready, so the walk shares the cycle with every other source // rather than holding it or being starved by it. () = core::future::ready(()), if unverified => Wake::Verify, + () = &mut resolve_wait, if resolve_ms.is_some() => Wake::Resolve, } }; match wake { @@ -1105,6 +1264,10 @@ where Wake::Local(event) => self.serve_local(event).await, Wake::Upstream(event) => self.serve_upstream(reconnect, event).await, Wake::Verify => self.verify_turn().await, + Wake::Resolve => { + expire_resolves(&mut self.state); + Ok(true) + } } } @@ -1313,7 +1476,6 @@ where Ok(true) } - /// Run the content driver now when anything is queued. fn wake_content(&mut self) -> Result<(), RelayError> { if content_sendable_files( &mut self.worker, @@ -1621,7 +1783,10 @@ where match event { HubEvent::Attached(id, out) => attach_tab(state, id, out), HubEvent::Frame(id, frame) => { - handle_tab_frame(worker, state, notices, id, frame).await?; + handle_tab_frame(worker, state, notices, content, walk, id, frame).await?; + } + HubEvent::Internal(id, frame, blob) => { + handle_tab_internal(worker, state, content, id, frame, blob).await?; } HubEvent::Unsynced(reply) => { let pending = PendingWork { @@ -1700,6 +1865,7 @@ fn attach_tab(state: &mut HubState, id: TabId, out: UnboundedSender) { local_watermark: None, credits: INITIAL_CREDITS, pending: VecDeque::new(), + staged: VecDeque::new(), }, ); } @@ -1763,6 +1929,7 @@ where U::Error: core::fmt::Display, { state.tabs.remove(&id); + state.pending_resolve.retain(|resolve| resolve.tab != id); let upstreams: Vec = state .agg_routes .iter() @@ -2195,6 +2362,7 @@ fn recovery_serves_idle(event: &HubEvent) -> bool { | HubEvent::ForgetRetired(_, _) | HubEvent::RefusedContent(_) | HubEvent::RetryRefused(_, _) + | HubEvent::Internal(_, _, _) ) } @@ -2212,6 +2380,7 @@ fn recovery_interrupts_attach(event: &HubEvent) -> bool { | HubEvent::ForgetRetired(_, _) | HubEvent::RefusedContent(_) | HubEvent::RetryRefused(_, _) + | HubEvent::Internal(_, ContentFrame::Resolve { .. }, _) ) } @@ -2240,6 +2409,8 @@ async fn handle_tab_frame( worker: &mut ConnettoConnection, state: &mut HubState, notices: &UnboundedSender, + content: Option<&ContentArchive>, + walk: &mut WalkState, id: TabId, frame: IncomingFrame, ) -> Result<(), RelayError> @@ -2251,13 +2422,14 @@ where IncomingFrame::Control(message) => { handle_tab_control(worker, state, notices, id, message).await } - IncomingFrame::Bulk(bulk) => handle_tab_bulk(worker, state, id, bulk).await, + IncomingFrame::Bulk(bulk) => handle_tab_bulk(worker, state, content, walk, id, bulk).await, }; match outcome { Ok(()) => Ok(()), Err(TabFault::Close(reason)) => { tracing::warn!(tab = %id, reason = %reason, "relay hub closed a tab"); state.tabs.remove(&id); + state.pending_resolve.retain(|resolve| resolve.tab != id); Ok(()) } Err(TabFault::Hub(err)) => Err(err), @@ -2522,6 +2694,8 @@ where async fn handle_tab_bulk( worker: &mut ConnettoConnection, state: &mut HubState, + content: Option<&ContentArchive>, + walk: &mut WalkState, id: TabId, bulk: BulkMessage, ) -> Result<(), TabFault> @@ -2585,7 +2759,15 @@ where &patch.patchset_zstd, ); } - handle_synced_mutation(worker, state, id, tab_seq, &changeset).await + // Content rides a mutation that names it: a changeset carrying a staged + // file's identity commits the manifest, the upload queue entry and the + // rows together. A changeset naming nothing staged is an ordinary + // mutation, and staged content it never names ages out. + let staged = take_staged(state, id, &changeset); + handle_synced_mutation( + worker, state, id, tab_seq, &changeset, staged, content, walk, + ) + .await } /// Bind one changeset value at its own SQLite storage class. @@ -2957,12 +3139,25 @@ where /// mutation. An apply failure rejects the mutation back to the tab and /// leaves the replica untouched, since the abort policy rolls the whole /// apply back. +/// +/// When the changeset paired with staged content, the file is chunked into +/// the encrypted store and its manifest and upload entry commit in the same +/// transaction as the rows, with the declared identity checked against what +/// the bytes actually hash to. A chunking, identity or apply failure all +/// reject the mutation with nothing committed. +#[expect( + clippy::too_many_arguments, + reason = "the mutation needs its tab, its content archive and the walk to wake, none of which belongs on another's struct" +)] async fn handle_synced_mutation( worker: &mut ConnettoConnection, state: &mut HubState, id: TabId, tab_seq: u64, changeset: &[u8], + staged: Option, + content: Option<&ContentArchive>, + walk: &mut WalkState, ) -> Result<(), TabFault> where U: Transport, @@ -2998,35 +3193,37 @@ where let Ok(seq) = i64::try_from(tab_seq) else { return Err(TabFault::Close("sequence overflows storage".to_owned())); }; - let applied = worker.conn().transaction::<_, TabApplyError, _>(|conn| { - conn.apply_changeset(&changeset, |_conflict| ConflictAction::Abort) - .map_err(|err| TabApplyError::Apply(err.to_string()))?; - // sql_query is kept because connetto_hub._tab_mutations is in an - // ATTACHED schema that diesel's table! macro does not model for SQLite. - diesel::sql_query( - "INSERT INTO connetto_hub._tab_mutations (client_id, last_seq) VALUES (?, ?) \ - ON CONFLICT (client_id) DO UPDATE SET \ - last_seq = MAX(last_seq, excluded.last_seq)", - ) - .bind::(client_id) - .bind::(seq) - .execute(conn)?; - Ok(()) - }); - match applied { - Ok(()) => {} - Err(TabApplyError::Apply(detail)) => { - let _ = out.send(TabOut::Control(ControlMessage::MutationReject( - MutationReject { - client_seq: tab_seq, - reason: MutationRejectReason::Other { - detail: format!("worker replica apply failed: {detail}"), - }, - }, - ))); - return Ok(()); + if let Some(staged) = staged { + match commit_staged_mutation(worker, content, staged, &changeset, client_id, seq).await { + Ok(()) => { + // The commit left an outbox entry the driver has not been + // told about, and no upstream event may arrive to tell it. + walk.outbox_wake = true; + } + Err(detail) => { + reject_mutation(&out, tab_seq, detail); + return Ok(()); + } + } + } else { + let applied = worker.conn().transaction::<_, TabApplyError, _>(|conn| { + conn.apply_changeset(&changeset, |_conflict| ConflictAction::Abort) + .map_err(|err| TabApplyError::Apply(err.to_string()))?; + record_tab_watermark(conn, client_id, seq)?; + Ok(()) + }); + match applied { + Ok(()) => {} + Err(TabApplyError::Apply(detail)) => { + reject_mutation( + &out, + tab_seq, + format!("worker replica apply failed: {detail}"), + ); + return Ok(()); + } + Err(TabApplyError::Db(err)) => return Err(RelayError::from(err).into()), } - Err(TabApplyError::Db(err)) => return Err(RelayError::from(err).into()), } if let Some(tab) = state.tabs.get_mut(&id) { tab.applied_watermark = Some(tab_seq); @@ -3040,6 +3237,346 @@ where Ok(()) } +/// Chunk a staged blob and commit its manifest, outbox entry, watermark and +/// row change in one transaction. +/// +/// The error is the refusal detail for the tab, never a hub fault. +async fn commit_staged_mutation( + worker: &mut ConnettoConnection, + content: Option<&ContentArchive>, + staged: StagedContent, + changeset: &[u8], + client_id: rosetta_uuid::Uuid, + seq: i64, +) -> Result<(), String> +where + U: Transport, + U::Error: core::fmt::Display, +{ + let Some(content) = content else { + return Err("this worker holds no content archive".to_owned()); + }; + let Ok(reader) = BlobSource::new(staged.blob) else { + return Err("staged content could not be opened for reading".to_owned()); + }; + // Chunking runs outside the transaction: a large blob must not hold the + // write lock, and bytes chunked for a commit that never lands are orphans + // the next sweep collects, exactly as with the native client's staging. + let manifest = content + .chunk_file(reader, staged.mime) + .await + .map_err(|err| format!("staged content could not be chunked: {err}"))?; + let declared = staged.file_id; + content + .commit_staged(worker, &manifest, |conn, worker_id| { + if worker_id != declared { + return Err(StageCommitError::Row(format!( + "staged bytes hash to {worker_id}, the mutation names {declared}" + ))); + } + conn.apply_changeset(changeset, |_conflict| ConflictAction::Abort) + .map_err(|err| { + StageCommitError::Row(format!("worker replica apply failed: {err}")) + })?; + record_tab_watermark(conn, client_id, seq)?; + Ok(()) + }) + .map_err(|err| match err { + StageCommitError::Row(detail) => detail, + StageCommitError::Bookkeeping(err) => { + format!("the worker could not record the staged file: {err}") + } + }) +} + +/// Answer a mutation with a refusal the tab can surface. +fn reject_mutation(out: &UnboundedSender, tab_seq: u64, detail: String) { + let _ = out.send(TabOut::Control(ControlMessage::MutationReject( + MutationReject { + client_seq: tab_seq, + reason: MutationRejectReason::Other { detail }, + }, + ))); +} + +/// Advance one tab's durable mutation watermark inside the transaction that +/// applied its mutation. +/// +/// `sql_query` is kept because `connetto_hub`._`tab_mutations` is in an ATTACHED +/// schema that diesel's table! macro does not model for SQLite. +fn record_tab_watermark( + conn: &mut SqliteConnection, + client_id: rosetta_uuid::Uuid, + seq: i64, +) -> Result<(), diesel::result::Error> { + diesel::sql_query( + "INSERT INTO connetto_hub._tab_mutations (client_id, last_seq) VALUES (?, ?) \ + ON CONFLICT (client_id) DO UPDATE SET \ + last_seq = MAX(last_seq, excluded.last_seq)", + ) + .bind::(client_id) + .bind::(seq) + .execute(conn)?; + Ok(()) +} + +/// Take the staged blob a synced changeset names, if this tab holds one, +/// dropping entries that went stale waiting for a mutation that never came. +fn take_staged(state: &mut HubState, id: TabId, changeset: &[u8]) -> Option { + let tab = state.tabs.get_mut(&id)?; + if tab.staged.is_empty() { + return None; + } + let now = js_sys::Date::now(); + tab.staged + .retain(|staged| now - staged.taken < STALE_CONTENT_MS); + let named = changeset_blob_values(changeset); + if named.is_empty() { + return None; + } + let index = tab + .staged + .iter() + .position(|staged| named.contains(&staged.file_id))?; + tab.staged.remove(index) +} + +/// The file identities a changeset could be naming: every 32-byte blob value +/// it writes, which is where staged content declares itself. Old images are +/// never scanned: a row's previous identity names content the hub holds +/// already, not content this mutation uploaded. +/// +/// Only changesets are scanned. A tab's own capture produces changesets, and +/// a patchset from a foreign tab is applied without content either way. +fn changeset_blob_values(bytes: &[u8]) -> Vec { + fn named(value: &Value>) -> Option { + let Value::Blob(blob) = value else { + return None; + }; + <[u8; 32]>::try_from(blob.as_slice()) + .ok() + .map(FileId::from_bytes) + } + let Ok(parsed) = ParsedDiffSet::parse(bytes) else { + return Vec::new(); + }; + let mut found = Vec::new(); + if let ParsedDiffSet::Changeset(diff) = parsed { + for op in diff.iter() { + match op { + ChangesetOp::Insert { values, .. } => { + found.extend(values.iter().filter_map(named)); + } + ChangesetOp::Update { values, .. } => { + found.extend(values.iter().filter_map(|pair| named(pair.1.as_ref()?))); + } + ChangesetOp::Delete { .. } => {} + } + } + } + found +} + +/// cannot hold the hub. Both serve wherever they arrive, like `Attached`: +/// they write hub state or read the replica, never the server on the +/// critical path a resolve cannot wait out. +async fn handle_tab_internal( + worker: &mut ConnettoConnection, + state: &mut HubState, + content: Option<&ContentArchive>, + id: TabId, + frame: ContentFrame, + blob: Option, +) -> Result<(), RelayError> +where + U: Transport, + U::Error: core::fmt::Display, +{ + if !state.tabs.get(&id).is_some_and(|tab| tab.handshaken) { + return Ok(()); + } + match frame { + ContentFrame::Stage { file_id, mime } => { + let Some(tab) = state.tabs.get_mut(&id) else { + return Ok(()); + }; + let Some(blob) = blob else { + tracing::warn!(tab = %id, "a stage message carried no blob"); + return Ok(()); + }; + let now = js_sys::Date::now(); + tab.staged + .retain(|staged| now - staged.taken < STALE_CONTENT_MS); + if tab.staged.len() >= MAX_STAGED_CONTENT { + tracing::warn!( + tab = %id, + "the tab already holds the maximum unpaired stages, refusing another; \ + content it announces from here on will not ride its mutation" + ); + return Ok(()); + } + tab.staged.push_back(StagedContent { + file_id: FileId::from_bytes(file_id), + mime: mime_from_code(mime), + blob, + taken: now, + }); + Ok(()) + } + ContentFrame::Resolve { + request_id, + file_id, + } => { + let start = match content { + Some(content) => { + match content + .start_resolve_connection(worker, FileId::from_bytes(file_id)) + .await + { + Ok(start) => Some(start), + Err(err) => { + tracing::warn!( + tab = %id, + ?err, + "the resolve's ticket request could not go out" + ); + None + } + } + } + None => None, + }; + match start { + Some(ResolveStart::Waiting(ticket)) => { + state.pending_resolve.push(PendingHubResolve { + tab: id, + request_id, + ticket, + deadline: js_sys::Date::now() + f64::from(RESOLVE_WAIT_MS), + }); + } + Some(ResolveStart::Answered(Resolved::Remote { url })) => { + answer_resolve(state, id, request_id, WireResolve::Remote { url }, None); + } + Some(ResolveStart::Answered(Resolved::Local { bytes, .. })) => { + answer_resolve(state, id, request_id, WireResolve::Local, Some(bytes)); + } + Some(ResolveStart::Answered(Resolved::Unavailable)) | None => { + answer_resolve(state, id, request_id, WireResolve::Unavailable, None); + } + } + Ok(()) + } + ContentFrame::ResolveReply { .. } => Ok(()), + } +} + +/// Answers a tab's resolve on its content lane. A `Local` answer attaches its +/// bytes as the reply's blob; a tab that has detached loses the answer with +/// the lane. +fn answer_resolve( + state: &mut HubState, + tab: TabId, + request_id: u64, + answer: WireResolve, + bytes: Option>, +) { + let reply = ContentFrame::ResolveReply { request_id, answer }; + let blob = bytes.and_then(|bytes| { + let parts = js_sys::Array::of1(&js_sys::Uint8Array::from(bytes.as_slice())); + match web_sys::Blob::new_with_u8_array_sequence(&parts) { + Ok(blob) => Some(blob), + Err(err) => { + tracing::warn!( + tab = %tab, + error = ?err, + "the browser refused a resolve reply blob" + ); + None + } + } + }); + if let Some(entry) = state.tabs.get(&tab) { + let _ = entry.out.send(TabOut::Internal(reply.to_json(), blob)); + } +} + +/// Runs one upstream event past the resolves waiting on tickets, answering +/// every tab whose wait it settles. The event goes on to the ordinary +/// handling: a grant is news the rest of the hub ignores, a refusal detail +/// names a request no tab write knows, and a closed link is news every +/// handler needs. +fn route_resolves(state: &mut HubState, event: &ClientEvent) { + if state.pending_resolve.is_empty() { + return; + } + let mut settled = Vec::new(); + let mut waiting = Vec::new(); + for resolve in std::mem::take(&mut state.pending_resolve) { + match resolve.ticket.route(event) { + ResolveRoute::Settled(result) => { + settled.push((resolve.tab, resolve.request_id, result)); + } + ResolveRoute::Other => waiting.push(resolve), + } + } + state.pending_resolve = waiting; + for (tab, request_id, result) in settled { + let answer = match result { + Ok(url) => WireResolve::Remote { url }, + Err(err) => { + tracing::warn!(tab = %tab, ?err, "the content ticket request did not reach a grant"); + WireResolve::Unavailable + } + }; + answer_resolve(state, tab, request_id, answer, None); + } +} + +/// Milliseconds until the earliest resolve wait must answer, `None` when no +/// resolve waits. Never zero, so a cycle never busy-spins on an elapsed +/// deadline the sweep has not yet collected. +fn resolve_deadline_ms(pending: &[PendingHubResolve]) -> Option { + let now = js_sys::Date::now(); + // The floor is 1 ms and the ceiling is RESOLVE_WAIT_MS, the widest wait + // this file ever queues, so the value is always in i32 range. + #[expect( + clippy::cast_possible_truncation, + reason = "clamped to [1.0, RESOLVE_WAIT_MS = 15_000]; sub-ms rounding is deliberate" + )] + let ms = pending + .iter() + .map(|resolve| { + (resolve.deadline - now) + .clamp(1.0, f64::from(RESOLVE_WAIT_MS)) + .ceil() as i32 + }) + .min(); + ms +} + +/// Answers `Unavailable` to every resolve whose wait ran out. +fn expire_resolves(state: &mut HubState) { + if state.pending_resolve.is_empty() { + return; + } + let now = js_sys::Date::now(); + let mut settled = Vec::new(); + let mut waiting = Vec::new(); + for resolve in std::mem::take(&mut state.pending_resolve) { + if now >= resolve.deadline { + settled.push((resolve.tab, resolve.request_id)); + } else { + waiting.push(resolve); + } + } + state.pending_resolve = waiting; + for (tab, request_id) in settled { + tracing::warn!(tab = %tab, request_id, "the content ticket did not answer in time"); + answer_resolve(state, tab, request_id, WireResolve::Unavailable, None); + } +} + /// Rewrite a tab's logical-named changeset onto the physical backing tables, /// returning the bytes unchanged when no table was split. fn rename_to_physical(changeset: &[u8], map: &PolicyTables) -> Result, RelayError> { @@ -3092,6 +3629,7 @@ where U: Transport, U::Error: core::fmt::Display, { + route_resolves(state, &event); match event { ClientEvent::LivePatch { cursor, @@ -3750,9 +4288,11 @@ fn session_err(err: E) -> RelayError { #[cfg(test)] mod tests { use super::{ - BrowserHttp, HubContent, HubEvent, TabApplyError, apply_local_changeset, - recovery_interrupts_attach, recovery_serves_idle, schedule_recovery_event, + BrowserHttp, ContentFrame, HubContent, HubEvent, TabApplyError, apply_local_changeset, + changeset_blob_values, recovery_interrupts_attach, recovery_serves_idle, + schedule_recovery_event, }; + use connetto_file_core::FileId; use diesel::connection::SimpleConnection; use diesel::{Connection, RunQueryDsl, SqliteConnection}; use diesel_sqlite_session::SqliteSessionExt; @@ -3880,6 +4420,125 @@ mod tests { assert_eq!(deferred.len(), 2); } + /// A staged blob only writes hub state, and a resolve is answered from + /// the replica or the chunk store, so both are served wherever they + /// arrive. A resolve has a tab waiting on the answer, so it interrupts a + /// replay. A stage waits on nobody and keeps its queue place. + #[wasm_bindgen_test] + fn internal_content_frames_follow_the_recovery_columns() { + let stage = HubEvent::Internal( + 1, + ContentFrame::Stage { + file_id: [1; 32], + mime: 0, + }, + None, + ); + let resolve = HubEvent::Internal( + 1, + ContentFrame::Resolve { + request_id: 7, + file_id: [2; 32], + }, + None, + ); + assert!( + recovery_serves_idle(&stage), + "a stage writes hub state only" + ); + assert!( + recovery_serves_idle(&resolve), + "a resolve never needs the server" + ); + assert!( + recovery_interrupts_attach(&resolve), + "a tab waits on its resolve, so it overtakes a replay" + ); + assert!( + !recovery_interrupts_attach(&stage), + "nothing waits on a stage, so it keeps its place" + ); + } + + /// Pairing reads the declared identity straight out of the tab's + /// changeset: every 32-byte blob an insert writes or an update names in + /// its new values. Old images name nothing, and a delete writes nothing. + #[wasm_bindgen_test] + fn a_changeset_names_its_thirty_two_byte_blobs() { + let mut conn = SqliteConnection::establish(":memory:").expect("open"); + conn.batch_execute( + "CREATE TABLE photos (id INTEGER PRIMARY KEY, content_id BLOB NOT NULL)", + ) + .expect("schema"); + let mut session = conn.create_session().expect("session"); + session.attach_all().expect("attach"); + conn.batch_execute("INSERT INTO photos VALUES (1, x'000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f')") + .expect("insert"); + let changeset = session.changeset().expect("changeset"); + assert_eq!( + changeset_blob_values(&changeset), + vec![FileId::from_bytes([ + 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, + 0x0e, 0x0f, 0x10, 0x11, 0x12, 0x13, 0x14, 0x15, 0x16, 0x17, 0x18, 0x19, 0x1a, 0x1b, + 0x1c, 0x1d, 0x1e, 0x1f + ])], + "an insert names its content" + ); + + conn.batch_execute("UPDATE photos SET content_id = x'202122232425262728292a2b2c2d2e2f303132333435363738393a3b3c3d3e3f' WHERE id = 1") + .expect("update"); + let changeset = session.changeset().expect("changeset"); + let named = changeset_blob_values(&changeset); + assert!( + named.contains(&FileId::from_bytes([ + 0x20, 0x21, 0x22, 0x23, 0x24, 0x25, 0x26, 0x27, 0x28, 0x29, 0x2a, 0x2b, 0x2c, 0x2d, + 0x2e, 0x2f, 0x30, 0x31, 0x32, 0x33, 0x34, 0x35, 0x36, 0x37, 0x38, 0x39, 0x3a, 0x3b, + 0x3c, 0x3d, 0x3e, 0x3f + ])), + "an update names its new content, got {named:?}" + ); + assert!( + !named.contains(&FileId::from_bytes([ + 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, + 0x0e, 0x0f, 0x10, 0x11, 0x12, 0x13, 0x14, 0x15, 0x16, 0x17, 0x18, 0x19, 0x1a, 0x1b, + 0x1c, 0x1d, 0x1e, 0x1f + ])), + "the old image a column sheds names nothing, got {named:?}" + ); + + conn.batch_execute("DELETE FROM photos WHERE id = 1") + .expect("delete"); + let changeset = session.changeset().expect("changeset"); + assert!( + changeset_blob_values(&changeset).is_empty(), + "a delete writes no content and stages nothing" + ); + } + + /// Thirty-one bytes is data, not an identity, and a blob value longer + /// than a hash is nobody's file. + #[wasm_bindgen_test] + fn only_a_thirty_two_byte_blob_names_a_file() { + let mut conn = SqliteConnection::establish(":memory:").expect("open"); + conn.batch_execute( + "CREATE TABLE photos (id INTEGER PRIMARY KEY, content_id BLOB NOT NULL)", + ) + .expect("schema"); + let mut session = conn.create_session().expect("session"); + session.attach_all().expect("attach"); + conn.batch_execute("INSERT INTO photos VALUES (1, x'010203')") + .expect("short blob"); + let changeset = session.changeset().expect("changeset"); + assert!( + changeset_blob_values(&changeset).is_empty(), + "a short blob is not an identity" + ); + assert!( + changeset_blob_values(b"not even a changeset").is_empty(), + "undecodable bytes name nothing" + ); + } + /// Chapter 18's job table: the walk owes a turn through the retry sleep and the /// connect, where the connection is idle. #[wasm_bindgen_test] @@ -4168,6 +4827,7 @@ mod tests { local_watermark: None, credits: super::INITIAL_CREDITS, pending: VecDeque::new(), + staged: VecDeque::new(), }; let mut agg_routes: HashMap = HashMap::new(); let subscribe = Subscribe { diff --git a/crates/connetto-web/src/workers.rs b/crates/connetto-web/src/workers.rs index ffd0cf48..c0d385cb 100644 --- a/crates/connetto-web/src/workers.rs +++ b/crates/connetto-web/src/workers.rs @@ -1,9 +1,9 @@ //! DB worker orchestration and page-side glue for the leader topology. mod archive_channel; -mod blob_io; +pub(crate) mod blob_io; mod boot; -mod helpers; +pub(crate) mod helpers; mod intake; mod logout; mod session; diff --git a/crates/connetto-web/src/workers/helpers.rs b/crates/connetto-web/src/workers/helpers.rs index ee88b134..1f874611 100644 --- a/crates/connetto-web/src/workers/helpers.rs +++ b/crates/connetto-web/src/workers/helpers.rs @@ -14,7 +14,7 @@ pub(super) fn content_store_namespace(seed: &str, replica_db_name: &str) -> Stri } /// Resolve after `ms` milliseconds, in a window or a worker context. -pub(super) async fn sleep_ms(ms: i32) { +pub(crate) async fn sleep_ms(ms: i32) { let promise = Promise::new(&mut |resolve, _reject| { let global = js_sys::global(); let set_timeout = js_sys::Reflect::get(&global, &JsValue::from_str("setTimeout")) diff --git a/crates/connetto-web/src/workers/intake.rs b/crates/connetto-web/src/workers/intake.rs index 1e2e9565..b9c6c772 100644 --- a/crates/connetto-web/src/workers/intake.rs +++ b/crates/connetto-web/src/workers/intake.rs @@ -451,7 +451,7 @@ pub(super) fn install_hello_intake(hub: RelayHub) -> Result<(), IntakeError> { } else if let Some(wire) = message.strip_prefix("tab:") { match MessageTransport::::new(wire) { Ok(transport) => { - hub.attach(transport); + hub.attach_with_content(transport); let _ = hello.post_message(&JsValue::from_str(&format!("attached:{wire}"))); } Err(err) => { diff --git a/crates/connetto-web/tests/content_relay.rs b/crates/connetto-web/tests/content_relay.rs new file mode 100644 index 00000000..4ee06127 --- /dev/null +++ b/crates/connetto-web/tests/content_relay.rs @@ -0,0 +1,660 @@ +//! The tab-to-worker content lane over a real message channel: staged files +//! commit with the mutation that names them, mismatches are refused, and +//! resolution answers from the worker's store, the server ticket, or nowhere. + +use core::convert::Infallible; +use core::future::{Future, ready}; +use std::cell::{Cell, RefCell}; +use std::rc::Rc; + +use connetto_client::reconnect::ReconnectPolicy; +use connetto_client::{ClientConfig, ConnettoConnection, Grant, Replica}; +use connetto_core::Cursor; +use connetto_core::PROTOCOL_VERSION; +use connetto_core::messages::{ + BulkMessage, CONTENT_TICKET_REFUSED, ContentTicketGrant, ControlMessage, Handshake, + HandshakeAck, MutationHeader, MutationPatch, MutationReject, MutationRejectReason, + NonFatalError, +}; +use connetto_core::traits::{IncomingFrame, Transport}; +use connetto_file_client::{BrowserHttp, BrowserStore, ContentArchive}; +use connetto_file_core::{FileId, FileIdHasher, MimeClass}; +use connetto_web::content_wire::{ContentFrame, WireResolve, mime_code}; +use connetto_web::relay::{HubReconnect, RelayHub}; +use connetto_web::workers::DB_ALIVE_LOCK; +use connetto_web::{InternalInbound, MessageTransport, locks}; +use diesel::Connection; +use diesel::connection::SimpleConnection; +use futures_channel::mpsc; +use futures_util::StreamExt; +use js_sys::{Array, Uint8Array}; +use wasm_bindgen::JsCast; +use wasm_bindgen_futures::spawn_local; +use wasm_bindgen_test::{wasm_bindgen_test, wasm_bindgen_test_configure}; +use web_sys::{DedicatedWorkerGlobalScope, MessagePort}; + +#[expect( + dead_code, + reason = "the fixture is shared with the archive suite, whose tests use helpers this one does not call" +)] +#[path = "content_archive/support.rs"] +mod support; + +use support::{install_content_transport, timeout_ms, until}; + +wasm_bindgen_test_configure!(run_in_dedicated_worker); + +const DDL: &str = "CREATE TABLE photos (id INTEGER PRIMARY KEY, content_id BLOB NOT NULL)"; +const TEST_LOCK: &str = "connetto-content-relay-test"; +/// A protocol-shaped id for the scripted tab; the hub stores it as the tab's +/// watermark key. +const TAB_CLIENT_ID: &str = "6f1c9d2e-8a4b-4c5d-9e6f-0a1b2c3d4e5f"; + +/// Distinct bytes per test: a relay a finished test left running keeps driving +/// its own outbox against whatever `fetch` the current test installed. +const COMMIT_PHOTO: &[u8] = b"a photograph whose row and bytes commit as one hub mutation"; +const REFUSE_PHOTO: &[u8] = b"a photograph whose declared identity is a lie"; +const UNPAIRED_PHOTO: &[u8] = b"a staged photograph no mutation ever names"; +const LOCAL_PHOTO: &[u8] = b"a photograph resolved back over the lane before upload"; +const REMOTE_PHOTO: &[u8] = b"a photograph resolved to its server address after upload"; +const HELD_PHOTO: &[u8] = b"a photograph committed while another file's ticket is held"; + +/// How the fake upstream answers content tickets. +#[derive(Clone)] +enum Tickets { + Grants, + Refused, + /// Every request goes unanswered and lands on the shared list, for a + /// test that releases one later. + Held(Rc>>), +} + +struct Upstream { + incoming: mpsc::UnboundedReceiver, + answers: mpsc::UnboundedSender, + mutations: Rc>, + tickets: Tickets, +} + +impl Upstream { + /// A live upstream answering the handshake, counting forwarded mutations, + /// and answering every content ticket per `tickets`. + fn live(mutations: Rc>, tickets: Tickets) -> Self { + let (answers, incoming) = mpsc::unbounded(); + answers + .unbounded_send(IncomingFrame::Control(ControlMessage::HandshakeAck( + HandshakeAck { + connection_id: "r69c-relay".to_owned(), + session_token: "r69c-relay".to_owned(), + resume_token: "r69c-relay".to_owned(), + current_cursor: Cursor::new(Vec::new()), + schema_version: None, + initial_credits: 64, + last_applied_seq: None, + }, + ))) + .expect("queue handshake"); + Self { + incoming, + answers, + mutations, + tickets, + } + } + + /// A sender onto this upstream's inbound queue, for a test that injects + /// an answer of its own after the connection has moved into the hub. + fn answers(&self) -> mpsc::UnboundedSender { + self.answers.clone() + } +} + +impl Transport for Upstream { + type Error = Infallible; + + fn send_control( + &mut self, + message: ControlMessage, + ) -> impl Future> { + if let ControlMessage::ContentTicketRequest(request) = message { + let reply = match &self.tickets { + Tickets::Grants => Some(ControlMessage::ContentTicketGrant(ContentTicketGrant { + request_id: request.request_id.clone(), + url: format!( + "https://content.invalid/files/{}/intent?t=r69c", + FileId::from_bytes(request.file_id) + ), + })), + Tickets::Refused => Some(ControlMessage::NonFatalError(NonFatalError { + related_to: Some(request.request_id.clone()), + detail: CONTENT_TICKET_REFUSED.to_owned(), + })), + Tickets::Held(held) => { + held.borrow_mut().push(request.request_id); + None + } + }; + if let Some(reply) = reply { + self.answers + .unbounded_send(IncomingFrame::Control(reply)) + .expect("queue ticket"); + } + } + ready(Ok(())) + } + + fn send_bulk(&mut self, message: BulkMessage) -> impl Future> { + if matches!(message, BulkMessage::MutationPatch(_)) { + self.mutations.set(self.mutations.get() + 1); + } + ready(Ok(())) + } + + async fn recv(&mut self) -> Result, Self::Error> { + Ok(self.incoming.next().await) + } + + fn close(&mut self) -> impl Future> { + ready(Ok(())) + } +} + +fn hub_config() -> ClientConfig { + ClientConfig::new("r69c-content-relay").with_login(Some(Grant::new("user:r69c"))) +} + +/// A content-aware hub on a live fake upstream plus the tab transport on +/// the other end of a real message channel, already handshaken. +/// +/// Also hands back the mutation counter, and a sender onto the upstream's +/// inbound queue for tests that inject an answer the fake never scripts. +async fn relay_tab( + store_name: &str, + key: [u8; 32], + tickets: Tickets, +) -> ( + RelayHub, + Rc>, + MessageTransport, + mpsc::UnboundedSender, +) { + let mutations = Rc::new(Cell::new(0u32)); + let upstream = Upstream::live(Rc::clone(&mutations), tickets.clone()); + let answers = upstream.answers(); + let worker = ConnettoConnection::::connect( + upstream, + &Replica::in_memory(), + DDL, + &hub_config(), + None, + ) + .await + .expect("connect relay replica"); + let scope = js_sys::global() + .dyn_into::() + .expect("dedicated worker"); + let store = BrowserStore::install(&scope, store_name) + .await + .expect("install store"); + let content = ContentArchive::new(store, key); + let (hub, pump, _notices) = RelayHub::with_reconnect_and_content( + worker, + ":memory:", + HubReconnect { + factory: { + let mutations = Rc::clone(&mutations); + let tickets = tickets.clone(); + move || { + ready(Ok::<_, Infallible>(Upstream::live( + Rc::clone(&mutations), + tickets.clone(), + ))) + } + }, + // A refused ticket retries on this timer. An instant sleeper + // turns that backoff into a busy loop that starves the worker, + // so the fixture sleeps for real, capped to keep tests bounded. + sleeper: |delay: core::time::Duration| async move { + let ms = delay.as_millis().clamp(1, 500); + timeout_ms(i32::try_from(ms).expect("capped above")).await; + }, + policy: ReconnectPolicy::new(), + upstream: Vec::new(), + }, + content, + BrowserHttp::new(), + ) + .expect("content relay"); + spawn_local(async move { + pump.await.expect("relay pump"); + }); + let channel = web_sys::MessageChannel::new().expect("message channel"); + hub.attach_with_content(MessageTransport::::new(channel.port1())); + let mut tab = MessageTransport::::new(channel.port2()); + tab.send_control(ControlMessage::Handshake(Handshake::new( + PROTOCOL_VERSION, + TAB_CLIENT_ID, + ))) + .await + .expect("post handshake"); + let ack = tokio::select! { + frame = tab.recv() => frame.expect("transport").expect("open"), + () = timeout_ms(2_000) => panic!("the hub never answered the handshake"), + }; + assert!( + matches!(ack, IncomingFrame::Control(ControlMessage::HandshakeAck(_))), + "the hub's first answer is the ack, got {ack:?}" + ); + (hub, mutations, tab, answers) +} + +/// The same relay with a scripted grant or refusal for every ticket. +async fn relay_with_tab( + store_name: &str, + key: [u8; 32], + grants: bool, +) -> (RelayHub, Rc>, MessageTransport) { + let (hub, mutations, tab, _answers) = relay_tab( + store_name, + key, + if grants { + Tickets::Grants + } else { + Tickets::Refused + }, + ) + .await; + (hub, mutations, tab) +} + +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") +} + +fn declared_id(bytes: &[u8]) -> FileId { + let mut hasher = FileIdHasher::new(); + hasher.update(bytes); + hasher.finalize() +} + +fn stage(tab: &MessageTransport, file_id: FileId, blob: &web_sys::Blob) { + let frame = ContentFrame::Stage { + file_id: *file_id.as_bytes(), + mime: mime_code(MimeClass::Jpeg), + }; + tab.post_internal(&frame.to_json(), Some(blob)) + .expect("post stage"); +} + +fn resolve(tab: &MessageTransport, request_id: u64, file_id: FileId) { + let frame = ContentFrame::Resolve { + request_id, + file_id: *file_id.as_bytes(), + }; + tab.post_internal(&frame.to_json(), None) + .expect("post resolve"); +} + +/// The changeset of one `photos` insert declaring `content_id`, captured from +/// a throwaway replica the way a tab's own capture would produce it. +fn captured_insert(content_id: &[u8]) -> Vec { + use diesel_sqlite_session::SqliteSessionExt; + let mut hex = String::new(); + for byte in content_id { + std::fmt::write(&mut hex, format_args!("{byte:02x}")).expect("string append"); + } + let mut conn = diesel::SqliteConnection::establish(":memory:").expect("open"); + conn.batch_execute(DDL).expect("schema"); + let mut session = conn.create_session().expect("session"); + session.attach_all().expect("attach"); + conn.batch_execute(&format!("INSERT INTO photos VALUES (1, x'{hex}')")) + .expect("insert"); + session.changeset().expect("changeset") +} + +async fn send_mutation(tab: &mut MessageTransport, seq: u64, changeset: &[u8]) { + tab.send_control(ControlMessage::MutationHeader(MutationHeader { + client_seq: seq, + op_count: 1, + })) + .await + .expect("post header"); + tab.send_bulk(BulkMessage::MutationPatch(MutationPatch { + client_seq: seq, + patchset_zstd: zstd::encode_all(changeset, 3).expect("compress"), + })) + .await + .expect("post patchset"); +} + +async fn recv_control( + tab: &mut MessageTransport, + patience: i32, +) -> Option { + tokio::select! { + frame = tab.recv() => match frame.expect("transport") { + Some(IncomingFrame::Control(message)) => Some(message), + Some(IncomingFrame::Bulk(bulk)) => panic!("unexpected bulk toward the tab: {bulk:?}"), + None => panic!("the hub closed the tab"), + }, + () = timeout_ms(patience) => None, + } +} + +async fn next_reply( + inbox: &mut mpsc::UnboundedReceiver, +) -> (u64, WireResolve, Option) { + let inbound = tokio::select! { + inbound = inbox.next() => inbound.expect("the lane stays open"), + () = timeout_ms(3_000) => panic!("the hub never answered the resolve"), + }; + match ContentFrame::from_json(&inbound.json).expect("decodable reply") { + ContentFrame::ResolveReply { request_id, answer } => (request_id, answer, inbound.blob), + frame => panic!("the lane answered a resolve with {frame:?}"), + } +} + +fn blob_bytes(blob: &web_sys::Blob) -> Vec { + let reader = web_sys::FileReaderSync::new().expect("FileReaderSync in worker"); + let buffer = reader.read_as_array_buffer(blob).expect("read blob"); + Uint8Array::new(&buffer).to_vec() +} + +fn uploads_of(uploaded: &RefCell>>, photo: &[u8]) -> usize { + uploaded + .borrow() + .iter() + .filter(|body| body.as_slice() == photo) + .count() +} + +/// A staged file plus a row naming its id: the hub commits the manifest with +/// the mutation, forwards it once, and the driver uploads the bytes. +#[wasm_bindgen_test] +async fn a_staged_file_and_its_row_reach_the_hub_as_one_mutation() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let uploaded = Rc::new(RefCell::new(Vec::new())); + let stubs = install_content_transport(&uploaded); + let (_hub, mutations, mut tab) = relay_with_tab("r69c-staged-commit", [11; 32], true).await; + let file_id = declared_id(COMMIT_PHOTO); + stage(&tab, file_id, &blob_of(COMMIT_PHOTO)); + send_mutation(&mut tab, 1, &captured_insert(file_id.as_bytes())).await; + assert!( + until(async || uploads_of(&uploaded, COMMIT_PHOTO) == 1).await, + "the staged file must upload" + ); + assert_eq!(mutations.get(), 1, "the mutation reaches the server"); + for _ in 0..4 { + let Some(message) = recv_control(&mut tab, 250).await else { + break; + }; + assert!( + !matches!(message, ControlMessage::MutationReject(_)), + "the staged mutation must be accepted, got {message:?}" + ); + } + drop(stubs); + serial.release(); +} + +#[wasm_bindgen_test] +async fn a_declaration_that_is_not_the_blob_is_refused_and_uploads_nothing() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let uploaded = Rc::new(RefCell::new(Vec::new())); + let stubs = install_content_transport(&uploaded); + let (_hub, mutations, mut tab) = relay_with_tab("r69c-staged-mismatch", [12; 32], true).await; + // The lane declares an identity the bytes do not hash to, and the row + // names that same lie, so pairing succeeds and the recomputed hash fails. + let lie = FileId::from_bytes([9; 32]); + stage(&tab, lie, &blob_of(REFUSE_PHOTO)); + send_mutation(&mut tab, 1, &captured_insert(lie.as_bytes())).await; + let mut reject = None; + for _ in 0..8 { + match recv_control(&mut tab, 1_000).await { + Some(ControlMessage::MutationReject(rejection)) => { + reject = Some(rejection); + break; + } + Some(_) => {} + None => break, + } + } + let detail = match reject { + Some(MutationReject { + reason: MutationRejectReason::Other { detail }, + .. + }) => detail, + other => panic!("the lying declaration must be refused, got {other:?}"), + }; + assert!( + detail.contains("hash to"), + "the refusal names the hash, got {detail}" + ); + assert_eq!( + mutations.get(), + 0, + "a refused mutation never reaches the server" + ); + timeout_ms(300).await; + assert!( + uploaded.borrow().is_empty(), + "a refused commit uploads nothing" + ); + drop(stubs); + serial.release(); +} + +#[wasm_bindgen_test] +async fn content_that_no_mutation_names_uploads_nothing() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let uploaded = Rc::new(RefCell::new(Vec::new())); + let stubs = install_content_transport(&uploaded); + // Tickets are refused, so even a wrongly paired stage cannot upload, and + // the age-out probe below cannot be answered from a server ticket. + let (_hub, mutations, mut tab) = relay_with_tab("r69c-staged-unpaired", [13; 32], false).await; + stage(&tab, declared_id(UNPAIRED_PHOTO), &blob_of(UNPAIRED_PHOTO)); + // The row names a three-byte value, not a file identity, so nothing + // pairs and the mutation goes out as an ordinary one. + send_mutation(&mut tab, 1, &captured_insert(&[0xAA, 0xBB, 0xCC])).await; + assert!( + until(async || mutations.get() == 1).await, + "the unnamed mutation is an ordinary one and must still forward" + ); + for _ in 0..4 { + let Some(message) = recv_control(&mut tab, 250).await else { + break; + }; + assert!( + !matches!(message, ControlMessage::MutationReject(_)), + "an unpaired stage is no reason to refuse the mutation, got {message:?}" + ); + } + timeout_ms(300).await; + assert!( + uploaded.borrow().is_empty() && mutations.get() == 1, + "the mutation forwards as an ordinary one and uploads nothing" + ); + // The stage must have aged out of the hub's books too: an ordinary + // mutation must not have left it resolvable as local content. + let mut inbox = tab.take_internal_inbox().expect("the lane is the tab's"); + resolve(&tab, 6, declared_id(UNPAIRED_PHOTO)); + let (request_id, answer, blob) = next_reply(&mut inbox).await; + assert_eq!(request_id, 6, "the reply answers the age-out probe"); + assert!( + matches!(answer, WireResolve::Unavailable), + "an unpaired stage must leave nothing resolvable, got {answer:?}" + ); + assert!(blob.is_none(), "an Unavailable answer carries no bytes"); + drop(stubs); + serial.release(); +} + +#[wasm_bindgen_test] +async fn an_unsent_staged_file_resolves_to_its_own_bytes() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + // Tickets are refused, so the upload can never complete and the file + // stays unsent, which is exactly the state a local answer comes from. + let (_hub, mutations, mut tab) = relay_with_tab("r69c-resolve-local", [14; 32], false).await; + let mut inbox = tab.take_internal_inbox().expect("the lane is the tab's"); + let file_id = declared_id(LOCAL_PHOTO); + stage(&tab, file_id, &blob_of(LOCAL_PHOTO)); + send_mutation(&mut tab, 1, &captured_insert(file_id.as_bytes())).await; + assert!( + until(async || mutations.get() == 1).await, + "the row commits" + ); + resolve(&tab, 7, file_id); + let (request_id, answer, blob) = next_reply(&mut inbox).await; + assert_eq!(request_id, 7, "the reply carries the request's own id"); + assert!( + matches!(answer, WireResolve::Local), + "an unsent file answers Local, got {answer:?}" + ); + let bytes = blob_bytes(&blob.expect("a Local answer carries the bytes")); + assert_eq!( + bytes, LOCAL_PHOTO, + "the lane carries the exact staged bytes" + ); + serial.release(); +} + +#[wasm_bindgen_test] +async fn an_uploaded_file_resolves_to_a_server_address() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let uploaded = Rc::new(RefCell::new(Vec::new())); + let stubs = install_content_transport(&uploaded); + let (_hub, _mutations, mut tab) = relay_with_tab("r69c-resolve-remote", [15; 32], true).await; + let mut inbox = tab.take_internal_inbox().expect("the lane is the tab's"); + let file_id = declared_id(REMOTE_PHOTO); + stage(&tab, file_id, &blob_of(REMOTE_PHOTO)); + send_mutation(&mut tab, 1, &captured_insert(file_id.as_bytes())).await; + assert!( + until(async || uploads_of(&uploaded, REMOTE_PHOTO) == 1).await, + "the file must finish uploading before the resolve" + ); + // The recorded upload is the request's send, not the dequeuing write + // that follows the server's answer, so the file can still read as + // unsent for a moment here. Retry until it reads sent. + let mut remote = None; + for attempt in 0..30 { + resolve(&tab, 8 + attempt, file_id); + let (reply_id, answer, blob) = next_reply(&mut inbox).await; + assert_eq!(reply_id, 8 + attempt, "the reply answers its own request"); + match answer { + WireResolve::Remote { url } => { + assert!(blob.is_none(), "a Remote answer carries no bytes"); + remote = Some(url); + break; + } + WireResolve::Local => timeout_ms(100).await, + WireResolve::Unavailable => { + panic!("an uploaded granted file must answer Remote, not Unavailable"); + } + } + } + let url = remote.expect("the uploaded file answers Remote within the retries"); + assert!( + url.contains(&file_id.to_string()), + "the ticket names the file, got {url}" + ); + drop(stubs); + serial.release(); +} + +#[wasm_bindgen_test] +async fn a_file_the_hub_never_saw_answers_unavailable() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let (_hub, _mutations, mut tab) = relay_with_tab("r69c-resolve-unknown", [16; 32], false).await; + let mut inbox = tab.take_internal_inbox().expect("the lane is the tab's"); + resolve(&tab, 9, FileId::from_bytes([5; 32])); + let (request_id, answer, blob) = next_reply(&mut inbox).await; + assert_eq!(request_id, 9); + assert!( + matches!(answer, WireResolve::Unavailable), + "no manifest and a refused ticket answer Unavailable, got {answer:?}" + ); + assert!(blob.is_none()); + serial.release(); +} + +/// Proves the hub's event loop is not parked on a ticket wait: a resolve +/// whose ticket never arrives stays unanswered, while mutations commit and +/// other resolves answer, and the grant released at last reaches only the +/// resolve that asked for it. +#[wasm_bindgen_test] +async fn the_hub_answers_around_a_ticket_that_never_arrives() { + let serial = locks::hold_lock(TEST_LOCK).await; + let _alive = locks::hold_lock(DB_ALIVE_LOCK).await; + let held: Rc>> = Rc::new(RefCell::new(Vec::new())); + let (_hub, mutations, mut tab, answers) = relay_tab( + "r69c-resolve-held", + [17; 32], + Tickets::Held(Rc::clone(&held)), + ) + .await; + let mut inbox = tab.take_internal_inbox().expect("the lane is the tab's"); + + let stuck = FileId::from_bytes([42u8; 32]); + resolve(&tab, 11, stuck); + assert!( + until(async || !held.borrow().is_empty()).await, + "the ticket request goes out and waits" + ); + timeout_ms(200).await; + assert!( + matches!( + inbox.try_recv(), + Err(futures_channel::mpsc::TryRecvError::Empty) + ), + "a held ticket must leave the resolve unanswered" + ); + + stage(&tab, declared_id(HELD_PHOTO), &blob_of(HELD_PHOTO)); + let file_id = declared_id(HELD_PHOTO); + send_mutation(&mut tab, 1, &captured_insert(file_id.as_bytes())).await; + assert!( + until(async || mutations.get() == 1).await, + "a mutation commits while another resolve waits on its ticket" + ); + resolve(&tab, 12, file_id); + let (reply_id, answer, _blob) = next_reply(&mut inbox).await; + assert_eq!( + reply_id, 12, + "the second resolve answers while the first waits" + ); + assert!( + matches!(answer, WireResolve::Local), + "an unsent staged file answers Local beside a waiting resolve, got {answer:?}" + ); + + let request_id = held.borrow_mut().remove(0); + answers + .unbounded_send(IncomingFrame::Control(ControlMessage::ContentTicketGrant( + ContentTicketGrant { + request_id, + url: format!("https://content.invalid/files/{stuck}?t=r69c-held"), + }, + ))) + .expect("queue the released grant"); + let (reply_id, answer, blob) = next_reply(&mut inbox).await; + assert_eq!( + reply_id, 11, + "the grant answers the resolve that waited on it" + ); + match answer { + WireResolve::Remote { url } => assert!( + url.contains(&stuck.to_string()), + "the released grant carries its URL, got {url}" + ), + other => panic!("a granted ticket answers Remote, got {other:?}"), + } + assert!(blob.is_none()); + serial.release(); +}