diff --git a/desktop/src-tauri/src/observed_unread.rs b/desktop/src-tauri/src/observed_unread.rs index 45bbbbbcd..acc64ccf8 100644 --- a/desktop/src-tauri/src/observed_unread.rs +++ b/desktop/src-tauri/src/observed_unread.rs @@ -60,6 +60,13 @@ struct IngestEvent { counts_toward_app_badge: bool, } +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +struct ChannelLatestUpdate { + channel_id: String, + created_at: u64, +} + #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] struct MarkerUpdate { @@ -101,6 +108,7 @@ pub(crate) struct IngestRequest { sequence: u64, base_revision: u64, events: Vec, + channel_latest: Vec, markers: Vec, membership: Vec, clear_channels: Vec, @@ -179,6 +187,9 @@ fn open_db(path: &Path) -> Result { counts_badge INTEGER NOT NULL, counts_app_badge INTEGER NOT NULL, PRIMARY KEY(scope,event_id)); CREATE INDEX IF NOT EXISTS observed_events_channel ON observed_events(scope,channel_id,created_at,event_id); + CREATE TABLE IF NOT EXISTS channel_latest( + scope TEXT NOT NULL, channel_id TEXT NOT NULL, created_at INTEGER NOT NULL, + PRIMARY KEY(scope,channel_id)); CREATE TABLE IF NOT EXISTS read_markers( scope TEXT NOT NULL, context_id TEXT NOT NULL, read_at INTEGER NOT NULL, PRIMARY KEY(scope,context_id)); @@ -292,6 +303,30 @@ fn projections(tx: &Transaction<'_>, scope: &str) -> Result>() .map_err(|e| format!("read unread markers: {e}"))?; + let mut by_channel: HashMap = HashMap::new(); + let mut latest_stmt = tx + .prepare("SELECT channel_id,created_at FROM channel_latest WHERE scope=?1") + .map_err(|e| format!("prepare channel latest: {e}"))?; + for row in latest_stmt + .query_map([scope], |r| { + Ok((r.get::<_, String>(0)?, r.get::<_, u64>(1)?)) + }) + .map_err(|e| format!("query channel latest: {e}"))? + { + let (channel_id, latest) = row.map_err(|e| format!("read channel latest: {e}"))?; + by_channel.insert( + channel_id.clone(), + ChannelProjection { + channel_id, + latest, + count: 0, + badge_count: 0, + app_badge_count: 0, + top_level_unread: false, + high_priority_unread: false, + }, + ); + } let mut stmt = tx.prepare("SELECT event_id,channel_id,created_at,root_id,high_priority,counts_badge,counts_app_badge FROM observed_events WHERE scope=?1 ORDER BY channel_id,created_at,event_id").map_err(|e| format!("prepare observed projection: {e}"))?; let rows = stmt .query_map([scope], |r| { @@ -306,7 +341,6 @@ fn projections(tx: &Transaction<'_>, scope: &str) -> Result = HashMap::new(); for row in rows { let (id, channel, created, root, high, badge, app) = row.map_err(|e| format!("read observed projection: {e}"))?; @@ -353,7 +387,7 @@ pub(crate) fn observed_unread_open_scope( .map_err(|e| format!("begin observed-unread open: {e}"))?; let scope = request.scope.key(); ensure_scope(&tx, &scope)?; - let (_, _, _, migration_complete, _) = state(&tx, &scope)?; + let (_, _, _, migration_complete, membership_seeded) = state(&tx, &scope)?; if !migration_complete { if let Some(payload) = &request.legacy_payload { if let Some(channels) = payload @@ -377,8 +411,10 @@ pub(crate) fn observed_unread_open_scope( ) .map_err(|e| format!("mark observed migration: {e}"))?; } - if let Some(seed) = &request.membership_seed { - seed_membership(&tx, &scope, seed)?; + if !membership_seeded { + if let Some(seed) = &request.membership_seed { + seed_membership(&tx, &scope, seed)?; + } } prune(&tx, &scope)?; let channels = projections(&tx, &scope)?; @@ -440,6 +476,8 @@ pub(crate) fn observed_unread_ingest( if request.clear_all { tx.execute("DELETE FROM observed_events WHERE scope=?1", [&scope]) .map_err(|e| format!("clear observed scope: {e}"))?; + tx.execute("DELETE FROM channel_latest WHERE scope=?1", [&scope]) + .map_err(|e| format!("clear channel latest scope: {e}"))?; } for channel in &request.clear_channels { tx.execute( @@ -447,10 +485,22 @@ pub(crate) fn observed_unread_ingest( params![scope, channel], ) .map_err(|e| format!("clear observed channel: {e}"))?; + tx.execute( + "DELETE FROM channel_latest WHERE scope=?1 AND channel_id=?2", + params![scope, channel], + ) + .map_err(|e| format!("clear channel latest: {e}"))?; } for event in &request.events { upsert_event(&tx, &scope, event)?; } + for update in &request.channel_latest { + tx.execute( + "INSERT INTO channel_latest(scope,channel_id,created_at) VALUES(?1,?2,?3) ON CONFLICT(scope,channel_id) DO UPDATE SET created_at=MAX(created_at,excluded.created_at)", + params![scope, update.channel_id, update.created_at], + ) + .map_err(|e| format!("advance channel latest: {e}"))?; + } for update in &request.membership { if update.present { tx.execute( @@ -577,6 +627,51 @@ mod tests { tx.commit().unwrap(); } #[test] + fn latest_anchor_survives_without_a_notify_event_and_seed_is_one_shot() { + let (_d, mut conn) = db(); + let tx = conn.transaction().unwrap(); + let key = scope().key(); + ensure_scope(&tx, &key).unwrap(); + let first = MembershipSeed { + participated_root_ids: vec!["kept".into()], + ..Default::default() + }; + seed_membership(&tx, &key, &first).unwrap(); + let empty = MembershipSeed::default(); + let (_, _, _, _, seeded) = state(&tx, &key).unwrap(); + if !seeded { + seed_membership(&tx, &key, &empty).unwrap(); + } + tx.execute( + "INSERT INTO channel_latest(scope,channel_id,created_at) VALUES(?1,'ch',42)", + [&key], + ) + .unwrap(); + let membership: i64 = tx + .query_row( + "SELECT COUNT(*) FROM unread_membership WHERE scope=?1 AND value='kept'", + [&key], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(membership, 1); + let projected = projections(&tx, &key).unwrap(); + assert_eq!(projected[0].latest, 42); + assert_eq!(projected[0].count, 0); + } + #[test] + fn ingest_request_wire_accepts_channel_latest() { + let request: IngestRequest = serde_json::from_value(serde_json::json!({ + "scope":{"pubkey":"PK","relayUrl":"wss://relay/"}, + "sequence":1,"baseRevision":0,"events":[], + "channelLatest":[{"channelId":"ch","createdAt":42}], + "markers":[],"membership":[],"clearChannels":[],"clearAll":false + })) + .unwrap(); + assert_eq!(request.channel_latest[0].channel_id, "ch"); + assert_eq!(request.channel_latest[0].created_at, 42); + } + #[test] fn serialized_response_matches_typescript_contract() { let actual = serde_json::to_value(ObservedUnreadResponse::Delta { scope: scope(), diff --git a/desktop/src/features/channels/observedUnreadNative.test.mjs b/desktop/src/features/channels/observedUnreadNative.test.mjs new file mode 100644 index 000000000..736ef628f --- /dev/null +++ b/desktop/src/features/channels/observedUnreadNative.test.mjs @@ -0,0 +1,430 @@ +/** + * Native-mode tests for the observed-unread store. + * + * Every other suite in this directory runs with no `window.__TAURI_INTERNALS__`, + * so `invokeTauri` throws and the hook takes the localStorage fallback. That + * makes the whole native protocol untested — the first test here fails if the + * native path is not entered, so the rest cannot silently become tautologies. + */ + +import assert from "node:assert/strict"; +import test from "node:test"; + +import { + installDOMShim, + installFreshStorage, + makeObservedEvent, + mountHook, + mountUnreadChannels, +} from "./observedUnreadTestHarness.mjs"; +import { + installNativeRig, + makeStubRelayClient, +} from "./observedUnreadNativeRig.mjs"; + +installDOMShim(); +installFreshStorage(); + +import { act } from "react"; + +const RELAY = "wss://relay.example.com"; +const NOW_S = Math.floor(Date.now() / 1_000); + +const DEFAULT_PROPS = { + relay: RELAY, + isReady: true, + readStateVersion: 0, + getTs: () => null, + getOwn: () => null, +}; + +function makeRefs() { + return { + eventsRef: { current: new Map() }, + latestRef: { current: new Map() }, + }; +} + +/** Let the hook's promise chain settle (open/ingest are async). */ +async function settle() { + await act(async () => { + await new Promise((resolve) => setTimeout(resolve, 0)); + }); +} + +// ── Entry: the boundary that makes every other test meaningful ──────────────── + +test("native mode is ENTERED: the hook opens the scope over the bridge", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + try { + harness = await mountHook( + { ...DEFAULT_PROPS, pubkey: "pk-entry" }, + makeRefs(), + ); + await settle(); + + assert.equal( + rig.requests("observed_unread_open_scope").length, + 1, + "the hook must call observed_unread_open_scope — if this fails, the suite is measuring the localStorage fallback and every assertion below is vacuous", + ); + assert.equal( + harness.api.isNative(), + true, + "isNative() must be true after a successful open", + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +test("native mode is NOT entered when the bridge fails, and the hook says so", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig({ + failCommands: new Set(["observed_unread_open_scope"]), + }); + try { + harness = await mountHook( + { ...DEFAULT_PROPS, pubkey: "pk-entry-fail" }, + makeRefs(), + ); + await settle(); + + assert.equal( + harness.api.isNative(), + false, + "a failed open must leave the hook on the declared fallback path", + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +// ── D1: local mark-read must reach the native store ────────────────────────── + +test("D1: local markChannelRead sends a read marker to the native store", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + try { + const PUBKEY = "pk-d1"; + const CHANNEL = "channel-d1"; + harness = await mountUnreadChannels({ + pubkey: PUBKEY, + relay: RELAY, + channels: [{ id: CHANNEL, name: "d1", channelType: "stream" }], + relayClient: makeStubRelayClient(), + }); + await settle(); + + const readAt = new Date(NOW_S * 1_000).toISOString(); + await act(async () => { + harness.markChannelRead(CHANNEL, readAt); + }); + await settle(); + + const markers = rig.markerUpdates(); + assert.ok( + markers.some((marker) => marker.contextId === CHANNEL), + `local mark-read must reach observed_unread_ingest as a marker for ${CHANNEL}; saw ${JSON.stringify(markers)}`, + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +test("D1 control: the rig DOES record markers when syncMarkers is called directly", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + try { + harness = await mountHook( + { + ...DEFAULT_PROPS, + pubkey: "pk-d1-control", + getTs: () => NOW_S, + }, + makeRefs(), + ); + await settle(); + + harness.api.syncMarkers(["channel-control"]); + await settle(); + + assert.deepEqual( + rig.markerUpdates(), + [{ contextId: "channel-control", readAt: NOW_S }], + "positive control: the marker path is observable through the rig, so a zero-marker result above means the code did not send one", + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +// ── D2: maxTrigger must survive in native mode ─────────────────────────────── + +test("D2: a catch-up maxTrigger with no notifying event still advances native latest", async () => { + installFreshStorage(); + const CHANNEL = "channel-d2"; + const MAX_TRIGGER = NOW_S - 10; + let harness; + const rig = installNativeRig({ + catchUpChannels: (request) => + request.channels.map((channel) => ({ + status: "success", + channelId: channel.id, + // The regression case: a trigger newer than the read marker that does + // NOT survive the notify filter, so it produces no observed event. + observedEvents: [], + maxTrigger: MAX_TRIGGER, + activityRows: [], + discovered: { participated: [], authored: [], mentioned: [] }, + })), + }); + try { + harness = await mountUnreadChannels({ + pubkey: "pk-d2", + relay: RELAY, + channels: [{ id: CHANNEL, name: "d2", channelType: "stream" }], + relayClient: makeStubRelayClient(), + }); + await settle(); + await settle(); + + assert.equal( + rig.requests("unread_catch_up").length >= 1, + true, + "catch-up must have run for this assertion to mean anything", + ); + assert.equal( + rig + .scope({ pubkey: "pk-d2", relayUrl: RELAY }) + .channelLatest.get(CHANNEL), + MAX_TRIGGER, + `maxTrigger ${MAX_TRIGGER} must survive as the channel latest anchor even when no observed row is returned`, + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +// ── D3: an empty seed must not wipe accumulated native membership ──────────── + +test("D3: reopening with an empty membership seed preserves discovered membership", async () => { + installFreshStorage(); + let first; + let second; + const rig = installNativeRig(); + const scope = { pubkey: "pk-d3", relayUrl: RELAY }; + const emptySeed = { + participatedRootIds: [], + authoredRootIds: [], + mentionedRootIds: [], + followedRootIds: [], + mutedRootIds: [], + mutedChannelIds: [], + }; + try { + first = await mountHook( + { ...DEFAULT_PROPS, pubkey: scope.pubkey, membershipSeed: emptySeed }, + makeRefs(), + ); + await settle(); + + // Catch-up discovery writes membership incrementally, as the commit + // message's ownership story describes. + first.api.updateMembership("participated", "root-discovered", true); + await settle(); + assert.ok( + rig.scope(scope).membership.has("participated\u0000root-discovered"), + "precondition: discovery must have written membership natively", + ); + await first.unmount(); + first = null; + + // Restart with an empty renderer seed (localStorage cleared / read failed). + second = await mountHook( + { ...DEFAULT_PROPS, pubkey: scope.pubkey, membershipSeed: emptySeed }, + makeRefs(), + ); + await settle(); + + assert.ok( + rig.scope(scope).membership.has("participated\u0000root-discovered"), + "an empty renderer seed must not delete membership the native store accumulated", + ); + } finally { + await first?.unmount(); + await second?.unmount(); + rig.restore(); + } +}); + +// ── Matrix rows that only became reachable once native mode was enterable ───── + +test("matrix: a replayed sequence is a no-op, not a second mutation", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + const scope = { pubkey: "pk-replay", relayUrl: RELAY }; + try { + harness = await mountHook( + { ...DEFAULT_PROPS, pubkey: scope.pubkey }, + makeRefs(), + ); + await settle(); + + harness.api.updateMembership("followed", "root-1", true); + await settle(); + const afterFirst = rig.scope(scope).revision; + + const replay = rig.requests("observed_unread_ingest").at(-1); + const response = await globalThis.window.__TAURI_INTERNALS__.invoke( + "observed_unread_ingest", + { request: replay }, + ); + + assert.equal( + response.kind, + "snapshot", + "replay must return a snapshot, not a delta", + ); + assert.equal( + rig.scope(scope).revision, + afterFirst, + "replaying an acked sequence must not advance the revision", + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +test("matrix: a sequence gap is rejected with snapshotRequired", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + const scope = { pubkey: "pk-gap", relayUrl: RELAY }; + try { + harness = await mountHook( + { ...DEFAULT_PROPS, pubkey: scope.pubkey }, + makeRefs(), + ); + await settle(); + + const current = rig.scope(scope); + const response = await globalThis.window.__TAURI_INTERNALS__.invoke( + "observed_unread_ingest", + { + request: { + scope, + sequence: current.lastSequence + 2, + baseRevision: current.revision, + events: [], + markers: [], + membership: [], + clearChannels: [], + clearAll: false, + }, + }, + ); + + assert.equal(response.kind, "snapshotRequired"); + assert.equal( + rig.scope(scope).lastSequence, + current.lastSequence, + "a gap must not advance the ack", + ); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +test("matrix: an ingested event reaches the projection the badge reads", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig(); + const scope = { pubkey: "pk-project", relayUrl: RELAY }; + try { + const refs = makeRefs(); + harness = await mountHook({ ...DEFAULT_PROPS, pubkey: scope.pubkey }, refs); + await settle(); + + harness.api.schedule( + harness.api.currentScope, + "channel-p", + makeObservedEvent({ id: "evt-p", createdAt: NOW_S }), + ); + harness.flushNative?.(); + await act(async () => { + globalThis.dispatchEvent( + new (class extends Event { + constructor() { + super("pagehide"); + } + })(), + ); + }); + await settle(); + + assert.equal( + harness.api.projectionsRef.current.get("channel-p")?.count, + 1, + "the native projection must carry the ingested event", + ); + assert.equal(harness.api.latestForChannel("channel-p"), NOW_S); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); + +test("matrix: a rebuilt store generation is reopened instead of wedging on the old revision", async () => { + installFreshStorage(); + let harness; + const rig = installNativeRig({ + newGeneration: (() => { + let generation = 0; + return () => `gen-${++generation}`; + })(), + }); + const scope = { pubkey: "pk-epoch", relayUrl: RELAY }; + try { + harness = await mountHook( + { ...DEFAULT_PROPS, pubkey: scope.pubkey }, + makeRefs(), + ); + await settle(); + harness.api.updateMembership("followed", "before-rebuild", true); + await settle(); + assert.equal( + rig.scope(scope).revision, + 1, + "precondition: renderer holds revision 1", + ); + + const rebuilt = rig.rebuildScope(scope); + harness.api.updateMembership("followed", "after-rebuild", true); + await settle(); + await settle(); + + assert.equal(rebuilt.generation, "gen-2"); + assert.ok( + rig.requests("observed_unread_open_scope").length >= 2, + "generation mismatch must reopen for a replacement snapshot", + ); + assert.equal(harness.api.isNative(), true); + } finally { + await harness?.unmount(); + rig.restore(); + } +}); diff --git a/desktop/src/features/channels/observedUnreadNativeRig.mjs b/desktop/src/features/channels/observedUnreadNativeRig.mjs new file mode 100644 index 000000000..80b133b4f --- /dev/null +++ b/desktop/src/features/channels/observedUnreadNativeRig.mjs @@ -0,0 +1,380 @@ +/** + * Fake native observed-unread store for tests. + * + * `invokeTauri` calls `window.__TAURI_INTERNALS__.invoke`, which does not exist + * under the Node test shim — so without this rig every mount silently takes the + * localStorage fallback and no test can enter the native path. Install this + * BEFORE mounting to make `observedPersistence.isNative()` true. + * + * The model mirrors `desktop/src-tauri/src/observed_unread.rs` closely enough to + * assert the protocol: one scope row (generation/revision/last_sequence/ + * migration_complete/membership_seeded), observed events, read markers, + * membership, deterministic pruning, and the same projection fold. Keep the two + * in step — a divergence here is a test that certifies the wrong contract. + * + * Exported from a non-test file so the `src/**\/*.test.mjs` glob never picks it + * up as a suite. + */ + +const HORIZON_SECONDS = 7 * 24 * 60 * 60; +const PER_CHANNEL_CAP = 1_000; +const GLOBAL_CAP = 5_000; + +const SEED_KINDS = [ + ["participated", "participatedRootIds"], + ["authored", "authoredRootIds"], + ["mentioned", "mentionedRootIds"], + ["followed", "followedRootIds"], + ["muted_root", "mutedRootIds"], + ["muted_channel", "mutedChannelIds"], +]; + +/** Mirror of `ObservedUnreadScope::key` in observed_unread.rs. */ +function scopeKey(scope) { + return `${scope.pubkey.trim().toLowerCase()}:${scope.relayUrl + .trim() + .replace(/\/+$/, "")}`; +} + +function validLegacyEvent(value, channelId) { + if (typeof value !== "object" || value === null) return null; + const { id, createdAt, rootId, highPriority } = value; + if (typeof id !== "string" || typeof createdAt !== "number") return null; + if (typeof highPriority !== "boolean") return null; + if (typeof value.countsTowardBadge !== "boolean") return null; + if (typeof value.countsTowardAppBadge !== "boolean") return null; + if (rootId !== null && rootId !== undefined && typeof rootId !== "string") + return null; + return { + channelId, + id, + createdAt, + rootId: rootId ?? null, + highPriority, + countsTowardBadge: value.countsTowardBadge, + countsTowardAppBadge: value.countsTowardAppBadge, + }; +} + +class Scope { + constructor(generation) { + this.generation = generation; + this.revision = 0; + this.lastSequence = 0; + this.migrationComplete = false; + this.membershipSeeded = false; + /** Map — `ON CONFLICT(scope,event_id) DO NOTHING`. */ + this.events = new Map(); + /** Map. */ + this.channelLatest = new Map(); + /** Map. */ + this.markers = new Map(); + /** Set<`${kind}\u0000${value}`>. */ + this.membership = new Set(); + } + + prune(nowSeconds) { + const cutoff = nowSeconds - HORIZON_SECONDS; + for (const [id, event] of this.events) { + if (event.createdAt <= cutoff) this.events.delete(id); + } + // Deterministic order, matching `ORDER BY created_at DESC, event_id DESC`. + const newestFirst = (a, b) => + b.createdAt - a.createdAt || (a.id < b.id ? 1 : a.id > b.id ? -1 : 0); + const byChannel = new Map(); + for (const event of this.events.values()) { + const bucket = byChannel.get(event.channelId) ?? []; + bucket.push(event); + byChannel.set(event.channelId, bucket); + } + for (const bucket of byChannel.values()) { + for (const event of bucket.sort(newestFirst).slice(PER_CHANNEL_CAP)) { + this.events.delete(event.id); + } + } + for (const event of [...this.events.values()] + .sort(newestFirst) + .slice(GLOBAL_CAP)) { + this.events.delete(event.id); + } + } + + /** Mirror of `projections()`. */ + projections() { + const marker = (key) => this.markers.get(key) ?? 0; + const byChannel = new Map( + [...this.channelLatest].map(([channelId, latest]) => [ + channelId, + { + channelId, + latest, + count: 0, + badgeCount: 0, + appBadgeCount: 0, + topLevelUnread: false, + highPriorityUnread: false, + }, + ]), + ); + for (const event of this.events.values()) { + let readAt = Math.max(marker(event.channelId), marker(`msg:${event.id}`)); + if (event.rootId) + readAt = Math.max(readAt, marker(`thread:${event.rootId}`)); + if (event.createdAt <= readAt) continue; + const entry = byChannel.get(event.channelId) ?? { + channelId: event.channelId, + latest: 0, + count: 0, + badgeCount: 0, + appBadgeCount: 0, + topLevelUnread: false, + highPriorityUnread: false, + }; + entry.latest = Math.max(entry.latest, event.createdAt); + entry.count += 1; + entry.badgeCount += event.countsTowardBadge ? 1 : 0; + entry.appBadgeCount += event.countsTowardAppBadge ? 1 : 0; + entry.topLevelUnread ||= !event.rootId; + entry.highPriorityUnread ||= event.highPriority; + byChannel.set(event.channelId, entry); + } + return [...byChannel.values()].sort((a, b) => + a.channelId < b.channelId ? -1 : a.channelId > b.channelId ? 1 : 0, + ); + } +} + +/** + * Install the fake native bridge on `window.__TAURI_INTERNALS__`. + * + * Returns a handle for asserting against the store and the recorded IPC calls. + * Call `restore()` in a finally block (or let the next install replace it). + */ +export function installNativeRig(options = {}) { + const { + now = () => Math.floor(Date.now() / 1_000), + newGeneration = () => `gen-${Math.random().toString(16).slice(2)}`, + catchUpChannels = () => [], + failCommands = new Set(), + } = options; + + const scopes = new Map(); + const calls = []; + const previous = globalThis.window?.__TAURI_INTERNALS__; + + const ensureScope = (key) => { + let scope = scopes.get(key); + if (!scope) { + scope = new Scope(newGeneration()); + scopes.set(key, scope); + } + return scope; + }; + + const snapshot = (scope, request) => ({ + kind: "snapshot", + scope: request.scope, + generation: scope.generation, + revision: scope.revision, + lastAckedSequence: scope.lastSequence, + migrationComplete: scope.migrationComplete, + membershipSeeded: scope.membershipSeeded, + channels: scope.projections(), + }); + + const openScope = (request) => { + const scope = ensureScope(scopeKey(request.scope)); + if (!scope.migrationComplete) { + const channels = request.legacyPayload?.eventsByChannel; + if (channels && typeof channels === "object") { + for (const [channelId, events] of Object.entries(channels)) { + if (!Array.isArray(events)) continue; + for (const value of events) { + const event = validLegacyEvent(value, channelId); + if (event && !scope.events.has(event.id)) + scope.events.set(event.id, event); + } + } + } + scope.migrationComplete = true; + } + if (!scope.membershipSeeded && request.membershipSeed) { + // The renderer seed establishes initial ownership once; subsequent opens + // preserve membership accumulated by native catch-up discovery. + scope.membership.clear(); + for (const [kind, field] of SEED_KINDS) { + for (const value of request.membershipSeed[field] ?? []) { + scope.membership.add(`${kind}\u0000${value}`); + } + } + scope.membershipSeeded = true; + } + scope.prune(now()); + return snapshot(scope, request); + }; + + const ingest = (request) => { + const scope = ensureScope(scopeKey(request.scope)); + if (request.sequence <= scope.lastSequence) { + return { + ...snapshot(scope, request), + migrationComplete: true, + membershipSeeded: true, + }; + } + if ( + request.sequence !== scope.lastSequence + 1 || + request.baseRevision !== scope.revision + ) { + return { + kind: "snapshotRequired", + scope: request.scope, + generation: scope.generation, + revision: scope.revision, + lastAckedSequence: scope.lastSequence, + }; + } + const before = new Map( + scope.projections().map((item) => [item.channelId, item]), + ); + if (request.clearAll) { + scope.events.clear(); + scope.channelLatest.clear(); + } + for (const channelId of request.clearChannels) { + scope.channelLatest.delete(channelId); + for (const [id, event] of scope.events) { + if (event.channelId === channelId) scope.events.delete(id); + } + } + for (const latest of request.channelLatest ?? []) { + scope.channelLatest.set( + latest.channelId, + Math.max( + scope.channelLatest.get(latest.channelId) ?? 0, + latest.createdAt, + ), + ); + } + for (const event of request.events) { + if (!scope.events.has(event.id)) { + scope.events.set(event.id, { ...event, rootId: event.rootId ?? null }); + } + } + for (const update of request.membership) { + const key = `${update.kind}\u0000${update.value}`; + if (update.present) scope.membership.add(key); + else scope.membership.delete(key); + } + for (const update of request.markers) { + if (update.readAt === null || update.readAt === undefined) { + scope.markers.delete(update.contextId); + } else { + scope.markers.set( + update.contextId, + Math.max(scope.markers.get(update.contextId) ?? 0, update.readAt), + ); + } + } + scope.prune(now()); + const after = scope.projections(); + const afterIds = new Set(after.map((item) => item.channelId)); + const baseRevision = scope.revision; + scope.revision += 1; + scope.lastSequence = request.sequence; + return { + kind: "delta", + scope: request.scope, + generation: scope.generation, + baseRevision, + revision: scope.revision, + ackedSequence: request.sequence, + upserts: after.filter( + (item) => + JSON.stringify(before.get(item.channelId)) !== JSON.stringify(item), + ), + removed: [...before.keys()].filter((id) => !afterIds.has(id)), + }; + }; + + const handlers = { + observed_unread_open_scope: (args) => openScope(args.request), + observed_unread_ingest: (args) => ingest(args.request), + unread_catch_up: (args) => ({ + channels: catchUpChannels(args.request), + }), + // ReadStateManager reaches the bridge for signing/encryption. Serve inert + // values so a real manager can initialize without a Tauri host. + sign_event: (args) => + JSON.stringify({ + id: `signed-${calls.length}`, + pubkey: "rig-pubkey", + created_at: args.createdAt ?? now(), + kind: args.kind, + tags: args.tags, + content: args.content, + sig: "rig-sig", + }), + nip44_encrypt_to_self: (args) => args.plaintext, + nip44_decrypt_from_self: (args) => args.ciphertext, + }; + + const invoke = async (command, args = {}) => { + calls.push({ command, args }); + if (failCommands.has(command)) { + throw new Error(`rig: ${command} configured to fail`); + } + const handler = handlers[command]; + if (!handler) throw new Error(`rig: unhandled command ${command}`); + return handler(args); + }; + + if (typeof globalThis.window === "undefined") { + Object.defineProperty(globalThis, "window", { + value: globalThis, + configurable: true, + }); + } + globalThis.window.__TAURI_INTERNALS__ = { invoke }; + + return { + calls, + /** Recorded requests for one command, in order. */ + requests: (command) => + calls + .filter((call) => call.command === command) + .map((call) => call.args.request), + /** Every marker update sent to the native store, flattened. */ + markerUpdates: () => + calls + .filter((call) => call.command === "observed_unread_ingest") + .flatMap((call) => call.args.request.markers), + scope: (scope) => scopes.get(scopeKey(scope)), + rebuildScope: (scope) => { + const rebuilt = new Scope(newGeneration()); + rebuilt.migrationComplete = true; + rebuilt.membershipSeeded = true; + scopes.set(scopeKey(scope), rebuilt); + return rebuilt; + }, + restore: () => { + if (previous === undefined) delete globalThis.window.__TAURI_INTERNALS__; + else globalThis.window.__TAURI_INTERNALS__ = previous; + }, + }; +} + +/** + * Minimal RelayClient stand-in so a real ReadStateManager can initialize. + * `useReadState` returns no-op markers unless a relayClient is supplied, and a + * no-op `markContextRead` cannot exercise the local read path at all. + */ +export function makeStubRelayClient() { + return { + fetchEvents: async () => [], + fetchFirstEvent: async () => null, + subscribeLive: async () => async () => {}, + subscribeToReconnects: () => () => {}, + publishEvent: async (event) => event, + }; +} diff --git a/desktop/src/features/channels/observedUnreadTestHarness.mjs b/desktop/src/features/channels/observedUnreadTestHarness.mjs index da3cfa3e9..5ae2a393c 100644 --- a/desktop/src/features/channels/observedUnreadTestHarness.mjs +++ b/desktop/src/features/channels/observedUnreadTestHarness.mjs @@ -256,6 +256,7 @@ export async function mountHook(props, refs) { getTs, getOwn, onPruned, + membershipSeed, }) { apiRef.current = useObservedUnreadPersistence( pubkey, @@ -266,7 +267,7 @@ export async function mountHook(props, refs) { getOwn, refs.eventsRef, refs.latestRef, - { onPruned: onPruned ?? (() => {}) }, + { onPruned: onPruned ?? (() => {}), membershipSeed }, ); return null; } @@ -311,6 +312,8 @@ export function seedStorage(pubkey, relay, channelId, eventId = "evt-1") { export async function mountUnreadChannels({ pubkey, relay = "wss://relay.example.com", + channels = [], + relayClient, }) { const qc = new QueryClient({ defaultOptions: { queries: { retry: false } }, @@ -318,13 +321,15 @@ export async function mountUnreadChannels({ let capturedMarkChannelRead = null; let capturedMarkAllChannelsRead = null; + let capturedResult = null; function Inner({ pubkey: pk }) { - const result = useUnreadChannels([], null, { + const result = useUnreadChannels(channels, null, { pubkey: pk, - relayClient: undefined, + relayClient, relayUrl: relay, }); + capturedResult = result; capturedMarkChannelRead = result.markChannelRead; capturedMarkAllChannelsRead = result.markAllChannelsRead; return null; @@ -350,6 +355,9 @@ export async function mountUnreadChannels({ await render(pubkey); return { + get result() { + return capturedResult; + }, get markChannelRead() { return capturedMarkChannelRead; }, diff --git a/desktop/src/features/channels/useObservedUnreadPersistence.ts b/desktop/src/features/channels/useObservedUnreadPersistence.ts index 6fb309fc6..d27a643b5 100644 --- a/desktop/src/features/channels/useObservedUnreadPersistence.ts +++ b/desktop/src/features/channels/useObservedUnreadPersistence.ts @@ -36,7 +36,11 @@ export type ObservedUnreadPersistence = { ) => void; removeChannel: (channelId: string) => void; updateMembership: (kind: string, value: string, present: boolean) => void; - syncMarkers: (contextIds: Iterable) => void; + syncMarkers: ( + contextIds: Iterable, + explicitReadAt?: ReadonlyMap, + ) => void; + advanceLatest: (channelId: string, createdAt: number) => void; latestForChannel: (channelId: string) => number | undefined; clearAll: () => void; }; @@ -201,6 +205,7 @@ export function useObservedUnreadPersistence( sequence: current.sequence + 1, baseRevision: current.revision, events, + channelLatest: [], markers: [], membership: [], clearChannels: [], @@ -310,15 +315,19 @@ export function useObservedUnreadPersistence( }, [normalizedPubkey, normalizedRelayUrl]); const syncMarkers = React.useCallback( - (contextIds: Iterable) => { + ( + contextIds: Iterable, + explicitReadAt?: ReadonlyMap, + ) => { const state = nativeRef.current; if (!state) return; const markers = [...new Set(contextIds)].map((contextId) => ({ contextId, readAt: - contextId.startsWith("thread:") || contextId.startsWith("msg:") + explicitReadAt?.get(contextId) ?? + (contextId.startsWith("thread:") || contextId.startsWith("msg:") ? getOwnTimestamp(contextId) - : getEffectiveTimestamp(contextId), + : getEffectiveTimestamp(contextId)), })); if (markers.length === 0) return; flushNative(); @@ -335,6 +344,7 @@ export function useObservedUnreadPersistence( sequence: current.sequence + 1, baseRevision: current.revision, events: [], + channelLatest: [], markers, membership: [], clearChannels: [], @@ -399,6 +409,7 @@ export function useObservedUnreadPersistence( sequence: current.sequence + 1, baseRevision: current.revision, events: [], + channelLatest: [], markers: [], membership: [], clearChannels, @@ -410,6 +421,41 @@ export function useObservedUnreadPersistence( [apply, flushNative], ); + // biome-ignore lint/correctness/useExhaustiveDependencies: mutable storage refs are stable containers + const advanceLatest = React.useCallback( + (channelId: string, createdAt: number) => { + const state = nativeRef.current; + if (!state) { + const current = latestByChannelRef.current.get(channelId) ?? 0; + if (createdAt > current) + latestByChannelRef.current.set(channelId, createdAt); + return; + } + flushNative(); + chainRef.current = chainRef.current.then(() => { + const current = nativeRef.current; + if ( + !current || + current.scope.pubkey !== state.scope.pubkey || + current.scope.relayUrl !== state.scope.relayUrl + ) + return; + return ingestObservedUnread({ + scope: current.scope, + sequence: current.sequence + 1, + baseRevision: current.revision, + events: [], + channelLatest: [{ channelId, createdAt }], + markers: [], + membership: [], + clearChannels: [], + clearAll: false, + }).then(apply); + }); + }, + [apply, flushNative], + ); + const updateMembership = React.useCallback( (kind: string, value: string, present: boolean) => { const state = nativeRef.current; @@ -428,6 +474,7 @@ export function useObservedUnreadPersistence( sequence: current.sequence + 1, baseRevision: current.revision, events: [], + channelLatest: [], markers: [], membership: [{ kind, value, present }], clearChannels: [], @@ -488,6 +535,7 @@ export function useObservedUnreadPersistence( removeChannel, updateMembership, syncMarkers, + advanceLatest, latestForChannel, clearAll, }), @@ -499,6 +547,7 @@ export function useObservedUnreadPersistence( removeChannel, updateMembership, syncMarkers, + advanceLatest, latestForChannel, clearAll, ], diff --git a/desktop/src/features/channels/useUnreadChannels.ts b/desktop/src/features/channels/useUnreadChannels.ts index 40778305b..b956b672d 100644 --- a/desktop/src/features/channels/useUnreadChannels.ts +++ b/desktop/src/features/channels/useUnreadChannels.ts @@ -249,7 +249,6 @@ export function useUnreadChannels( mutedChannelIdsRef.current, ); - // Persistence layer: hydration, pagehide flush, scope fence, write-through, marker-prune. const observedPersistence = useObservedUnreadPersistence( normalizedPubkey, normalizedRelayUrl, @@ -265,7 +264,6 @@ export function useUnreadChannels( }, ); - // Forward only changed NIP-RS contexts; native owns the event read model. // biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion is the intentional drain trigger React.useEffect(() => { const advanced = drainSyncedAdvances(); @@ -350,11 +348,13 @@ export function useUnreadChannels( ); if (markAt === null) return; markContextRead(channelId, markAt); - // Delegate destructive observed-ref removal to the fenced owner operation — + observedPersistence.syncMarkers( + [channelId], + new Map([[channelId, markAt]]), + ); // the parent must not delete from latestByChannelRef or // observedUnreadEventsByChannelRef directly on the clear-observed path, // or a stale scope-A callback could corrupt scope B before the fence rejects. - // (Fenced record writes in handleChannelMessage and catch-up remain in the parent.) if (clearObserved) { observedPersistence.removeChannel(channelId); bumpLatestVersion(); @@ -721,11 +721,8 @@ export function useUnreadChannels( ) { const current = observedPersistence.latestForChannel(result.channelId) ?? 0; - if ( - !observedPersistence.isNative() && - result.maxTrigger > current - ) { - latestByChannelRef.current.set( + if (result.maxTrigger > current) { + observedPersistence.advanceLatest( result.channelId, result.maxTrigger, ); @@ -923,6 +920,7 @@ export function useUnreadChannels( unreadChannelIdsRef.current = unreadChannelIds; const markAllChannelsRead = React.useCallback(() => { + const marked = new Map(); for (const channelId of unreadChannelIdsRef.current) { delete forcedUnreadRef.current[channelId]; const unixSeconds = @@ -931,12 +929,13 @@ export function useUnreadChannels( null; if (unixSeconds !== null) { markContextRead(channelId, unixSeconds); + marked.set(channelId, unixSeconds); } } + observedPersistence.syncMarkers(marked.keys(), marked); if (pubkey) { forcedUnreadStore.write(pubkey, forcedUnreadRef.current); } - // Delegate destructive observed-ref clearing to the fenced owner operation — // the parent must not reset the observed Maps directly on this path, or a // stale scope-A callback could corrupt scope B before the fence rejects. // (Fenced record writes in handleChannelMessage and catch-up remain in the parent.) diff --git a/desktop/src/shared/api/tauriObservedUnread.ts b/desktop/src/shared/api/tauriObservedUnread.ts index 7a02ece72..73999aeae 100644 --- a/desktop/src/shared/api/tauriObservedUnread.ts +++ b/desktop/src/shared/api/tauriObservedUnread.ts @@ -64,6 +64,7 @@ export function ingestObservedUnread(request: { sequence: number; baseRevision: number; events: ObservedUnreadWireEvent[]; + channelLatest: Array<{ channelId: string; createdAt: number }>; markers: Array<{ contextId: string; readAt: number | null }>; membership: Array<{ kind: string; value: string; present: boolean }>; clearChannels: string[];