From 21931c38d41d1cd05806d12b61e3e999cd61d62d Mon Sep 17 00:00:00 2001 From: randogoth Date: Tue, 29 Sep 2026 00:08:03 +0300 Subject: [PATCH] feat: check fetch cancellation before each envelope --- README.md | 2 +- core/src/client.rs | 28 +++++++++++++++++++++++----- 2 files changed, 24 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index 3651e4e..ee40e36 100644 --- a/README.md +++ b/README.md @@ -76,7 +76,7 @@ The crate is a workspace: `core/` is the `fumi-core` library — everything but - **A storeless build.** `default-features = false` drops rusqlite entirely: identities, addresses, seal/open and the raw transport operations remain, for hosts with their own database. `store` and `bundled-sqlite` compose with it as separate features. - **One shareable store handle.** `Store` is `Send + Sync`: every operation locks an inner mutex, and the database runs in WAL mode, so a host holds one `Store` behind an executor and reads it while a fetch writes. Each envelope is committed as it is verified — a concurrent reader listing the inbox sees fetch progress without any callback. `Store::open_in_memory()` gives a scratch store with the schema and version already in place — for tests, foreign importers assembling rows in code, and FFI round-trips that never touch disk. - **Decisions, not prose.** `error::Error` is an enum: `NotPinned`, `PinMismatch`, `KeyChanged { address, known, offered }`, `NotRegistered`, `AuthFailed`, `RateLimited`, `QuotaExceeded`, `SchemaVersion`, ... — the distinctions a UI routes on. `Display` still produces the message the CLI prints. -- **Long operations can stop.** `fetch_with` takes a `FetchOptions` with a cancellation token checked between pages; the cursor is saved per page, so a cancelled fetch resumes rather than repeats. +- **Long operations can stop.** `fetch_with` takes a `FetchOptions` with a cancellation token checked between pages and before each envelope; the cursor covers exactly what was processed, so a cancelled fetch resumes rather than repeats. - **Composable flows.** `restore` remains the one-call convenience, but its pieces are public: `connect` + `resolve`, `account::find_rotation_index`, then the `Store` setters — a wizard can hold the master and the address before it has connectivity, and show each step. Trust outcomes arrive as values: `TrustChange` from `trust_key` and `send`, `Session::unpinned_static` for the sec 4 warning. - **A versioned store.** The schema carries a `user_version`; a store written by a newer build is refused (`Error::SchemaVersion`) rather than misread, and within a version the schema is stable. The host, not the library, decides when to upgrade the binary. - **Binding from other languages.** The API is `Send`-friendly and free of CLI globals, so plain Rust bindings (flutter_rust_bridge, uniffi) work against it. `core/vectors.json` holds the reference vectors the unit tests pin — derivations, seals, a full envelope — so a port or FFI binding can verify itself end to end without the reference stack. diff --git a/core/src/client.rs b/core/src/client.rs index c659ca5..f8a18c8 100644 --- a/core/src/client.rs +++ b/core/src/client.rs @@ -535,8 +535,8 @@ pub struct Fetched { /// cursor (the default). pub acknowledged: bool, /// `cancel` was raised mid-fetch; the summary covers what completed. - /// Every finished page is accounted for locally — cursor moved or - /// messages deleted — so fetching again continues rather than repeats. + /// Every envelope the summary counted is accounted for locally, so + /// fetching again continues rather than repeats. pub cancelled: bool, } @@ -551,8 +551,10 @@ pub struct FetchOptions<'a> { pub reset: bool, /// Network timeout in seconds. pub timeout: u64, - /// Checked between pages; when raised, `fetch_with` stops and returns - /// the partial summary with `cancelled` set. + /// Checked between pages and before each envelope; when raised, + /// `fetch_with` stops and returns the partial summary with `cancelled` + /// set. What was processed is accounted for locally — cursor moved or + /// messages deleted — so fetching again continues rather than repeats. pub cancel: Option<&'a std::sync::atomic::AtomicBool>, } @@ -632,7 +634,16 @@ pub fn fetch_with( break; } let mut acked: Vec<[u8; ID_LEN]> = Vec::new(); + let mut page_cancelled = false; for _ in 0..count { + // Checked before each envelope as well as between pages, so a + // cancel lands within one envelope; the cursor and acks below + // then cover exactly the envelopes that were processed. + if cancelled() { + page_cancelled = true; + summary.cancelled = true; + break; + } let mid: [u8; ID_LEN] = r.take(ID_LEN)?.try_into().unwrap(); let received_at = r.i64()?; let flags = r.u8()?; @@ -662,12 +673,19 @@ pub fn fetch_with( } acked.push(mid); } - r.done()?; + // A cancelled page is not fully parsed, so trailing bytes are + // expected and not an error; everything processed is accounted for. + if !page_cancelled { + r.done()?; + } if opts.keep { store.set_cursor(after_time, &after_id)?; } else if !acked.is_empty() { delete_ids(transport, &acked)?; } + if page_cancelled { + break; + } } session.close(); if !opts.keep {