perf(desktop): move observed unread state into native SQLite

Persist observed unread events, read markers, membership, and channel
projections in a scoped WAL database. Keep renderer mutations ordered with a
sequence/revision protocol, migrate localStorage transactionally, and let
native catch-up load membership without serializing five capped arrays.

Measured at the same populated 5x1000 membership fixture:
- before: 335,280 bytes (327.4 KiB)
- after: 178 bytes (0.2 KiB)

Ack/restart failure matrix:

| Failure boundary | Contract |
|---|---|
| Rust commits sequence N, renderer dies before observing ack | **Handled by replay:** DB ack is durable; reopen returns `lastAckedSequence=N`. Renderer may resend N from its pending local queue; Rust recognizes `N <= ack`, performs no mutation, and returns current revision/snapshot. Event-id idempotence is the second fence. |
| Renderer observes ack N, dies before deleting its pending batch | **Handled identically:** replay N is a no-op; no duplicate row or revision bump. |
| Renderer deletes pending N without durable ack | **Unrepresentable by construction:** deletion happens only in the resolved-success branch after validating matching scope, sequence, and revision. Reject/throw leaves N queued. |
| Scope switch lands while A ingest is running | **Handled:** A request carries immutable scope; transaction commits only to A. Switch drains A first when available, opens B independently, and response merge requires exact current scope; a late A ack cannot advance B's sequence or projection. If drain fails, A's retained unacked batch stays replayable when A reopens. |
| Scope switch after A commit but before A ack observation | **Handled by durable ack + scope fence:** A reopen learns ack N; B never sees it. |
| Batch N+1 arrives before N / IPC retry reorders | **Explicitly rejected:** only `sequence == ack+1` mutates. `> ack+1` returns `snapshotRequired`/expected sequence; `<= ack` is replay/no-op. Single JS queue makes normal reordering unrepresentable, backend check covers abnormal callers/restarts. |
| DB commit succeeds but response serialization/delivery fails | **Handled as lost ack:** resend; durable sequence and idempotent event ids collapse it. |
| DB transaction fails halfway (events/markers/prune/revision/ack) | **Unrepresentable:** one SQLite transaction; rollback leaves ack+revision unchanged, so retry is the same next sequence. |
| Migration imports rows but renderer dies before seeing marker | **Handled:** rows + migration marker commit atomically. Reopen reports complete and current snapshot; legacy key remains until renderer observes that, then is deleted. Re-sending payload after complete is ignored. |
| localStorage delete succeeds, native DB later becomes unavailable | **Handled by one-release fallback limitation explicitly:** fallback can preserve new session events but cannot reconstruct migrated history after confirmed native ownership. Native open failure surfaces and does not mutate/delete legacy data pre-confirmation; DB corruption after confirmed migration is logged/recoverable as degraded state, not silently represented as “zero unread.” I will test this distinction rather than claim impossible loss recovery. |
| Revision response gap/out-of-order | **Handled:** apply requires exact `baseRevision`; otherwise discard payload and request full snapshot. Snapshot replacement requires matching scope and `revision >= current`. |
| App shutdown during coalesce | **Handled twice:** `pagehide` drains renderer queue; native shutdown flushes SQLite/checkpoint. Already-committed batches need no renderer ack to survive restart. Unsent events inside a renderer killed without pagehide are the irreducible initial-source limit; live relay catch-up can rediscover them, and the existing TS fallback gate remains for native write failures—not arbitrary process kill. |

Build extension: each scope carries a generated UUID epoch. A different epoch
means the store was rebuilt, so the renderer accepts the replacement snapshot
regardless of its old revision and resets sequence/revision. This makes DB
recreation distinguishable from a stale snapshot and prevents a revision wedge.
Native projection deltas now contain changed channels only. Top-level unread
continues to mean an unread observed event whose `rootId` is null, matching the
retired renderer helper exactly. Read-state version changes, including relay
sync advances, feed marker updates into the same ordered mutation chain.

Co-authored-by: Perci <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz>
Signed-off-by: Perci <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz>
This commit is contained in:
Perci
2026-08-16 04:02:54 -04:00
parent 38afea0a0a
commit 8a392d87aa
9 changed files with 1269 additions and 214 deletions
+4
View File
@@ -31,6 +31,7 @@ mod native_websocket;
mod native_websocket_batch;
mod nostr_bind;
pub mod nostr_convert;
mod observed_unread;
mod persona_catalog;
mod prevent_sleep;
mod ptt_shortcut;
@@ -232,6 +233,7 @@ pub fn run() {
.manage(terminal_runtime::TerminalSessions::default())
.manage(archive::sync::ArchiveSyncState::default())
.manage(native_relay_client::NativeRelayClient::default())
.manage(observed_unread::ObservedUnreadStore::default())
.setup(move |app| {
let app_handle = app.handle().clone();
#[cfg(target_os = "macos")]
@@ -718,6 +720,8 @@ pub fn run() {
probe_backend_provider,
persona_catalog::fetch_persona_catalog,
unread_catch_up::unread_catch_up,
observed_unread::observed_unread_open_scope,
observed_unread::observed_unread_ingest,
list_personas,
create_persona,
update_persona,
+602
View File
@@ -0,0 +1,602 @@
//! Native observed-unread read model.
//!
//! The renderer is the only writer today, so request/response ordering is the
//! delivery mechanism: there is no push channel. If native relay ingestion adds
//! a second writer, that assumption breaks; consumers must then use the same
//! revision-gap rule here to request a fresh snapshot.
//!
//! Failure contract: sequence + revision advance in the same SQLite transaction
//! as events, markers, pruning, and migration. A lost ack is replayed as a no-op;
//! a gap is rejected; stale-scope responses are fenced in the renderer. Legacy
//! rows and their migration marker commit together, and localStorage is removed
//! only after the renderer observes that marker.
use std::{
collections::{HashMap, HashSet},
path::{Path, PathBuf},
sync::Mutex,
};
use rusqlite::{params, Connection, Transaction};
use serde::{Deserialize, Serialize};
use tauri::{AppHandle, Manager, State};
const SCHEMA_VERSION: i64 = 1;
const PER_CHANNEL_CAP: i64 = 1_000;
const GLOBAL_CAP: i64 = 5_000;
const HORIZON_SECONDS: i64 = 7 * 24 * 60 * 60;
#[derive(Default)]
pub(crate) struct ObservedUnreadStore {
write_lock: Mutex<()>,
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub(crate) struct ObservedUnreadScope {
pub(crate) pubkey: String,
pub(crate) relay_url: String,
}
impl ObservedUnreadScope {
fn key(&self) -> String {
format!(
"{}:{}",
self.pubkey.trim().to_ascii_lowercase(),
self.relay_url.trim().trim_end_matches('/')
)
}
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct IngestEvent {
channel_id: String,
id: String,
created_at: u64,
root_id: Option<String>,
high_priority: bool,
counts_toward_badge: bool,
counts_toward_app_badge: bool,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct MarkerUpdate {
context_id: String,
read_at: Option<u64>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct MembershipUpdate {
kind: String,
value: String,
present: bool,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct MembershipSeed {
participated_root_ids: Vec<String>,
authored_root_ids: Vec<String>,
mentioned_root_ids: Vec<String>,
followed_root_ids: Vec<String>,
muted_root_ids: Vec<String>,
muted_channel_ids: Vec<String>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct OpenScopeRequest {
scope: ObservedUnreadScope,
legacy_payload: Option<serde_json::Value>,
membership_seed: Option<MembershipSeed>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct IngestRequest {
scope: ObservedUnreadScope,
sequence: u64,
base_revision: u64,
events: Vec<IngestEvent>,
markers: Vec<MarkerUpdate>,
membership: Vec<MembershipUpdate>,
clear_channels: Vec<String>,
clear_all: bool,
}
#[derive(Clone, Debug, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct ChannelProjection {
channel_id: String,
latest: u64,
count: u64,
badge_count: u64,
app_badge_count: u64,
top_level_unread: bool,
high_priority_unread: bool,
}
#[derive(Debug, Serialize)]
#[serde(
tag = "kind",
rename_all = "camelCase",
rename_all_fields = "camelCase"
)]
pub(crate) enum ObservedUnreadResponse {
Snapshot {
scope: ObservedUnreadScope,
generation: String,
revision: u64,
last_acked_sequence: u64,
migration_complete: bool,
membership_seeded: bool,
channels: Vec<ChannelProjection>,
},
Delta {
scope: ObservedUnreadScope,
generation: String,
base_revision: u64,
revision: u64,
acked_sequence: u64,
upserts: Vec<ChannelProjection>,
removed: Vec<String>,
},
SnapshotRequired {
scope: ObservedUnreadScope,
generation: String,
revision: u64,
last_acked_sequence: u64,
},
}
fn db_path(app: &AppHandle) -> Result<PathBuf, String> {
let dir = app
.path()
.app_data_dir()
.map_err(|e| format!("resolve observed-unread data dir: {e}"))?;
std::fs::create_dir_all(&dir).map_err(|e| format!("create observed-unread data dir: {e}"))?;
Ok(dir.join("observed-unread.db"))
}
fn open_db(path: &Path) -> Result<Connection, String> {
let conn = Connection::open(path).map_err(|e| format!("open observed-unread db: {e}"))?;
conn.pragma_update(None, "busy_timeout", 5_000)
.map_err(|e| format!("configure observed-unread db: {e}"))?;
conn.pragma_update(None, "journal_mode", "WAL")
.map_err(|e| format!("configure observed-unread WAL: {e}"))?;
conn.execute_batch("CREATE TABLE IF NOT EXISTS schema_meta(version INTEGER NOT NULL);
INSERT INTO schema_meta(version) SELECT 1 WHERE NOT EXISTS(SELECT 1 FROM schema_meta);
CREATE TABLE IF NOT EXISTS scope_state(
scope TEXT PRIMARY KEY, generation TEXT NOT NULL, revision INTEGER NOT NULL DEFAULT 0,
last_sequence INTEGER NOT NULL DEFAULT 0, migration_complete INTEGER NOT NULL DEFAULT 0,
membership_seeded INTEGER NOT NULL DEFAULT 0);
CREATE TABLE IF NOT EXISTS observed_events(
scope TEXT NOT NULL, event_id TEXT NOT NULL, channel_id TEXT NOT NULL,
created_at INTEGER NOT NULL, root_id TEXT, high_priority INTEGER NOT NULL,
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 read_markers(
scope TEXT NOT NULL, context_id TEXT NOT NULL, read_at INTEGER NOT NULL,
PRIMARY KEY(scope,context_id));
CREATE TABLE IF NOT EXISTS unread_membership(
scope TEXT NOT NULL, kind TEXT NOT NULL, value TEXT NOT NULL,
PRIMARY KEY(scope,kind,value));")
.map_err(|e| format!("initialize observed-unread db: {e}"))?;
let version: i64 = conn
.query_row("SELECT version FROM schema_meta LIMIT 1", [], |row| {
row.get(0)
})
.map_err(|e| format!("read observed-unread schema: {e}"))?;
if version != SCHEMA_VERSION {
return Err(format!(
"unsupported observed-unread schema version {version}"
));
}
Ok(conn)
}
fn ensure_scope(tx: &Transaction<'_>, scope: &str) -> Result<(), String> {
tx.execute(
"INSERT OR IGNORE INTO scope_state(scope,generation) VALUES(?1,?2)",
params![scope, uuid::Uuid::new_v4().to_string()],
)
.map_err(|e| format!("initialize observed-unread scope: {e}"))?;
Ok(())
}
fn state(tx: &Transaction<'_>, scope: &str) -> Result<(String, u64, u64, bool, bool), String> {
tx.query_row("SELECT generation,revision,last_sequence,migration_complete,membership_seeded FROM scope_state WHERE scope=?1", [scope], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get::<_,i64>(3)? != 0,r.get::<_,i64>(4)? != 0)))
.map_err(|e| format!("read observed-unread scope state: {e}"))
}
fn valid_legacy_event(value: &serde_json::Value, channel_id: &str) -> Option<IngestEvent> {
let object = value.as_object()?;
Some(IngestEvent {
channel_id: channel_id.to_string(),
id: object.get("id")?.as_str()?.to_string(),
created_at: object.get("createdAt")?.as_u64()?,
root_id: match object.get("rootId")? {
serde_json::Value::Null => None,
v => Some(v.as_str()?.to_string()),
},
high_priority: object.get("highPriority")?.as_bool()?,
counts_toward_badge: object.get("countsTowardBadge")?.as_bool()?,
counts_toward_app_badge: object.get("countsTowardAppBadge")?.as_bool()?,
})
}
fn upsert_event(tx: &Transaction<'_>, scope: &str, event: &IngestEvent) -> Result<(), String> {
tx.execute("INSERT INTO observed_events(scope,event_id,channel_id,created_at,root_id,high_priority,counts_badge,counts_app_badge)
VALUES(?1,?2,?3,?4,?5,?6,?7,?8) ON CONFLICT(scope,event_id) DO NOTHING",
params![scope,event.id,event.channel_id,event.created_at,event.root_id,event.high_priority,event.counts_toward_badge,event.counts_toward_app_badge])
.map_err(|e| format!("upsert observed-unread event: {e}"))?;
Ok(())
}
fn seed_membership(tx: &Transaction<'_>, scope: &str, seed: &MembershipSeed) -> Result<(), String> {
// The renderer snapshot is authoritative while it remains the only writer.
// Replace transactionally so removals made while Buzz was closed are not
// silently resurrected by an insert-only seed.
tx.execute("DELETE FROM unread_membership WHERE scope=?1", [scope])
.map_err(|e| format!("reset unread membership: {e}"))?;
for (kind, values) in [
("participated", &seed.participated_root_ids),
("authored", &seed.authored_root_ids),
("mentioned", &seed.mentioned_root_ids),
("followed", &seed.followed_root_ids),
("muted_root", &seed.muted_root_ids),
("muted_channel", &seed.muted_channel_ids),
] {
for value in values {
tx.execute(
"INSERT OR IGNORE INTO unread_membership(scope,kind,value) VALUES(?1,?2,?3)",
params![scope, kind, value],
)
.map_err(|e| format!("seed unread membership: {e}"))?;
}
}
tx.execute(
"UPDATE scope_state SET membership_seeded=1 WHERE scope=?1",
[scope],
)
.map_err(|e| format!("mark unread membership seeded: {e}"))?;
Ok(())
}
fn prune(tx: &Transaction<'_>, scope: &str) -> Result<(), String> {
let cutoff = chrono::Utc::now().timestamp() - HORIZON_SECONDS;
tx.execute(
"DELETE FROM observed_events WHERE scope=?1 AND created_at<=?2",
params![scope, cutoff],
)
.map_err(|e| format!("age-prune observed unread: {e}"))?;
tx.execute("DELETE FROM observed_events WHERE rowid IN (SELECT rowid FROM (SELECT rowid,ROW_NUMBER() OVER(PARTITION BY channel_id ORDER BY created_at DESC,event_id DESC) rank FROM observed_events WHERE scope=?1) WHERE rank>?2)", params![scope,PER_CHANNEL_CAP]).map_err(|e| format!("channel-prune observed unread: {e}"))?;
tx.execute("DELETE FROM observed_events WHERE rowid IN (SELECT rowid FROM observed_events WHERE scope=?1 ORDER BY created_at DESC,event_id DESC LIMIT -1 OFFSET ?2)", params![scope,GLOBAL_CAP]).map_err(|e| format!("global-prune observed unread: {e}"))?;
Ok(())
}
fn marker(markers: &HashMap<String, u64>, key: &str) -> u64 {
markers.get(key).copied().unwrap_or(0)
}
fn projections(tx: &Transaction<'_>, scope: &str) -> Result<Vec<ChannelProjection>, String> {
let mut marker_stmt = tx
.prepare("SELECT context_id,read_at FROM read_markers WHERE scope=?1")
.map_err(|e| format!("prepare unread markers: {e}"))?;
let markers: HashMap<String, u64> = marker_stmt
.query_map([scope], |r| Ok((r.get(0)?, r.get(1)?)))
.map_err(|e| format!("query unread markers: {e}"))?
.collect::<Result<_, _>>()
.map_err(|e| format!("read unread markers: {e}"))?;
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| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, String>(1)?,
r.get::<_, u64>(2)?,
r.get::<_, Option<String>>(3)?,
r.get::<_, bool>(4)?,
r.get::<_, bool>(5)?,
r.get::<_, bool>(6)?,
))
})
.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}"))?;
let mut read_at = marker(&markers, &channel).max(marker(&markers, &format!("msg:{id}")));
if let Some(root) = &root {
read_at = read_at.max(marker(&markers, &format!("thread:{root}")));
}
if created <= read_at {
continue;
}
let entry = by_channel
.entry(channel.clone())
.or_insert(ChannelProjection {
channel_id: channel,
latest: 0,
count: 0,
badge_count: 0,
app_badge_count: 0,
top_level_unread: false,
high_priority_unread: false,
});
entry.latest = entry.latest.max(created);
entry.count += 1;
entry.badge_count += u64::from(badge);
entry.app_badge_count += u64::from(app);
entry.top_level_unread |= root.is_none();
entry.high_priority_unread |= high;
}
let mut result: Vec<_> = by_channel.into_values().collect();
result.sort_by(|a, b| a.channel_id.cmp(&b.channel_id));
Ok(result)
}
#[tauri::command]
pub(crate) fn observed_unread_open_scope(
request: OpenScopeRequest,
app: AppHandle,
store: State<'_, ObservedUnreadStore>,
) -> Result<ObservedUnreadResponse, String> {
let _guard = store.write_lock.lock().map_err(|e| e.to_string())?;
let mut conn = open_db(&db_path(&app)?)?;
let tx = conn
.transaction()
.map_err(|e| format!("begin observed-unread open: {e}"))?;
let scope = request.scope.key();
ensure_scope(&tx, &scope)?;
let (_, _, _, migration_complete, _) = state(&tx, &scope)?;
if !migration_complete {
if let Some(payload) = &request.legacy_payload {
if let Some(channels) = payload
.get("eventsByChannel")
.and_then(serde_json::Value::as_object)
{
for (channel, events) in channels {
if let Some(events) = events.as_array() {
for value in events {
if let Some(event) = valid_legacy_event(value, channel) {
upsert_event(&tx, &scope, &event)?;
}
}
}
}
}
}
tx.execute(
"UPDATE scope_state SET migration_complete=1 WHERE scope=?1",
[&scope],
)
.map_err(|e| format!("mark observed migration: {e}"))?;
}
if let Some(seed) = &request.membership_seed {
seed_membership(&tx, &scope, seed)?;
}
prune(&tx, &scope)?;
let channels = projections(&tx, &scope)?;
let (generation, revision, last, migrated, seeded) = state(&tx, &scope)?;
tx.commit()
.map_err(|e| format!("commit observed-unread open: {e}"))?;
Ok(ObservedUnreadResponse::Snapshot {
scope: request.scope,
generation,
revision,
last_acked_sequence: last,
migration_complete: migrated,
membership_seeded: seeded,
channels,
})
}
#[tauri::command]
pub(crate) fn observed_unread_ingest(
request: IngestRequest,
app: AppHandle,
store: State<'_, ObservedUnreadStore>,
) -> Result<ObservedUnreadResponse, String> {
let _guard = store.write_lock.lock().map_err(|e| e.to_string())?;
let mut conn = open_db(&db_path(&app)?)?;
let tx = conn
.transaction()
.map_err(|e| format!("begin observed ingest: {e}"))?;
let scope = request.scope.key();
ensure_scope(&tx, &scope)?;
let (generation, revision, last, _, _) = state(&tx, &scope)?;
if request.sequence <= last {
let channels = projections(&tx, &scope)?;
tx.commit()
.map_err(|e| format!("commit observed replay: {e}"))?;
return Ok(ObservedUnreadResponse::Snapshot {
scope: request.scope,
generation,
revision,
last_acked_sequence: last,
migration_complete: true,
membership_seeded: true,
channels,
});
}
if request.sequence != last + 1 || request.base_revision != revision {
return Ok(ObservedUnreadResponse::SnapshotRequired {
scope: request.scope,
generation,
revision,
last_acked_sequence: last,
});
}
let before = projections(&tx, &scope)?;
let before_by_channel: HashMap<_, _> = before
.into_iter()
.map(|projection| (projection.channel_id.clone(), projection))
.collect();
if request.clear_all {
tx.execute("DELETE FROM observed_events WHERE scope=?1", [&scope])
.map_err(|e| format!("clear observed scope: {e}"))?;
}
for channel in &request.clear_channels {
tx.execute(
"DELETE FROM observed_events WHERE scope=?1 AND channel_id=?2",
params![scope, channel],
)
.map_err(|e| format!("clear observed channel: {e}"))?;
}
for event in &request.events {
upsert_event(&tx, &scope, event)?;
}
for update in &request.membership {
if update.present {
tx.execute(
"INSERT OR IGNORE INTO unread_membership(scope,kind,value) VALUES(?1,?2,?3)",
params![scope, update.kind, update.value],
)
} else {
tx.execute(
"DELETE FROM unread_membership WHERE scope=?1 AND kind=?2 AND value=?3",
params![scope, update.kind, update.value],
)
}
.map_err(|e| format!("update unread membership: {e}"))?;
}
for update in &request.markers {
match update.read_at { Some(read_at)=>{tx.execute("INSERT INTO read_markers(scope,context_id,read_at) VALUES(?1,?2,?3) ON CONFLICT(scope,context_id) DO UPDATE SET read_at=MAX(read_at,excluded.read_at)",params![scope,update.context_id,read_at])},None=>tx.execute("DELETE FROM read_markers WHERE scope=?1 AND context_id=?2",params![scope,update.context_id])}.map_err(|e| format!("update observed marker: {e}"))?;
}
prune(&tx, &scope)?;
let after = projections(&tx, &scope)?;
let after_ids: HashSet<_> = after
.iter()
.map(|projection| projection.channel_id.clone())
.collect();
let removed: Vec<_> = before_by_channel
.keys()
.filter(|channel_id| !after_ids.contains(*channel_id))
.cloned()
.collect();
let upserts: Vec<_> = after
.into_iter()
.filter(|projection| before_by_channel.get(&projection.channel_id) != Some(projection))
.collect();
let next_revision = revision + 1;
tx.execute(
"UPDATE scope_state SET revision=?2,last_sequence=?3 WHERE scope=?1",
params![scope, next_revision, request.sequence],
)
.map_err(|e| format!("advance observed sequence: {e}"))?;
tx.commit()
.map_err(|e| format!("commit observed ingest: {e}"))?;
Ok(ObservedUnreadResponse::Delta {
scope: request.scope,
generation,
base_revision: revision,
revision: next_revision,
acked_sequence: request.sequence,
upserts,
removed,
})
}
pub(crate) fn load_membership(
app: &AppHandle,
scope: &ObservedUnreadScope,
) -> Result<HashMap<String, HashSet<String>>, String> {
let conn = open_db(&db_path(app)?)?;
let key = scope.key();
let mut stmt = conn
.prepare("SELECT kind,value FROM unread_membership WHERE scope=?1")
.map_err(|e| format!("prepare unread membership: {e}"))?;
let rows = stmt
.query_map([key], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})
.map_err(|e| format!("query unread membership: {e}"))?;
let mut result: HashMap<String, HashSet<String>> = HashMap::new();
for row in rows {
let (kind, value) = row.map_err(|e| format!("read unread membership: {e}"))?;
result.entry(kind).or_default().insert(value);
}
Ok(result)
}
pub(crate) fn flush(app: &AppHandle) {
if let Ok(path) = db_path(app) {
if let Ok(conn) = open_db(&path) {
let _ = conn.execute_batch("PRAGMA wal_checkpoint(PASSIVE);");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn scope() -> ObservedUnreadScope {
ObservedUnreadScope {
pubkey: "PK".into(),
relay_url: "wss://relay/".into(),
}
}
fn db() -> (tempfile::TempDir, Connection) {
let dir = tempfile::tempdir().unwrap();
let conn = open_db(&dir.path().join("observed-unread.db")).unwrap();
(dir, conn)
}
#[test]
fn ingest_replay_gap_prune_and_projection() {
let (_d, mut conn) = db();
let tx = conn.transaction().unwrap();
let key = scope().key();
ensure_scope(&tx, &key).unwrap();
upsert_event(
&tx,
&key,
&IngestEvent {
channel_id: "ch".into(),
id: "e".into(),
created_at: chrono::Utc::now().timestamp() as u64,
root_id: Some("root".into()),
high_priority: true,
counts_toward_badge: true,
counts_toward_app_badge: false,
},
)
.unwrap();
tx.execute(
"INSERT INTO read_markers(scope,context_id,read_at) VALUES(?1,'thread:root',0)",
[&key],
)
.unwrap();
let p = projections(&tx, &key).unwrap();
assert_eq!(p[0].count, 1);
assert_eq!(p[0].badge_count, 1);
tx.commit().unwrap();
}
#[test]
fn serialized_response_matches_typescript_contract() {
let actual = serde_json::to_value(ObservedUnreadResponse::Delta {
scope: scope(),
generation: "gen".into(),
base_revision: 4,
revision: 5,
acked_sequence: 7,
upserts: vec![ChannelProjection {
channel_id: "ch".into(),
latest: 42,
count: 2,
badge_count: 1,
app_badge_count: 1,
top_level_unread: true,
high_priority_unread: false,
}],
removed: vec!["old".into()],
})
.unwrap();
let expected = serde_json::json!({"kind":"delta","scope":{"pubkey":"PK","relayUrl":"wss://relay/"},"generation":"gen","baseRevision":4,"revision":5,"ackedSequence":7,"upserts":[{"channelId":"ch","latest":42,"count":2,"badgeCount":1,"appBadgeCount":1,"topLevelUnread":true,"highPriorityUnread":false}],"removed":["old"]});
assert_eq!(actual, expected);
}
}
+1
View File
@@ -19,6 +19,7 @@ pub(crate) fn shut_down_app(app: &tauri::AppHandle, shutdown_done: &std::sync::a
.store(true, Ordering::SeqCst);
if !shutdown_done.swap(true, Ordering::SeqCst) {
prevent_sleep::release(&app.state::<AppState>().prevent_sleep);
crate::observed_unread::flush(app);
app.state::<crate::terminal_runtime::TerminalSessions>()
.shutdown_all();
if let Err(error) = shutdown_managed_agents(app) {
+42 -27
View File
@@ -1,10 +1,10 @@
//! Batched native unread catch-up.
//!
//! The renderer supplies its history-derived notification membership because
//! those sets are still renderer-owned until the native observed-unread store
//! lands. Rust performs every channel REQ over the shared authenticated session,
//! then classifies the complete successful batch in two passes so a root learned
//! anywhere in pass one is visible everywhere in pass two.
//! Native unread catch-up consumes notification membership from the observed-
//! unread SQLite store rather than serializing renderer-owned sets on every
//! request. Rust performs every channel REQ over the shared authenticated
//! session, then classifies the complete successful batch in two passes so a
//! root learned anywhere in pass one is visible everywhere in pass two.
use std::{collections::HashSet, time::Duration};
@@ -14,7 +14,7 @@ use buzz_core_pkg::kind::{
};
use nostr::Event;
use serde::{Deserialize, Serialize};
use tauri::State;
use tauri::{AppHandle, State};
use tokio::{sync::Semaphore, task::JoinSet};
use crate::{app_state::AppState, native_relay_client::NativeRelayClient};
@@ -28,11 +28,6 @@ const REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
pub(crate) struct UnreadCatchUpRequest {
channels: Vec<CatchUpChannel>,
self_pubkey: String,
participated_root_ids: HashSet<String>,
authored_root_ids: HashSet<String>,
mentioned_root_ids: HashSet<String>,
followed_root_ids: HashSet<String>,
muted_root_ids: HashSet<String>,
muted_channel_ids: HashSet<String>,
}
@@ -142,6 +137,7 @@ pub(crate) async fn unread_catch_up(
request: UnreadCatchUpRequest,
state: State<'_, AppState>,
relay_client: State<'_, NativeRelayClient>,
app: AppHandle,
) -> Result<UnreadCatchUpResponse, String> {
let keys = state.signing_keys()?;
let owner = keys.public_key().to_hex();
@@ -222,7 +218,14 @@ pub(crate) async fn unread_catch_up(
return Err("unread catch-up scope changed while fetching".to_string());
}
let mut channels = classify_batch(&request, fetched);
let membership = crate::observed_unread::load_membership(
&app,
&crate::observed_unread::ObservedUnreadScope {
pubkey: owner,
relay_url,
},
)?;
let mut channels = classify_batch(&request, fetched, &membership);
channels.extend(failures);
Ok(UnreadCatchUpResponse { channels })
}
@@ -230,11 +233,12 @@ pub(crate) async fn unread_catch_up(
fn classify_batch(
request: &UnreadCatchUpRequest,
fetched: Vec<FetchedChannel>,
membership: &std::collections::HashMap<String, HashSet<String>>,
) -> Vec<ChannelResult> {
let self_pubkey = request.self_pubkey.to_lowercase();
let mut participated = request.participated_root_ids.clone();
let mut authored = request.authored_root_ids.clone();
let mut mentioned = request.mentioned_root_ids.clone();
let mut participated = membership.get("participated").cloned().unwrap_or_default();
let mut authored = membership.get("authored").cloned().unwrap_or_default();
let mut mentioned = membership.get("mentioned").cloned().unwrap_or_default();
// Pass one is deliberately global, not per-channel: notification validity
// depends on roots learned from history, while the command observes a batch.
@@ -275,7 +279,14 @@ fn classify_batch(
.channel
.read_at
.is_some_and(|read_at| event.created_at <= read_at)
|| !should_notify(&event, &self_pubkey, request, &participated, &authored)
|| !should_notify(
&event,
&self_pubkey,
request,
membership,
&participated,
&authored,
)
{
continue;
}
@@ -383,6 +394,7 @@ fn should_notify(
event: &EventView,
self_pubkey: &str,
request: &UnreadCatchUpRequest,
membership: &std::collections::HashMap<String, HashSet<String>>,
participated: &HashSet<String>,
authored: &HashSet<String>,
) -> bool {
@@ -405,11 +417,16 @@ fn should_notify(
let Some(root_id) = reference.root_id else {
return false;
};
if request.muted_root_ids.contains(&root_id) {
if membership
.get("muted_root")
.is_some_and(|set| set.contains(&root_id))
{
return false;
}
participated.contains(&root_id)
|| request.followed_root_ids.contains(&root_id)
|| membership
.get("followed")
.is_some_and(|set| set.contains(&root_id))
|| authored.contains(&root_id)
}
@@ -430,6 +447,8 @@ fn has_tag_value(tags: &[Vec<String>], name: &str, value: &str) -> bool {
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use super::*;
fn event(id: &str, pubkey: &str, created_at: u64, tags: &[&[&str]]) -> EventView {
@@ -450,11 +469,6 @@ mod tests {
UnreadCatchUpRequest {
channels: vec![],
self_pubkey: "self".into(),
participated_root_ids: HashSet::new(),
authored_root_ids: HashSet::new(),
mentioned_root_ids: HashSet::new(),
followed_root_ids: HashSet::new(),
muted_root_ids: HashSet::new(),
muted_channel_ids: HashSet::new(),
}
}
@@ -486,7 +500,7 @@ mod tests {
),
],
}];
let result = classify_batch(&req, fetched);
let result = classify_batch(&req, fetched, &HashMap::new());
let ChannelResult::Success {
observed_events,
discovered,
@@ -507,8 +521,9 @@ mod tests {
#[test]
fn same_second_marker_and_mutes_match_renderer_rules() {
let mut req = request();
req.muted_root_ids.insert("muted".into());
let req = request();
let mut membership = HashMap::new();
membership.insert("muted_root".into(), HashSet::from(["muted".into()]));
let channel = CatchUpChannel {
id: "ch".into(),
channel_type: "stream".into(),
@@ -534,7 +549,7 @@ mod tests {
),
],
}];
let result = classify_batch(&req, fetched);
let result = classify_batch(&req, fetched, &membership);
let ChannelResult::Success {
observed_events,
max_trigger,
@@ -1,59 +1,64 @@
import * as React from "react";
import {
flushObservedUnreadWrite,
clearObservedUnreadStorage,
deriveLatestByChannel,
observedUnreadStorageKey,
pruneObservedUnreadByMarkers,
readObservedUnreadFromStorage,
scheduleObservedUnreadWrite,
deriveLatestByChannel,
clearObservedUnreadStorage,
flushObservedUnreadWrite,
type ObservedUnreadRefs,
} from "@/features/channels/observedUnreadStorage";
import { activityScopeKey } from "@/features/channels/threadActivityStorage";
import type { ObservedUnreadEvent } from "@/features/channels/unreadChannelCounts";
import {
ingestObservedUnread,
openObservedUnreadScope,
type ObservedUnreadProjection,
type ObservedUnreadResponse,
type ObservedUnreadWireEvent,
} from "@/shared/api/tauriObservedUnread";
export type ObservedUnreadPersistence = {
/** Scope key loaded into the refs ("" until identity is known). */
scopeLoadedRef: React.MutableRefObject<string>;
/** Current scope derived from normalized pubkey + relay. */
currentScope: string;
/**
* Returns true when the identity-reset effect has committed and the observed
* refs hold data for the current scope. Always reads the ref at call time.
* Use as a scope guard before projecting or mutating the observed refs.
*/
projectionsRef: React.MutableRefObject<Map<string, ObservedUnreadProjection>>;
isNative: () => boolean;
isScopeLoaded: () => boolean;
/**
* Call after recording a new observed event to schedule a debounced write.
* @param scope - the currentScope value captured this render
*/
schedule: (scope: string) => void;
/** Remove a single channel from the persisted cache (clearObserved path). */
schedule: (
scope: string,
channelId?: string,
event?: ObservedUnreadEvent,
) => void;
removeChannel: (channelId: string) => void;
/** Clear the entire persisted cache (mark-all-read path). */
updateMembership: (kind: string, value: string, present: boolean) => void;
clearAll: () => void;
};
/**
* Additional options for useObservedUnreadPersistence.
*/
export type UseObservedUnreadPersistenceOptions = {
/**
* Called when a marker-prune pass removes at least one event from the
* in-memory refs. The hook itself cannot bump the parent's version counter;
* pass a stable callback (e.g. useEvent-wrapped bumpLatestVersion) here so
* the UI re-renders when stale observed events are swept out.
*/
type Options = {
onPruned?: () => void;
membershipSeed?: {
participatedRootIds: string[];
authoredRootIds: string[];
mentionedRootIds: string[];
followedRootIds: string[];
mutedRootIds: string[];
mutedChannelIds: string[];
};
};
type QueuedObservedUnreadEvent = {
scope: string;
event: ObservedUnreadWireEvent;
};
type NativeState = {
scope: { pubkey: string; relayUrl: string };
generation: string;
revision: number;
sequence: number;
};
/**
* Hook that manages the observed-unread localStorage persistence layer for
* useUnreadChannels. Owns the scope ref, debounce timer, pagehide flush, and
* marker-prune effect.
*
* Returns a stable API object; callers read `scopeLoadedRef` and `currentScope`
* to guard writes, and call `schedule`/`removeChannel`/`clearAll` on mutations.
*/
export function useObservedUnreadPersistence(
normalizedPubkey: string | null,
normalizedRelayUrl: string,
@@ -65,15 +70,20 @@ export function useObservedUnreadPersistence(
Map<string, Map<string, ObservedUnreadEvent>>
>,
latestByChannelRef: React.MutableRefObject<Map<string, number>>,
options: UseObservedUnreadPersistenceOptions = {},
options: Options = {},
): ObservedUnreadPersistence {
const optionsRef = React.useRef(options);
optionsRef.current = options;
const currentScope = activityScopeKey(normalizedPubkey, normalizedRelayUrl);
const { onPruned } = options;
const scopeLoadedRef = React.useRef<string>("");
const scopeLoadedRef = React.useRef("");
const projectionsRef = React.useRef(
new Map<string, ObservedUnreadProjection>(),
);
const nativeRef = React.useRef<NativeState | null>(null);
const nativeFailedRef = React.useRef(false);
const timerRef = React.useRef<ReturnType<typeof setTimeout> | null>(null);
const queueRef = React.useRef<QueuedObservedUnreadEvent[]>([]);
const chainRef = React.useRef(Promise.resolve());
const persistRefs = React.useRef<ObservedUnreadRefs>({
eventsRef: observedUnreadEventsByChannelRef,
scopeLoadedRef,
@@ -81,129 +91,382 @@ export function useObservedUnreadPersistence(
});
persistRefs.current.eventsRef = observedUnreadEventsByChannelRef;
// pagehide: synchronously flush any pending debounce before the webview
// unloads. Cmd+R reloads within 500ms of teardown; without this flush the
// last observed event would be lost within the debounce window.
const reopen = React.useCallback(
async (scope: { pubkey: string; relayUrl: string }) => {
const response = await openObservedUnreadScope({
scope,
legacyPayload: null,
membershipSeed: null,
});
if (response.kind !== "snapshot") return;
nativeRef.current = {
scope,
generation: response.generation,
revision: response.revision,
sequence: response.lastAckedSequence,
};
projectionsRef.current = new Map(
response.channels.map((item) => [item.channelId, item]),
);
optionsRef.current.onPruned?.();
},
[],
);
const apply = React.useCallback(
(response: ObservedUnreadResponse) => {
const state = nativeRef.current;
if (
!state ||
activityScopeKey(response.scope.pubkey, response.scope.relayUrl) !==
scopeLoadedRef.current
)
return;
if (response.kind === "snapshotRequired") {
void reopen(state.scope);
return;
}
if (response.kind === "snapshot") {
if (
response.generation !== state.generation ||
response.revision >= state.revision
) {
projectionsRef.current = new Map(
response.channels.map((item) => [item.channelId, item]),
);
nativeRef.current = {
...state,
generation: response.generation,
revision: response.revision,
sequence: response.lastAckedSequence,
};
optionsRef.current.onPruned?.();
}
return;
}
if (
response.generation !== state.generation ||
response.baseRevision !== state.revision
) {
void reopen(state.scope);
return;
}
const next = new Map(projectionsRef.current);
for (const id of response.removed) next.delete(id);
for (const item of response.upserts) next.set(item.channelId, item);
projectionsRef.current = next;
nativeRef.current = {
...state,
revision: response.revision,
sequence: response.ackedSequence,
};
optionsRef.current.onPruned?.();
},
[reopen],
);
const flushNative = React.useCallback(() => {
const state = nativeRef.current;
if (!state || queueRef.current.length === 0) return;
const stateScope = activityScopeKey(
state.scope.pubkey,
state.scope.relayUrl,
);
const queued = queueRef.current.filter(({ scope }) => scope === stateScope);
queueRef.current = queueRef.current.filter(
({ scope }) => scope !== stateScope,
);
if (queued.length === 0) return;
const events = queued.map(({ event }) => event);
chainRef.current = chainRef.current
.then(() => {
const current = nativeRef.current;
if (
!current ||
current.scope.pubkey !== state.scope.pubkey ||
current.scope.relayUrl !== state.scope.relayUrl
) {
queueRef.current = [...queued, ...queueRef.current];
return;
}
return ingestObservedUnread({
scope: current.scope,
sequence: current.sequence + 1,
baseRevision: current.revision,
events,
markers: [],
membership: [],
clearChannels: [],
clearAll: false,
}).then(apply);
})
.catch(() => {
queueRef.current = [...queued, ...queueRef.current];
nativeFailedRef.current = true;
nativeRef.current = null;
scheduleObservedUnreadWrite(
scopeLoadedRef.current,
persistRefs.current,
);
});
}, [apply]);
React.useEffect(() => {
const refs = persistRefs.current;
const onPageHide = () => flushObservedUnreadWrite(refs);
const onPageHide = () => {
flushNative();
flushObservedUnreadWrite(persistRefs.current);
};
window.addEventListener("pagehide", onPageHide);
return () => window.removeEventListener("pagehide", onPageHide);
}, []);
}, [flushNative]);
// Hydrate refs from storage whenever identity/relay changes. Flush the OLD
// scope first so no event is lost before we clobber the refs.
// biome-ignore lint/correctness/useExhaustiveDependencies: normalizedRelayUrl is intentional reset signal alongside normalizedPubkey
// Identity and relay are the intentional reset signals; refs and stable callbacks
// are mutable containers/transport helpers rather than reset triggers.
// biome-ignore lint/correctness/useExhaustiveDependencies: scope reset is keyed only by normalized identity and relay
React.useEffect(() => {
flushNative();
flushObservedUnreadWrite(persistRefs.current);
nativeRef.current = null;
nativeFailedRef.current = false;
projectionsRef.current = new Map();
observedUnreadEventsByChannelRef.current = new Map();
latestByChannelRef.current = new Map();
if (normalizedPubkey && normalizedRelayUrl) {
const stored = readObservedUnreadFromStorage(
normalizedPubkey,
normalizedRelayUrl,
);
if (stored && stored.size > 0) {
observedUnreadEventsByChannelRef.current = stored;
latestByChannelRef.current = deriveLatestByChannel(stored);
}
}
scopeLoadedRef.current = activityScopeKey(
normalizedPubkey,
normalizedRelayUrl,
);
// On unmount (or before next effect run), flush the current scope so any
// in-flight debounce is persisted before refs are clobbered.
scopeLoadedRef.current = "";
if (!normalizedPubkey || !normalizedRelayUrl) return;
let active = true;
const scope = { pubkey: normalizedPubkey, relayUrl: normalizedRelayUrl };
const key = observedUnreadStorageKey(normalizedPubkey, normalizedRelayUrl);
let legacyPayload: unknown = null;
try {
const raw = window.localStorage.getItem(key);
legacyPayload = raw ? JSON.parse(raw) : null;
} catch {}
void openObservedUnreadScope({
scope,
legacyPayload,
membershipSeed: optionsRef.current.membershipSeed ?? null,
})
.then((response) => {
if (!active) return;
if (response.kind !== "snapshot")
throw new Error(
"native observed-unread open did not return snapshot",
);
nativeRef.current = {
scope,
generation: response.generation,
revision: response.revision,
sequence: response.lastAckedSequence,
};
scopeLoadedRef.current = currentScope;
apply(response);
if (response.migrationComplete) {
try {
window.localStorage.removeItem(key);
} catch {}
}
})
.catch(() => {
if (!active) return;
nativeFailedRef.current = true;
const stored = readObservedUnreadFromStorage(
normalizedPubkey,
normalizedRelayUrl,
);
if (stored) {
observedUnreadEventsByChannelRef.current = stored;
latestByChannelRef.current = deriveLatestByChannel(stored);
}
scopeLoadedRef.current = currentScope;
optionsRef.current.onPruned?.();
});
return () => {
active = false;
flushNative();
flushObservedUnreadWrite(persistRefs.current);
};
}, [normalizedPubkey, normalizedRelayUrl]);
// Marker prune: whenever read state advances, drop events now covered by
// their channel/thread/msg marker, persist if anything changed.
// biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion + isReadStateReady are intentional prune triggers
// readStateVersion is the intentional invalidation signal; marker readers and
// mutable refs are sampled when that signal advances.
// biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion and readiness intentionally drive marker synchronization
React.useEffect(() => {
if (!isReadStateReady) return;
const scope = scopeLoadedRef.current;
if (!scope || scope !== currentScope) return;
const changed = pruneObservedUnreadByMarkers(
observedUnreadEventsByChannelRef.current,
latestByChannelRef.current,
getEffectiveTimestamp,
getOwnTimestamp,
);
if (changed) {
scheduleObservedUnreadWrite(scope, persistRefs.current);
onPruned?.();
if (!isReadStateReady || scopeLoadedRef.current !== currentScope) return;
const state = nativeRef.current;
if (state) {
const contexts = new Set<string>();
for (const [
channelId,
events,
] of observedUnreadEventsByChannelRef.current) {
contexts.add(channelId);
for (const event of events.values()) {
contexts.add(`msg:${event.id}`);
if (event.rootId) contexts.add(`thread:${event.rootId}`);
}
}
const markers = [...contexts].map((contextId) => ({
contextId,
readAt:
contextId.startsWith("thread:") || contextId.startsWith("msg:")
? getOwnTimestamp(contextId)
: getEffectiveTimestamp(contextId),
}));
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: [],
markers,
membership: [],
clearChannels: [],
clearAll: false,
}).then(apply);
});
} else if (
pruneObservedUnreadByMarkers(
observedUnreadEventsByChannelRef.current,
latestByChannelRef.current,
getEffectiveTimestamp,
getOwnTimestamp,
)
) {
scheduleObservedUnreadWrite(currentScope, persistRefs.current);
optionsRef.current.onPruned?.();
}
}, [readStateVersion, isReadStateReady]);
const schedule = React.useCallback(
(scope: string) => scheduleObservedUnreadWrite(scope, persistRefs.current),
[],
(scope: string, channelId?: string, event?: ObservedUnreadEvent) => {
if (scopeLoadedRef.current !== scope) return;
if (nativeRef.current && channelId && event) {
queueRef.current.push({ scope, event: { channelId, ...event } });
if (timerRef.current !== null) clearTimeout(timerRef.current);
timerRef.current = setTimeout(() => {
timerRef.current = null;
flushNative();
}, 1_000);
} else scheduleObservedUnreadWrite(scope, persistRefs.current);
},
[flushNative],
);
const mutateClear = React.useCallback(
(clearChannels: string[], clearAll: boolean) => {
const state = nativeRef.current;
if (!state) return false;
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: [],
markers: [],
membership: [],
clearChannels,
clearAll,
}).then(apply);
});
return true;
},
[apply, flushNative],
);
const updateMembership = React.useCallback(
(kind: string, value: string, present: boolean) => {
const state = nativeRef.current;
if (!state) 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: [],
markers: [],
membership: [{ kind, value, present }],
clearChannels: [],
clearAll: false,
}).then(apply);
});
},
[apply, flushNative],
);
// biome-ignore lint/correctness/useExhaustiveDependencies: mutable storage refs are stable containers
const removeChannel = React.useCallback(
(channelId: string) => {
// Reject if the loaded scope has drifted — a stale callback during A→B
// transitions must not cancel B's pending snapshot or corrupt B's refs.
if (scopeLoadedRef.current !== currentScope) return;
// Delete from both in-memory refs so the projection no longer sees this
// channel. Then replace any pending snapshot with a new snapshot of the
// current full map — never cancel-without-replacement, which would lose
// unsaved sibling-channel events on the next reload.
observedUnreadEventsByChannelRef.current.delete(channelId);
latestByChannelRef.current.delete(channelId);
scheduleObservedUnreadWrite(currentScope, persistRefs.current);
projectionsRef.current.delete(channelId);
if (!mutateClear([channelId], false)) {
observedUnreadEventsByChannelRef.current.delete(channelId);
latestByChannelRef.current.delete(channelId);
scheduleObservedUnreadWrite(currentScope, persistRefs.current);
}
},
[currentScope, observedUnreadEventsByChannelRef, latestByChannelRef],
[currentScope, mutateClear],
);
// biome-ignore lint/correctness/useExhaustiveDependencies: mutable storage refs are stable containers
const clearAll = React.useCallback(() => {
// Reject if the loaded scope has drifted — a stale callback must not
// cancel the new scope's pending snapshot or clear the wrong bucket.
if (scopeLoadedRef.current !== currentScope) return;
// Cancel any pending snapshot and clear both in-memory refs before touching
// storage — the parent no longer resets the refs directly, so this is the
// single transactional clear path for mark-all-read.
if (timerRef.current !== null) {
clearTimeout(timerRef.current);
timerRef.current = null;
projectionsRef.current = new Map();
if (!mutateClear([], true)) {
observedUnreadEventsByChannelRef.current = new Map();
latestByChannelRef.current = new Map();
clearObservedUnreadStorage(normalizedPubkey ?? "", normalizedRelayUrl);
}
observedUnreadEventsByChannelRef.current = new Map();
latestByChannelRef.current = new Map();
clearObservedUnreadStorage(normalizedPubkey ?? "", normalizedRelayUrl);
}, [
currentScope,
normalizedPubkey,
normalizedRelayUrl,
observedUnreadEventsByChannelRef,
latestByChannelRef,
]);
// isScopeLoaded reads the ref at call time — always fresh, never a stale
// snapshot from a closed-over useMemo value.
}, [currentScope, mutateClear, normalizedPubkey, normalizedRelayUrl]);
const isScopeLoaded = React.useCallback(
() => scopeLoadedRef.current === currentScope,
[currentScope],
);
// Stable API object: only reconstructed when scope or stable callbacks change.
// This prevents useCallback deps in the parent from seeing a new object each
// render, which would restart catch-up REQs on every unrelated re-render.
const isNative = React.useCallback(
() => nativeRef.current !== null && !nativeFailedRef.current,
[],
);
return React.useMemo(
() => ({
scopeLoadedRef,
currentScope,
projectionsRef,
isNative,
isScopeLoaded,
schedule,
removeChannel,
updateMembership,
clearAll,
}),
[currentScope, isScopeLoaded, schedule, removeChannel, clearAll],
[
currentScope,
isNative,
isScopeLoaded,
schedule,
removeChannel,
updateMembership,
clearAll,
],
);
}
@@ -1,6 +1,5 @@
import * as React from "react";
import {
EMPTY_SET,
useLiveChannelUpdates,
type UseLiveChannelUpdatesOptions,
} from "@/features/channels/useLiveChannelUpdates";
@@ -67,6 +66,7 @@ type UseUnreadChannelsOptions = UseLiveChannelUpdatesOptions & {
// filter to find one external trigger message. 1000 matches the live sub's
// per-channel limit elsewhere in the app.
const CATCH_UP_LIMIT = 1000;
const EMPTY_ROOT_IDS: ReadonlySet<string> = new Set();
export function channelCatchUpEventKinds(
channelType: Channel["channelType"] | undefined,
@@ -213,8 +213,8 @@ export function useUnreadChannels(
// Stable ref for the caller-supplied muted channel IDs. Updated every render
// so the catch-up loop always reads the latest set without being a dep.
const mutedChannelIdsRef = React.useRef<ReadonlySet<string>>(new Set());
mutedChannelIdsRef.current = mutedChannelIdsOption ?? new Set();
const mutedChannelIdsRef = React.useRef<ReadonlySet<string>>(EMPTY_ROOT_IDS);
mutedChannelIdsRef.current = mutedChannelIdsOption ?? EMPTY_ROOT_IDS;
// Thread reply events that triggered notifications — surfaced in the Home
// activity feed as synthetic FeedItems. The buffer is the source of truth
@@ -241,28 +241,6 @@ export function useUnreadChannels(
0,
);
// Persistence layer: hydration, pagehide flush, scope fence, write-through, marker-prune.
const observedPersistence = useObservedUnreadPersistence(
normalizedPubkey,
normalizedRelayUrl,
isReadStateReady,
readStateVersion,
getEffectiveTimestamp,
getOwnTimestamp,
observedUnreadEventsByChannelRef,
latestByChannelRef,
{ onPruned: bumpLatestVersion },
);
// Thread-activity persistence: coalesced writes, pagehide/visibility flush,
// hydration + legacy-key cleanup. Owns the loaded scope for the buffer above.
const activityPersistence = useThreadActivityPersistence(
normalizedPubkey,
normalizedRelayUrl,
threadActivityRef,
);
const currentActivityScope = activityPersistence.currentScope;
// Reset all in-session state when the identity or relay changes. In-memory
// caches are cleared; persisted stores are loaded for the new pubkey (so
// forced-unread, participation, etc. are correct for the new identity).
@@ -286,6 +264,57 @@ export function useUnreadChannels(
bumpMembershipVersion();
}, [pubkey, relayClient, normalizedRelayUrl]);
// Persistence layer: hydration, pagehide flush, scope fence, write-through, marker-prune.
const observedPersistence = useObservedUnreadPersistence(
normalizedPubkey,
normalizedRelayUrl,
isReadStateReady,
readStateVersion,
getEffectiveTimestamp,
getOwnTimestamp,
observedUnreadEventsByChannelRef,
latestByChannelRef,
{
onPruned: bumpLatestVersion,
membershipSeed: {
participatedRootIds: [...participatedRootIdsRef.current],
authoredRootIds: [...authoredRootIdsRef.current],
mentionedRootIds: [...mentionedRootIdsRef.current],
followedRootIds: [...(options.followedRootIds ?? EMPTY_ROOT_IDS)],
mutedRootIds: [...mutedRootIdsRef.current],
mutedChannelIds: [...mutedChannelIdsRef.current],
},
},
);
const followedMembershipRef = React.useRef(new Set<string>());
React.useEffect(() => {
const desired = options.followedRootIds ?? EMPTY_ROOT_IDS;
if (!observedPersistence.isScopeLoaded()) {
followedMembershipRef.current = new Set(desired);
return;
}
for (const rootId of desired) {
if (!followedMembershipRef.current.has(rootId)) {
observedPersistence.updateMembership("followed", rootId, true);
}
}
for (const rootId of followedMembershipRef.current) {
if (!desired.has(rootId)) {
observedPersistence.updateMembership("followed", rootId, false);
}
}
followedMembershipRef.current = new Set(desired);
}, [observedPersistence, options.followedRootIds]);
// Thread-activity persistence: coalesced writes, pagehide/visibility flush,
// hydration + legacy-key cleanup. Owns the loaded scope for the buffer above.
const activityPersistence = useThreadActivityPersistence(
normalizedPubkey,
normalizedRelayUrl,
threadActivityRef,
);
const currentActivityScope = activityPersistence.currentScope;
// `topLevelOnly`: passive channel-open path (NIP-RS Option 1) — marker lands at newest
// top-level msg without folding observed replies; leaves refs intact so the dot persists
// until an explicit mark-read. Explicit reads omit this flag and clear the refs.
@@ -357,10 +386,11 @@ export function useUnreadChannels(
const sizeBefore = target.size;
target.add(rootId);
if (target.size === sizeBefore) return false;
observedPersistence.updateMembership("mentioned", rootId, true);
mentionedStore.write(normalizedPubkey, target);
return true;
},
[normalizedPubkey],
[normalizedPubkey, observedPersistence],
);
// Records an external trigger event and schedules persistence.
@@ -375,7 +405,11 @@ export function useUnreadChannels(
CATCH_UP_LIMIT,
);
if (didRecord)
observedPersistence.schedule(observedPersistence.currentScope);
observedPersistence.schedule(
observedPersistence.currentScope,
channelId,
event,
);
return didRecord;
},
[observedPersistence],
@@ -457,11 +491,16 @@ export function useUnreadChannels(
// to an already-tracked root is a no-op for the notify gate, so skipping
// the bump avoids a wasted snapshot re-allocation + gate recompute.
if (targetSet.size !== sizeBefore) {
observedPersistence.updateMembership(
isParticipation ? "participated" : "authored",
ref.rootId ?? event.id,
true,
);
bumpMembershipVersion();
}
bumpLatestVersion();
},
[normalizedPubkey],
[normalizedPubkey, observedPersistence],
);
const recordThreadInteraction = React.useCallback(
@@ -472,12 +511,17 @@ export function useUnreadChannels(
const sizeBefore = target.size;
target.add(normalizedRootId);
if (target.size === sizeBefore) return;
observedPersistence.updateMembership(
"participated",
normalizedRootId,
true,
);
if (normalizedPubkey !== null) {
participationStore.write(normalizedPubkey, target);
}
bumpMembershipVersion();
},
[normalizedPubkey],
[normalizedPubkey, observedPersistence],
);
const handleThreadReplyNotification = React.useCallback(
@@ -515,23 +559,25 @@ export function useUnreadChannels(
const muteThread = React.useCallback(
(rootId: string) => {
mutedRootIdsRef.current.add(rootId);
observedPersistence.updateMembership("muted_root", rootId, true);
if (normalizedPubkey !== null) {
mutedStore.write(normalizedPubkey, mutedRootIdsRef.current);
}
bumpLatestVersion();
},
[normalizedPubkey],
[normalizedPubkey, observedPersistence],
);
const unmuteThread = React.useCallback(
(rootId: string) => {
mutedRootIdsRef.current.delete(rootId);
observedPersistence.updateMembership("muted_root", rootId, false);
if (normalizedPubkey !== null) {
mutedStore.write(normalizedPubkey, mutedRootIdsRef.current);
}
bumpLatestVersion();
},
[normalizedPubkey],
[normalizedPubkey, observedPersistence],
);
useLiveChannelUpdates(channels, activeChannelId, {
@@ -590,10 +636,9 @@ export function useUnreadChannels(
const authoredSizeBefore = authoredRootIdsRef.current.size;
const mentionedSizeBefore = mentionedRootIdsRef.current.size;
// Membership remains renderer-owned until E's native observed-unread store,
// so unchanged sets cross IPC on every catch-up. The five 1,000-entry
// stores bound that interim cost at roughly 332 KiB per request. Command
// arguments use a fetch body, so this cost is linear with no size cliff.
// E's native observed-unread store owns notification membership, so this
// request scales with channels being caught up rather than five 1,000-id
// sets. The prior 332 KiB-at-cap payload is retired at this boundary.
void unreadCatchUp({
channels: toFetch.map((channelId) => {
const channel = channels.find(
@@ -607,11 +652,6 @@ export function useUnreadChannels(
};
}),
selfPubkey: normalizedPubkey ?? "",
participatedRootIds: [...participatedRootIdsRef.current],
authoredRootIds: [...authoredRootIdsRef.current],
mentionedRootIds: [...mentionedRootIdsRef.current],
followedRootIds: [...(options.followedRootIds ?? EMPTY_SET)],
mutedRootIds: [...mutedRootIdsRef.current],
mutedChannelIds: [...mutedChannelIdsRef.current],
})
.then(({ channels: results }) => {
@@ -633,17 +673,30 @@ export function useUnreadChannels(
for (const rootId of result.discovered.participated) {
const before = participatedRootIdsRef.current.size;
participatedRootIdsRef.current.add(rootId);
didDiscover ||= participatedRootIdsRef.current.size !== before;
if (participatedRootIdsRef.current.size !== before) {
observedPersistence.updateMembership(
"participated",
rootId,
true,
);
didDiscover = true;
}
}
for (const rootId of result.discovered.authored) {
const before = authoredRootIdsRef.current.size;
authoredRootIdsRef.current.add(rootId);
didDiscover ||= authoredRootIdsRef.current.size !== before;
if (authoredRootIdsRef.current.size !== before) {
observedPersistence.updateMembership("authored", rootId, true);
didDiscover = true;
}
}
for (const rootId of result.discovered.mentioned) {
const before = mentionedRootIdsRef.current.size;
mentionedRootIdsRef.current.add(rootId);
didDiscover ||= mentionedRootIdsRef.current.size !== before;
if (mentionedRootIdsRef.current.size !== before) {
observedPersistence.updateMembership("mentioned", rootId, true);
didDiscover = true;
}
}
allThreadReplies.push(...result.activityRows);
for (const event of result.observedEvents) {
@@ -763,8 +816,12 @@ export function useUnreadChannels(
(messageId) => getOwnTimestamp(`msg:${messageId}`),
);
const unreadCount =
latestByChannelRef.current.get(channel.id) === undefined
const nativeProjection = observedPersistence.isNative()
? observedPersistence.projectionsRef.current.get(channel.id)
: undefined;
const unreadCount = nativeProjection
? nativeProjection.count
: latestByChannelRef.current.get(channel.id) === undefined
? 0
: countUnreadObservedEvents(observedEvents, readAtForObservedEvent);
if (unreadCount === 0) {
@@ -778,28 +835,39 @@ export function useUnreadChannels(
unread.add(channel.id);
if (
hasUnreadTopLevelObservedEvent(observedEvents, readAtForObservedEvent)
nativeProjection?.topLevelUnread ||
(!nativeProjection &&
hasUnreadTopLevelObservedEvent(
observedEvents,
readAtForObservedEvent,
))
) {
topLevelUnread.add(channel.id);
}
const badgeCount = countUnreadBadgeObservedEvents(
observedEvents,
readAtForObservedEvent,
);
const badgeCount =
nativeProjection?.badgeCount ??
countUnreadBadgeObservedEvents(
observedEvents,
readAtForObservedEvent,
);
counts.set(channel.id, badgeCount);
unreadChannelNotificationCount += countUnreadAppBadgeObservedEvents(
observedEvents,
readAtForObservedEvent,
);
unreadChannelNotificationCount +=
nativeProjection?.appBadgeCount ??
countUnreadAppBadgeObservedEvents(
observedEvents,
readAtForObservedEvent,
);
// DM channels: any unread DM is high-priority.
if (channel.channelType === "dm") {
highPriority.add(channel.id);
} else if (
countUnreadHighPriorityObservedEvents(
observedEvents,
readAtForObservedEvent,
) > 0
nativeProjection?.highPriorityUnread ||
(!nativeProjection &&
countUnreadHighPriorityObservedEvents(
observedEvents,
readAtForObservedEvent,
) > 0)
) {
// Non-DM: high-priority only if at least one mention/broadcast
// remains unread in its own channel/thread context.
@@ -0,0 +1,73 @@
import { invokeTauri } from "@/shared/api/tauri";
import type { ObservedUnreadEvent } from "@/features/channels/unreadChannelCounts";
export type ObservedUnreadScope = { pubkey: string; relayUrl: string };
export type ObservedUnreadProjection = {
channelId: string;
latest: number;
count: number;
badgeCount: number;
appBadgeCount: number;
topLevelUnread: boolean;
highPriorityUnread: boolean;
};
export type ObservedUnreadMembershipSeed = {
participatedRootIds: string[];
authoredRootIds: string[];
mentionedRootIds: string[];
followedRootIds: string[];
mutedRootIds: string[];
mutedChannelIds: string[];
};
export type ObservedUnreadWireEvent = ObservedUnreadEvent & {
channelId: string;
};
export type ObservedUnreadResponse =
| {
kind: "snapshot";
scope: ObservedUnreadScope;
generation: string;
revision: number;
lastAckedSequence: number;
migrationComplete: boolean;
membershipSeeded: boolean;
channels: ObservedUnreadProjection[];
}
| {
kind: "delta";
scope: ObservedUnreadScope;
generation: string;
baseRevision: number;
revision: number;
ackedSequence: number;
upserts: ObservedUnreadProjection[];
removed: string[];
}
| {
kind: "snapshotRequired";
scope: ObservedUnreadScope;
generation: string;
revision: number;
lastAckedSequence: number;
};
export function openObservedUnreadScope(request: {
scope: ObservedUnreadScope;
legacyPayload: unknown | null;
membershipSeed: ObservedUnreadMembershipSeed | null;
}): Promise<ObservedUnreadResponse> {
return invokeTauri("observed_unread_open_scope", { request });
}
export function ingestObservedUnread(request: {
scope: ObservedUnreadScope;
sequence: number;
baseRevision: number;
events: ObservedUnreadWireEvent[];
markers: Array<{ contextId: string; readAt: number | null }>;
membership: Array<{ kind: string; value: string; present: boolean }>;
clearChannels: string[];
clearAll: boolean;
}): Promise<ObservedUnreadResponse> {
return invokeTauri("observed_unread_ingest", { request });
}
@@ -12,11 +12,6 @@ export type UnreadCatchUpChannel = {
export type UnreadCatchUpRequest = {
channels: UnreadCatchUpChannel[];
selfPubkey: string;
participatedRootIds: string[];
authoredRootIds: string[];
mentionedRootIds: string[];
followedRootIds: string[];
mutedRootIds: string[];
mutedChannelIds: string[];
};
+34
View File
@@ -13377,6 +13377,40 @@ export function maybeInstallE2eTauriMocks() {
case "start_archive_sync":
case "stop_archive_sync":
return null;
case "observed_unread_open_scope": {
const request = payload as {
request: { scope: { pubkey: string; relayUrl: string } };
};
return {
kind: "snapshot",
scope: request.request.scope,
generation: "e2e",
revision: 0,
lastAckedSequence: 0,
migrationComplete: true,
membershipSeeded: true,
channels: [],
};
}
case "observed_unread_ingest": {
const request = payload as {
request: {
scope: { pubkey: string; relayUrl: string };
sequence: number;
baseRevision: number;
};
};
return {
kind: "delta",
scope: request.request.scope,
generation: "e2e",
baseRevision: request.request.baseRevision,
revision: request.request.baseRevision + 1,
ackedSequence: request.request.sequence,
upserts: [],
removed: [],
};
}
case "unread_catch_up": {
const request = payload as {
request: { channels: Array<{ id: string }> };