diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index e578e8277..6e3d8399c 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -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, diff --git a/desktop/src-tauri/src/observed_unread.rs b/desktop/src-tauri/src/observed_unread.rs new file mode 100644 index 000000000..45bbbbbcd --- /dev/null +++ b/desktop/src-tauri/src/observed_unread.rs @@ -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, + 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, +} + +#[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, + authored_root_ids: Vec, + mentioned_root_ids: Vec, + followed_root_ids: Vec, + muted_root_ids: Vec, + muted_channel_ids: Vec, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct OpenScopeRequest { + scope: ObservedUnreadScope, + legacy_payload: Option, + membership_seed: Option, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct IngestRequest { + scope: ObservedUnreadScope, + sequence: u64, + base_revision: u64, + events: Vec, + markers: Vec, + membership: Vec, + clear_channels: Vec, + 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, + }, + Delta { + scope: ObservedUnreadScope, + generation: String, + base_revision: u64, + revision: u64, + acked_sequence: u64, + upserts: Vec, + removed: Vec, + }, + SnapshotRequired { + scope: ObservedUnreadScope, + generation: String, + revision: u64, + last_acked_sequence: u64, + }, +} + +fn db_path(app: &AppHandle) -> Result { + 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 { + 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 { + 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, key: &str) -> u64 { + markers.get(key).copied().unwrap_or(0) +} + +fn projections(tx: &Transaction<'_>, scope: &str) -> Result, 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 = marker_stmt + .query_map([scope], |r| Ok((r.get(0)?, r.get(1)?))) + .map_err(|e| format!("query unread markers: {e}"))? + .collect::>() + .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>(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 = 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 { + 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 { + 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>, 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> = 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); + } +} diff --git a/desktop/src-tauri/src/shutdown.rs b/desktop/src-tauri/src/shutdown.rs index efd88f3ca..17ca7a7bb 100644 --- a/desktop/src-tauri/src/shutdown.rs +++ b/desktop/src-tauri/src/shutdown.rs @@ -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::().prevent_sleep); + crate::observed_unread::flush(app); app.state::() .shutdown_all(); if let Err(error) = shutdown_managed_agents(app) { diff --git a/desktop/src-tauri/src/unread_catch_up.rs b/desktop/src-tauri/src/unread_catch_up.rs index c7de5dae8..5e31df1e6 100644 --- a/desktop/src-tauri/src/unread_catch_up.rs +++ b/desktop/src-tauri/src/unread_catch_up.rs @@ -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, self_pubkey: String, - participated_root_ids: HashSet, - authored_root_ids: HashSet, - mentioned_root_ids: HashSet, - followed_root_ids: HashSet, - muted_root_ids: HashSet, muted_channel_ids: HashSet, } @@ -142,6 +137,7 @@ pub(crate) async fn unread_catch_up( request: UnreadCatchUpRequest, state: State<'_, AppState>, relay_client: State<'_, NativeRelayClient>, + app: AppHandle, ) -> Result { 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, + membership: &std::collections::HashMap>, ) -> Vec { 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>, participated: &HashSet, authored: &HashSet, ) -> 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], 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, diff --git a/desktop/src/features/channels/useObservedUnreadPersistence.ts b/desktop/src/features/channels/useObservedUnreadPersistence.ts index c1a94570b..6bde48c3f 100644 --- a/desktop/src/features/channels/useObservedUnreadPersistence.ts +++ b/desktop/src/features/channels/useObservedUnreadPersistence.ts @@ -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; - /** 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>; + 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> >, latestByChannelRef: React.MutableRefObject>, - 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(""); + const scopeLoadedRef = React.useRef(""); + const projectionsRef = React.useRef( + new Map(), + ); + const nativeRef = React.useRef(null); + const nativeFailedRef = React.useRef(false); const timerRef = React.useRef | null>(null); - + const queueRef = React.useRef([]); + const chainRef = React.useRef(Promise.resolve()); const persistRefs = React.useRef({ 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(); + 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, + ], ); } diff --git a/desktop/src/features/channels/useUnreadChannels.ts b/desktop/src/features/channels/useUnreadChannels.ts index 34b2d9450..8a8743525 100644 --- a/desktop/src/features/channels/useUnreadChannels.ts +++ b/desktop/src/features/channels/useUnreadChannels.ts @@ -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 = 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>(new Set()); - mutedChannelIdsRef.current = mutedChannelIdsOption ?? new Set(); + const mutedChannelIdsRef = React.useRef>(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()); + 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. diff --git a/desktop/src/shared/api/tauriObservedUnread.ts b/desktop/src/shared/api/tauriObservedUnread.ts new file mode 100644 index 000000000..7a02ece72 --- /dev/null +++ b/desktop/src/shared/api/tauriObservedUnread.ts @@ -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 { + 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 { + return invokeTauri("observed_unread_ingest", { request }); +} diff --git a/desktop/src/shared/api/tauriUnreadCatchUp.ts b/desktop/src/shared/api/tauriUnreadCatchUp.ts index 0f0c125b0..445fe465e 100644 --- a/desktop/src/shared/api/tauriUnreadCatchUp.ts +++ b/desktop/src/shared/api/tauriUnreadCatchUp.ts @@ -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[]; }; diff --git a/desktop/src/testing/e2eBridge.ts b/desktop/src/testing/e2eBridge.ts index 8a31d2ad2..b27bf91ec 100644 --- a/desktop/src/testing/e2eBridge.ts +++ b/desktop/src/testing/e2eBridge.ts @@ -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 }> };