fix(desktop): cover and repair native unread mode

Exercise the native IPC boundary with an in-memory protocol rig, route local read markers with their exact timestamps, persist catch-up latest anchors independently of notify rows, and seed renderer membership only once.

Co-authored-by: Tyler Longwell <tlongwell@squareup.com>
Signed-off-by: Tyler Longwell <tlongwell@squareup.com>
This commit is contained in:
Perci
2026-08-16 04:57:08 -04:00
co-authored by Tyler Longwell
parent a77c6c0679
commit d1901c1c8a
7 changed files with 983 additions and 21 deletions
+99 -4
View File
@@ -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<IngestEvent>,
channel_latest: Vec<ChannelLatestUpdate>,
markers: Vec<MarkerUpdate>,
membership: Vec<MembershipUpdate>,
clear_channels: Vec<String>,
@@ -179,6 +187,9 @@ fn open_db(path: &Path) -> Result<Connection, String> {
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<Vec<ChannelProjectio
.map_err(|e| format!("query unread markers: {e}"))?
.collect::<Result<_, _>>()
.map_err(|e| format!("read unread markers: {e}"))?;
let mut by_channel: HashMap<String, ChannelProjection> = 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<Vec<ChannelProjectio
))
})
.map_err(|e| format!("query observed projection: {e}"))?;
let mut by_channel: HashMap<String, ChannelProjection> = 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(),
@@ -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();
}
});
@@ -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<eventId, event> — `ON CONFLICT(scope,event_id) DO NOTHING`. */
this.events = new Map();
/** Map<channelId, latest catch-up trigger>. */
this.channelLatest = new Map();
/** Map<contextId, readAt>. */
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,
};
}
@@ -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;
},
@@ -36,7 +36,11 @@ export type ObservedUnreadPersistence = {
) => void;
removeChannel: (channelId: string) => void;
updateMembership: (kind: string, value: string, present: boolean) => void;
syncMarkers: (contextIds: Iterable<string>) => void;
syncMarkers: (
contextIds: Iterable<string>,
explicitReadAt?: ReadonlyMap<string, number>,
) => 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<string>) => {
(
contextIds: Iterable<string>,
explicitReadAt?: ReadonlyMap<string, number>,
) => {
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,
],
@@ -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<string, number>();
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.)
@@ -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[];