From 5d82f93bc3e41e4e879a238b066b629f7f014d6a Mon Sep 17 00:00:00 2001 From: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 Date: Wed, 8 Jul 2026 10:15:59 -0400 Subject: [PATCH] =?UTF-8?q?fix(desktop):=20scope=20observer=20feed=20by=20?= =?UTF-8?q?channel=20and=20add=20session-boundary=20dividers=20Two=20obser?= =?UTF-8?q?ver-feed=20defects=20fixed=20in=20a=20single=20PR:=20Slice=201?= =?UTF-8?q?=20=E2=80=94=20channel-scoped=20archive=20(index=20at=20ingest?= =?UTF-8?q?=20+=20idempotent=20backfill)=20The=20archive=20read=20path=20q?= =?UTF-8?q?ueried=20all=20owner=5Fp=20kind=2024200=20rows=20for=20an=20ide?= =?UTF-8?q?ntity,=20then=20relied=20on=20a=20display-time=20scopeByChannel?= =?UTF-8?q?=20filter.=20When=20the=20filter=20was=20missing=20or=20skipped?= =?UTF-8?q?,=20frames=20from=20other=20channels=20leaked=20into=20the=20cu?= =?UTF-8?q?rrent=20channel's=20observer=20feed.=20Fix:=20add=20an=20observ?= =?UTF-8?q?er=5Fchannel=5Findex=20table=20that=20maps=20(identity,=20relay?= =?UTF-8?q?,=20event=5Fid)=20to=20channel=5Fid.=20Three=20new=20Tauri=20co?= =?UTF-8?q?mmands=20drive=20this:=20-=20read=5Farchived=5Fobserver=5Fevent?= =?UTF-8?q?s=5Ffor=5Fchannel:=20channel-scoped=20paginated=20read=20via=20?= =?UTF-8?q?=20=20the=20index=20(replaces=20the=20raw=20owner=5Fp=20scan=20?= =?UTF-8?q?in=20useLoadArchivedObserverEvents)=20-=20index=5Fobserver=5Fch?= =?UTF-8?q?annel=5Fid:=20TS=20calls=20this=20after=20decrypting=20events?= =?UTF-8?q?=20to=20write=20=20=20index=20entries=20-=20read=5Funindexed=5F?= =?UTF-8?q?observer=5Frows:=20returns=20all=20owner=5Fp=20kind=2024200=20r?= =?UTF-8?q?ows=20not=20yet=20=20=20indexed,=20for=20the=20one-shot=20idemp?= =?UTF-8?q?otent=20backfill=20Backfill=20strategy=20(choice=20i,=20backfil?= =?UTF-8?q?l-before-read):=20a=20React.useEffect=20in=20useLoadArchivedObs?= =?UTF-8?q?erverEvents=20runs=20once=20on=20mount,=20reads=20all=20unindex?= =?UTF-8?q?ed=20rows,=20decrypts=20each,=20and=20batch-writes=20index=20en?= =?UTF-8?q?tries.=20After=20backfill=20the=20channel=20index=20is=20author?= =?UTF-8?q?itative=20=E2=80=94=20the=20scoped=20read=20pages=20it=20direct?= =?UTF-8?q?ly=20with=20no=20zero-match=20raw-page=20guard=20needed.=20Null?= =?UTF-8?q?/decrypt-failed=20channelId=20rows=20stay=20unscoped=20and=20ar?= =?UTF-8?q?e=20hidden=20from=20every=20scoped=20channel=20view=20(Will's?= =?UTF-8?q?=20ruling:=20option=20(a)=20=E2=80=94=20display=20decision=20on?= =?UTF-8?q?ly,=20rows=20stay=20on=20disk).=20Slice=202=20=E2=80=94=20sessi?= =?UTF-8?q?on-boundary=20rendering=20Items=20rendered=20flat=20with=20no?= =?UTF-8?q?=20session=20grouping;=20archived=20and=20live=20frames=20inter?= =?UTF-8?q?leaved=20without=20any=20visual=20marker=20of=20where=20one=20a?= =?UTF-8?q?gent=20session=20ended=20and=20another=20began,=20causing=20use?= =?UTF-8?q?rs=20to=20think=20archived=20historical=20context=20is=20still?= =?UTF-8?q?=20live.=20Fix:=20split=20TranscriptItem[]=20into=20contiguous?= =?UTF-8?q?=20session=20runs=20before=20calling=20buildTranscriptDisplayBl?= =?UTF-8?q?ocks.=20This=20also=20fixes=20the=20pendingSystemPrompt=20singl?= =?UTF-8?q?e-slot=20clobber=20(each=20run=20gets=20its=20own=20lifecycle).?= =?UTF-8?q?=20A=20new=20session-boundary=20display=20block=20kind=20is=20i?= =?UTF-8?q?nterleaved=20between=20runs;=20SessionBoundaryDivider=20renders?= =?UTF-8?q?=20a=20horizontal=20rule=20with=20session=20start=20time.=20No?= =?UTF-8?q?=20boundary=20emitted=20for=20single-session=20transcripts.=20L?= =?UTF-8?q?iveness=20=E2=80=94=20"current=20session"=20gating=20Track=20la?= =?UTF-8?q?test-live-session-id=20per=20(normalized=20agent=20pubkey,=20ch?= =?UTF-8?q?annelId)=20in=20observerRelayStore,=20set=20on=20the=20handleRe?= =?UTF-8?q?layObserverEvent=20path,=20cleared=20in=20resetAgentObserverSto?= =?UTF-8?q?re.=20Label=20a=20boundary=20"current=20session"=20only=20when?= =?UTF-8?q?=20it=20is=20both=20the=20newest-visible=20session=20AND=20equa?= =?UTF-8?q?ls=20that=20latest-live=20id;=20otherwise=20label=20it=20"most?= =?UTF-8?q?=20recent=20observed=20session".?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- desktop/src-tauri/src/archive/mod.rs | 114 +++++++++ desktop/src-tauri/src/archive/pipeline.rs | 20 ++ desktop/src-tauri/src/archive/store.rs | 205 ++++++++++++++++ desktop/src-tauri/src/archive/store_tests.rs | 122 ++++++++++ desktop/src-tauri/src/lib.rs | 3 + .../src/features/agents/observerRelayStore.ts | 75 ++++++ .../agents/ui/AgentSessionTranscriptList.tsx | 84 ++++++- .../agentSessionTranscriptGrouping.test.mjs | 221 ++++++++++++++++++ .../ui/agentSessionTranscriptGrouping.ts | 130 ++++++++++- .../features/agents/ui/useObserverEvents.ts | 149 +++++++++++- .../channels/ui/AgentSessionThreadPanel.tsx | 2 +- desktop/src/shared/api/tauriArchive.ts | 86 +++++++ 12 files changed, 1194 insertions(+), 17 deletions(-) diff --git a/desktop/src-tauri/src/archive/mod.rs b/desktop/src-tauri/src/archive/mod.rs index 684910ea3..a65b126a4 100644 --- a/desktop/src-tauri/src/archive/mod.rs +++ b/desktop/src-tauri/src/archive/mod.rs @@ -502,6 +502,120 @@ pub fn delete_save_subscription( // ── read_archived_events ───────────────────────────────────────────────────── +// ── read_archived_observer_events_for_channel ──────────────────────────────── + +/// Read a paginated page of archived kind 24200 events scoped to one channel, +/// using the `observer_channel_index` as the primary lookup. +/// +/// The index only contains rows with a known, non-null `channelId` (frames +/// pre-dating the channelId stamp, or rows where decryption failed, are +/// absent from the index and thus absent from every scoped channel view — +/// Will's (a) ruling, 2026-07-08). +/// +/// Returns at most `limit` events (default `DEFAULT_READ_LIMIT`) in +/// newest-first order. Compound cursor `(before_created_at, before_id)` works +/// identically to `read_archived_events`. +#[tauri::command] +pub fn read_archived_observer_events_for_channel( + state: State<'_, AppState>, + channel_id: String, + before_created_at: Option, + before_id: Option, + limit: Option, +) -> Result, String> { + let identity_pk = identity_pubkey(&state)?; + let relay_url = relay_ws_url_with_override(&state); + let conn = open_db()?; + store::read_archived_observer_events_for_channel( + &conn, + &identity_pk, + &relay_url, + &channel_id, + before_created_at, + before_id.as_deref(), + limit.unwrap_or(DEFAULT_READ_LIMIT), + ) +} + +// ── index_observer_channel_id ───────────────────────────────────────────────── + +/// Index one or more archived observer frame ids with their decoded channelId. +/// +/// Called from the TS-side backfill after attempting `decryptObserverEvent`. +/// `channel_id` is `Some` for frames with a non-null channelId; `None` for +/// frames where decryption yielded no channelId or failed entirely. Both cases +/// write a status row so a re-run skips the frame (INSERT OR IGNORE on PK). +/// +/// Idempotent: rows that are already indexed are left unchanged. +#[tauri::command] +pub fn index_observer_channel_id( + state: State<'_, AppState>, + entries: Vec, +) -> Result<(), String> { + let identity_pk = identity_pubkey(&state)?; + let relay_url = relay_ws_url_with_override(&state); + let conn = open_db()?; + for entry in &entries { + store::upsert_observer_channel_index( + &conn, + &identity_pk, + &relay_url, + &entry.event_id, + entry.channel_id.as_deref(), + entry.created_at, + )?; + } + Ok(()) +} + +/// A single (event_id, channel_id?, created_at) record used by +/// `index_observer_channel_id`. +/// +/// `channel_id` is `None` for frames where decryption found no channelId or +/// failed; those rows are written to `observer_channel_index` with a NULL +/// channel_id so the frame is treated as processed (no re-decrypt on re-run). +#[derive(Debug, Deserialize)] +pub struct ObserverChannelIndexEntry { + pub event_id: String, + pub channel_id: Option, + pub created_at: i64, +} + +// ── read_unindexed_observer_rows ───────────────────────────────────────────── + +/// Return raw event JSON + id + created_at for all `owner_p` kind 24200 rows +/// that are NOT yet in `observer_channel_index`. +/// +/// The TS-side backfill driver calls this once, decrypts each row, and sends +/// the (id, channelId, created_at) triples back via `index_observer_channel_id`. +/// Together these constitute the one-shot idempotent backfill required by the +/// Slice 1 acceptance criteria (Thufir Pass 4). +#[tauri::command] +pub fn read_unindexed_observer_rows( + state: State<'_, AppState>, +) -> Result, String> { + let identity_pk = identity_pubkey(&state)?; + let relay_url = relay_ws_url_with_override(&state); + let conn = open_db()?; + let rows = store::read_unindexed_observer_rows(&conn, &identity_pk, &relay_url)?; + Ok(rows + .into_iter() + .map(|(id, raw_json, created_at)| RawObserverRow { + id, + raw_json, + created_at, + }) + .collect()) +} + +/// Wire type returned by `read_unindexed_observer_rows`. +#[derive(Debug, Serialize)] +pub struct RawObserverRow { + pub id: String, + pub raw_json: String, + pub created_at: i64, +} + /// Default page size for `read_archived_events`. const DEFAULT_READ_LIMIT: i64 = 50; diff --git a/desktop/src-tauri/src/archive/pipeline.rs b/desktop/src-tauri/src/archive/pipeline.rs index ec3de94cc..2bd149dce 100644 --- a/desktop/src-tauri/src/archive/pipeline.rs +++ b/desktop/src-tauri/src/archive/pipeline.rs @@ -423,6 +423,26 @@ pub(super) fn commit_archive( &p.matched_scope.scope_value, now, )?; + + // Index at ingest: attempt to decrypt and extract channelId. + // Write a status row regardless of outcome so backfill never + // re-processes this frame (INSERT OR IGNORE on PK is a no-op if + // the row is already present from a prior run). + let channel_id_for_index: Option = + buzz_core_pkg::observer::decrypt_observer_payload::( + owner_keys, &p.event, + ) + .ok() + .and_then(|v| v.get("channelId")?.as_str().map(|s| s.to_owned())); + store::upsert_observer_channel_index( + &tx, + identity_pk, + relay_url, + eid, + channel_id_for_index.as_deref(), + p.event.created_at.as_secs() as i64, + )?; + persisted += 1; } diff --git a/desktop/src-tauri/src/archive/store.rs b/desktop/src-tauri/src/archive/store.rs index 083f9dc43..24d21c0b1 100644 --- a/desktop/src-tauri/src/archive/store.rs +++ b/desktop/src-tauri/src/archive/store.rs @@ -46,8 +46,62 @@ CREATE TABLE IF NOT EXISTS save_subscriptions ( created_at INTEGER NOT NULL, PRIMARY KEY (identity_pubkey, relay_url, scope_type, scope_value) ); + +-- Processing-status index for kind 24200 (observer) frames. +-- +-- Every examined observer row — whether its channelId was successfully +-- decrypted or not — gets exactly one row here keyed by (identity, relay, id). +-- `channel_id` is nullable: +-- NOT NULL → frame was attributed to a channel; visible in that channel's feed. +-- NULL → frame was examined but had no channelId (pre-stamp or decrypt +-- failure); excluded from all scoped views per Will's (a) ruling +-- (2026-07-08), but the row exists so a re-run is a no-op. +-- +-- Scoped reads filter `channel_id = ?`; NULL rows never match, staying hidden. +CREATE TABLE IF NOT EXISTS observer_channel_index ( + identity_pubkey TEXT NOT NULL, + relay_url TEXT NOT NULL, + id TEXT NOT NULL, + channel_id TEXT, + created_at INTEGER NOT NULL, + PRIMARY KEY (identity_pubkey, relay_url, id) +); +CREATE INDEX IF NOT EXISTS idx_observer_channel + ON observer_channel_index (identity_pubkey, relay_url, channel_id, created_at DESC, id DESC); + +-- One-row migration state table: tracks which idempotent migrations have run. +CREATE TABLE IF NOT EXISTS archive_migrations ( + name TEXT PRIMARY KEY, + applied_at INTEGER NOT NULL +); "; +// ── Migration helpers ──────────────────────────────────────────────────────── + +/// Returns true if a named migration has already been applied. +pub fn migration_applied(conn: &Connection, name: &str) -> Result { + // The table is created in SCHEMA above, so it always exists after open. + let count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM archive_migrations WHERE name = ?1", + params![name], + |row| row.get(0), + ) + .map_err(|e| format!("migration_applied query: {e}"))?; + Ok(count > 0) +} + +/// Mark a named migration as applied (idempotent: no-op if already recorded). +pub fn mark_migration_applied(conn: &Connection, name: &str, now: i64) -> Result<(), String> { + conn.execute( + "INSERT INTO archive_migrations (name, applied_at) VALUES (?1, ?2) + ON CONFLICT (name) DO NOTHING", + params![name, now], + ) + .map_err(|e| format!("mark_migration_applied: {e}"))?; + Ok(()) +} + // ── Open / init ───────────────────────────────────────────────────────────── /// Open (or create) the archive database at the given path. @@ -568,6 +622,157 @@ pub fn read_archived_events( .map_err(|e| format!("read read_archived_events row: {e}")) } +/// Upsert a row in `observer_channel_index` for an examined kind 24200 frame. +/// +/// `channel_id` is `Some` for frames with a successfully-decrypted, non-null +/// channelId; `None` for frames that had no channelId or where decryption +/// failed. Null-channel_id rows are excluded from all scoped channel reads +/// (scoped queries filter `channel_id = ?`; NULLs never match), satisfying +/// Will's (a) ruling (2026-07-08), while ensuring the row exists so a re-run +/// is a no-op. +/// +/// Idempotent: if the (identity, relay, id) PK already exists the row is left +/// unchanged (INSERT OR IGNORE). Called both at ingest time (new frames) and +/// from the one-shot backfill for existing archived rows. +pub fn upsert_observer_channel_index( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + event_id: &str, + channel_id: Option<&str>, + created_at: i64, +) -> Result<(), String> { + conn.execute( + "INSERT OR IGNORE INTO observer_channel_index + (identity_pubkey, relay_url, id, channel_id, created_at) + VALUES (?1, ?2, ?3, ?4, ?5)", + params![identity_pubkey, relay_url, event_id, channel_id, created_at], + ) + .map_err(|e| format!("failed to upsert observer_channel_index: {e}"))?; + Ok(()) +} + +/// Read all `owner_p` kind 24200 archived event rows (id + raw_json + +/// created_at) that have NOT yet been processed (i.e., have no row in +/// `observer_channel_index`, regardless of what channel_id would be). +/// +/// Used by the one-shot backfill migration to discover which existing rows +/// still need channel attribution attempted. A row is "processed" as soon as +/// we attempt decryption — whether it succeeds (non-null channel_id written) +/// or fails (null channel_id written) — so re-runs skip those rows. +pub fn read_unindexed_observer_rows( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, +) -> Result, String> { + let mut stmt = conn + .prepare( + "SELECT ae.id, ae.raw_json, ae.created_at + FROM archived_events ae + INNER JOIN archived_event_scopes aes + ON aes.identity_pubkey = ae.identity_pubkey + AND aes.relay_url = ae.relay_url + AND aes.id = ae.id + WHERE ae.identity_pubkey = ?1 + AND ae.relay_url = ?2 + AND ae.kind = 24200 + AND aes.scope_type = 'owner_p' + AND ae.id NOT IN ( + SELECT id FROM observer_channel_index + WHERE identity_pubkey = ?1 + AND relay_url = ?2 + ) + ORDER BY ae.created_at DESC, ae.id DESC", + ) + .map_err(|e| format!("prepare read_unindexed_observer_rows: {e}"))?; + + let rows = stmt + .query_map(params![identity_pubkey, relay_url], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, i64>(2)?, + )) + }) + .map_err(|e| format!("query read_unindexed_observer_rows: {e}"))?; + + rows.collect::, _>>() + .map_err(|e| format!("read read_unindexed_observer_rows row: {e}")) +} + +/// Read a paginated page of archived kind 24200 events scoped to one channel, +/// using the `observer_channel_index` as the primary lookup. +/// +/// Returns `raw_json` strings in newest-first order. The compound cursor +/// `(before_created_at, before_id)` works identically to `read_archived_events`. +/// A page shorter than `limit` signals the archive is exhausted for this channel. +#[allow(clippy::too_many_arguments)] +pub fn read_archived_observer_events_for_channel( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + channel_id: &str, + before_created_at: Option, + before_id: Option<&str>, + limit: i64, +) -> Result, String> { + let mut next_slot: usize = 4; + let mut extra_clauses = String::new(); + let mut before_at_val: Option = None; + let mut before_id_val: Option = None; + + if let (Some(bat), Some(bid)) = (before_created_at, before_id) { + before_at_val = Some(bat); + before_id_val = Some(bid.to_owned()); + extra_clauses.push_str(&format!( + " AND (oci.created_at < ?{next_slot} \ + OR (oci.created_at = ?{next_slot} AND oci.id < ?{}))", + next_slot + 1, + )); + next_slot += 2; + } + let limit_slot = next_slot; + + let sql = format!( + "SELECT ae.raw_json \ + FROM observer_channel_index oci \ + INNER JOIN archived_events ae \ + ON ae.identity_pubkey = oci.identity_pubkey \ + AND ae.relay_url = oci.relay_url \ + AND ae.id = oci.id \ + WHERE oci.identity_pubkey = ?1 \ + AND oci.relay_url = ?2 \ + AND oci.channel_id = ?3\ + {extra_clauses}\ + ORDER BY oci.created_at DESC, oci.id DESC \ + LIMIT ?{limit_slot}", + ); + + let mut params_vec: Vec> = vec![ + Box::new(identity_pubkey.to_owned()), + Box::new(relay_url.to_owned()), + Box::new(channel_id.to_owned()), + ]; + if let (Some(bat), Some(bid)) = (before_at_val, before_id_val) { + params_vec.push(Box::new(bat)); + params_vec.push(Box::new(bid)); + } + params_vec.push(Box::new(limit)); + + let param_refs: Vec<&dyn rusqlite::ToSql> = params_vec.iter().map(|p| p.as_ref()).collect(); + + let mut stmt = conn + .prepare(&sql) + .map_err(|e| format!("prepare read_archived_observer_events_for_channel: {e}"))?; + + let rows = stmt + .query_map(param_refs.as_slice(), |row| row.get::<_, String>(0)) + .map_err(|e| format!("query read_archived_observer_events_for_channel: {e}"))?; + + rows.collect::, _>>() + .map_err(|e| format!("read read_archived_observer_events_for_channel row: {e}")) +} + /// GC: delete orphaned event rows whose last scope row was just removed. /// /// Called after any batch deletion of scope rows. Uses a LEFT JOIN so only diff --git a/desktop/src-tauri/src/archive/store_tests.rs b/desktop/src-tauri/src/archive/store_tests.rs index 18a80af60..0a07811e7 100644 --- a/desktop/src-tauri/src/archive/store_tests.rs +++ b/desktop/src-tauri/src/archive/store_tests.rs @@ -728,3 +728,125 @@ fn test_read_archived_events_same_second_cursor_no_skip() { // No overlap with page1. assert!(page2.iter().all(|r| !r.contains("\"z\""))); } + +// ── observer_channel_index ─────────────────────────────────────────────────── + +#[test] +fn test_upsert_observer_channel_index_non_null_is_idempotent() { + let conn = in_memory(); + let pk = "owner"; + let relay = "wss://r"; + + upsert_observer_channel_index(&conn, pk, relay, "ev1", Some("ch-abc"), 1000).unwrap(); + // Second call with same PK → INSERT OR IGNORE, no error, row unchanged. + upsert_observer_channel_index(&conn, pk, relay, "ev1", Some("ch-abc"), 2000).unwrap(); + + // Only one row should exist and created_at stays at 1000 (first write wins). + let count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM observer_channel_index WHERE id = 'ev1'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(count, 1, "INSERT OR IGNORE must not duplicate rows"); + + let at: i64 = conn + .query_row( + "SELECT created_at FROM observer_channel_index WHERE id = 'ev1'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(at, 1000, "first write's created_at must be preserved"); +} + +#[test] +fn test_upsert_observer_channel_index_null_channel_id_is_idempotent() { + let conn = in_memory(); + let pk = "owner"; + let relay = "wss://r"; + + // Write a NULL channel_id status row (unscoped / decrypt-failed frame). + upsert_observer_channel_index(&conn, pk, relay, "ev-null", None, 500).unwrap(); + // Re-run → no-op. + upsert_observer_channel_index(&conn, pk, relay, "ev-null", None, 600).unwrap(); + + let count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM observer_channel_index WHERE id = 'ev-null'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!( + count, 1, + "null-channel_id row must not be duplicated on re-run" + ); + + let channel_id: Option = conn + .query_row( + "SELECT channel_id FROM observer_channel_index WHERE id = 'ev-null'", + [], + |r| r.get(0), + ) + .unwrap(); + assert!( + channel_id.is_none(), + "channel_id must be NULL for unscoped frames" + ); +} + +#[test] +fn test_read_archived_observer_excludes_null_channel_id_rows() { + // Scoped reads filter channel_id = '?'; NULL rows must never appear. + let conn = in_memory(); + let pk = "owner"; + let relay = "wss://r"; + + // Seed an archived event row. + upsert_archived_event(&conn, pk, relay, "ev-null", 24200, "agent", 1000, "{}", 0).unwrap(); + upsert_event_scope(&conn, pk, relay, "ev-null", "owner_p", pk, 0).unwrap(); + + // Index it with a NULL channel_id. + upsert_observer_channel_index(&conn, pk, relay, "ev-null", None, 1000).unwrap(); + + // Scoped read for ANY channel must return nothing (NULL != 'some-channel'). + let rows = + read_archived_observer_events_for_channel(&conn, pk, relay, "some-channel", None, None, 50) + .unwrap(); + assert!( + rows.is_empty(), + "scoped read must not return null-channel_id rows" + ); +} + +#[test] +fn test_read_unindexed_observer_rows_excludes_processed_rows() { + // After a NULL channel_id row is written, the event must no longer appear + // in read_unindexed_observer_rows (it is "processed"). + let conn = in_memory(); + let pk = "owner"; + let relay = "wss://r"; + + // Seed an archived observer event with owner_p scope. + upsert_archived_event(&conn, pk, relay, "ev-old", 24200, "agent", 1000, "{}", 0).unwrap(); + upsert_event_scope(&conn, pk, relay, "ev-old", "owner_p", pk, 0).unwrap(); + + // Before indexing: ev-old must appear in unindexed rows. + let unindexed = read_unindexed_observer_rows(&conn, pk, relay).unwrap(); + assert!( + unindexed.iter().any(|(id, _, _)| id == "ev-old"), + "ev-old must appear in unindexed rows before indexing" + ); + + // Index it with NULL channel_id (decrypt-failed / unscoped). + upsert_observer_channel_index(&conn, pk, relay, "ev-old", None, 1000).unwrap(); + + // After indexing: ev-old must NOT appear in unindexed rows (re-run is a no-op). + let unindexed2 = read_unindexed_observer_rows(&conn, pk, relay).unwrap(); + assert!( + !unindexed2.iter().any(|(id, _, _)| id == "ev-old"), + "ev-old must be excluded from unindexed rows after null-channel_id indexing" + ); +} diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index 6a550b21e..3249aa0c4 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -635,6 +635,9 @@ pub fn run() { archive::list_save_subscriptions, archive::delete_save_subscription, archive::read_archived_events, + archive::read_archived_observer_events_for_channel, + archive::index_observer_channel_id, + archive::read_unindexed_observer_rows, is_auto_update_supported, ]) .build(tauri::generate_context!()) diff --git a/desktop/src/features/agents/observerRelayStore.ts b/desktop/src/features/agents/observerRelayStore.ts index adf6ea92d..7fdcb9504 100644 --- a/desktop/src/features/agents/observerRelayStore.ts +++ b/desktop/src/features/agents/observerRelayStore.ts @@ -41,6 +41,41 @@ const eventsByAgent = new Map(); const transcriptByAgent = new Map(); const snapshotByAgent = new Map(); +// Per-agent, per-channel latest-live-session-id. +// Key: `${normalizePubkey(agentPubkey)}:${channelId}`. +// Set when a live relay observer event with a sessionId arrives. +// Cleared in resetAgentObserverStore. +// +// "Latest-live" means: the sessionId that most recently appeared via the +// live relay path (handleRelayObserverEvent). It is NOT derived from +// connectionState or an ever-live Set — an ever-live Set would incorrectly +// mark session A as "current" after session B has started (Thufir Pass 3). +// +// Stored as `{ sessionId, timestamp, seq }` so that late-arriving live frames +// from an older session never regress the latest-live id. We only advance when +// the parsed event sorts strictly AFTER the stored one, using the same +// two-key ordering as `compareObserverEvents`: timestamp first, then seq on a +// tie — so a higher-seq frame at equal timestamp still advances the entry. +type LatestLiveEntry = { sessionId: string; timestamp: string; seq: number }; +const latestLiveSessionByAgentChannel = new Map(); + +function liveSessionKey(agentPubkey: string, channelId: string | null): string { + return `${normalizePubkey(agentPubkey)}:${channelId ?? ""}`; +} + +/** Read the latest-live-session-id for a (agent, channel) pair. */ +export function getLatestLiveSessionId( + agentPubkey: string | null | undefined, + channelId: string | null | undefined, +): string | null { + if (!agentPubkey) return null; + return ( + latestLiveSessionByAgentChannel.get( + liveSessionKey(agentPubkey, channelId ?? null), + )?.sessionId ?? null + ); +} + // Per-agent listeners for `control_result` frames. The ModelPicker subscribes // here to learn the async outcome of a `switch_model` frame (the send is // fire-and-forget; the harness replies out-of-band over the observer relay). @@ -184,6 +219,26 @@ export function compareObserverEvents( return left.seq - right.seq; } +/** + * Returns true if `candidate` sorts strictly after `stored` using the same + * two-key ordering as `compareObserverEvents`: later timestamp wins; equal + * timestamp falls back to higher seq. Extracted so latest-live advancement + * cannot drift from transcript ordering. + */ +export function isObserverEventAfter( + candidate: { timestamp: string; seq: number }, + stored: { timestamp: string; seq: number }, +): boolean { + const candidateTime = Date.parse(candidate.timestamp); + const storedTime = Date.parse(stored.timestamp); + if (Number.isFinite(candidateTime) && Number.isFinite(storedTime)) { + if (candidateTime !== storedTime) { + return candidateTime > storedTime; + } + } + return candidate.seq > stored.seq; +} + async function handleRelayObserverEvent( event: RelayEvent, activeGeneration: number, @@ -211,6 +266,25 @@ async function handleRelayObserverEvent( if (activeGeneration !== generation) { return; } + // Track the latest-live-session-id per (agent, channel) on the live path. + // Only set when the parsed event carries both a sessionId and channelId, + // so we never attribute a session to the wrong channel. + if (parsed.sessionId && parsed.channelId) { + const key = liveSessionKey(agentPubkey, parsed.channelId); + const stored = latestLiveSessionByAgentChannel.get(key); + // Advance only when this event sorts strictly AFTER the stored one via + // isObserverEventAfter (timestamp then seq — same ordering as + // compareObserverEvents). This prevents late-arriving live frames from + // older sessions from regressing the latest-live id, while also + // correctly advancing on a same-timestamp frame with a higher seq. + if (!stored || isObserverEventAfter(parsed, stored)) { + latestLiveSessionByAgentChannel.set(key, { + sessionId: parsed.sessionId, + timestamp: parsed.timestamp, + seq: parsed.seq, + }); + } + } appendAgentEvent(agentPubkey, parsed); if (parsed.kind === "session_config_captured") { void putAgentSessionConfig(agentPubkey, parsed.payload); @@ -513,6 +587,7 @@ export function resetAgentObserverStore() { snapshotByAgent.clear(); knownAgentPubkeys.clear(); knownAgentsBySubscription.clear(); + latestLiveSessionByAgentChannel.clear(); onSessionConfigCaptured = null; connectionState = "idle"; errorMessage = null; diff --git a/desktop/src/features/agents/ui/AgentSessionTranscriptList.tsx b/desktop/src/features/agents/ui/AgentSessionTranscriptList.tsx index 00eecd0d0..38143aaac 100644 --- a/desktop/src/features/agents/ui/AgentSessionTranscriptList.tsx +++ b/desktop/src/features/agents/ui/AgentSessionTranscriptList.tsx @@ -1,11 +1,15 @@ import * as React from "react"; import { motion, useReducedMotion } from "motion/react"; -import { CheckCheck, Radio } from "lucide-react"; +import { CheckCheck, Clock, Radio } from "lucide-react"; import { useActiveAgentTurns, type ActiveTurnSummary, } from "@/features/agents/activeAgentTurnsStore"; +import { + subscribeAgentObserverStore, + getLatestLiveSessionId, +} from "@/features/agents/observerRelayStore"; import type { UserProfileLookup } from "@/features/profile/lib/identity"; import { useAnchoredScroll } from "@/features/messages/ui/useAnchoredScroll"; import { cn } from "@/shared/lib/cn"; @@ -134,9 +138,21 @@ export function AgentSessionTranscriptList({ () => isAgentTurnLive(activeTurns, channelId), [activeTurns, channelId], ); + + // Subscribe to the observer relay store so we read the latest-live-session-id + // reactively. We don't need the full snapshot — only the key for boundary labeling. + const getLatestLive = React.useCallback( + () => getLatestLiveSessionId(agentPubkey, channelId), + [agentPubkey, channelId], + ); + const latestLiveSessionId = React.useSyncExternalStore( + subscribeAgentObserverStore, + getLatestLive, + ); + const displayBlocks = React.useMemo( - () => buildTranscriptDisplayBlocks(items), - [items], + () => buildTranscriptDisplayBlocks(items, latestLiveSessionId), + [items, latestLiveSessionId], ); const scrollContainerRef = React.useRef(null); const contentRef = React.useRef(null); @@ -282,6 +298,11 @@ function hasRenderableCompactBlock(block: TranscriptDisplayBlock) { return isRenderableCompactItem(block.item); } + // session-boundary dividers are not renderable content in compact view. + if (block.kind === "session-boundary") { + return false; + } + return block.segments.some((segment) => { if (segment.kind === "item") { return isRenderableCompactItem(segment.item); @@ -316,6 +337,9 @@ function getDisplayBlockKey(block: TranscriptDisplayBlock) { if (block.kind === "single") { return block.item.id; } + if (block.kind === "session-boundary") { + return `session-boundary:${block.sessionId}`; + } return `turn:${block.turnId}`; } @@ -341,6 +365,15 @@ function TranscriptDisplayBlockView({ const hasCompletedInitialRenderRef = useHasCompletedInitialRender(); const animateSegmentEnter = animationPreferenceEnabled && !shouldReduceMotion; + if (block.kind === "session-boundary") { + return ( + + ); + } + if (block.kind === "single") { return ( +
+ + {labelState === "current" ? ( + +
+
+ ); +} + const TranscriptItemView = React.memo(function TranscriptItemView({ agentAvatarUrl, agentName, diff --git a/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.test.mjs b/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.test.mjs index 599cc9244..569895556 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.test.mjs +++ b/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.test.mjs @@ -6,6 +6,7 @@ import { flattenDisplayBlocks, formatTurnSetupLabel, } from "./agentSessionTranscriptGrouping.ts"; +import { isObserverEventAfter } from "../observerRelayStore.ts"; const baseTimestamp = "2026-06-14T22:20:23.000Z"; @@ -639,3 +640,223 @@ function mkTool(id, label, renderClass = "generic", groupKey = label) { channelId: "chan-1", }; } + +// ── Session-run splitting and session-boundary blocks ────────────────────────── + +/** + * Build a minimal tool-call item stamped with a specific session. + */ +function sessionItem(id, sessionId, ts = "2026-07-08T00:00:00.000Z") { + return { + id, + type: "tool", + renderClass: "generic", + descriptor: { + renderClass: "generic", + label: id, + preview: id, + source: "harness", + groupKey: id, + }, + title: id, + toolName: id, + buzzToolName: null, + status: "completed", + args: {}, + result: "", + isError: false, + timestamp: ts, + startedAt: ts, + completedAt: ts, + turnId: `turn-${id}`, + sessionId, + channelId: "chan-1", + }; +} + +// ── Single session — no boundary injected ────────────────────────────────────── + +test("buildTranscriptDisplayBlocks_singleSession_noBoundaryBlock", () => { + const items = [sessionItem("a", "sess-1"), sessionItem("b", "sess-1")]; + const blocks = buildTranscriptDisplayBlocks(items); + const boundaryBlocks = blocks.filter((b) => b.kind === "session-boundary"); + assert.equal( + boundaryBlocks.length, + 0, + "no session-boundary blocks for a single session", + ); +}); + +// ── Two sessions — one boundary between them ─────────────────────────────────── + +test("buildTranscriptDisplayBlocks_twoSessions_oneBoundaryBetween", () => { + // items ordered oldest-first: sess-1 then sess-2 + const items = [ + sessionItem("a", "sess-1", "2026-07-08T00:00:01.000Z"), + sessionItem("b", "sess-2", "2026-07-08T00:00:02.000Z"), + ]; + const blocks = buildTranscriptDisplayBlocks(items); + const boundaryBlocks = blocks.filter((b) => b.kind === "session-boundary"); + assert.equal( + boundaryBlocks.length, + 1, + "exactly one boundary for two sessions", + ); + // The boundary is inserted BEFORE the newer run (sess-2). + const boundaryIndex = blocks.indexOf(boundaryBlocks[0]); + const prevBlock = blocks[boundaryIndex - 1]; + const nextBlock = blocks[boundaryIndex + 1]; + // Previous content block belongs to sess-1 items, next to sess-2. + const flatPrev = flattenDisplayBlocks([prevBlock]).map((i) => i.id); + const flatNext = flattenDisplayBlocks([nextBlock]).map((i) => i.id); + assert.ok(flatPrev.includes("a"), "content before boundary is sess-1"); + assert.ok(flatNext.includes("b"), "content after boundary is sess-2"); +}); + +// ── Newest session labeled correctly relative to latestLiveSessionId ────────── + +test("buildTranscriptDisplayBlocks_newestMatchesLive_labelStateCurrent", () => { + const items = [ + sessionItem("old", "sess-1", "2026-07-08T00:00:01.000Z"), + sessionItem("new", "sess-2", "2026-07-08T00:00:02.000Z"), + ]; + const blocks = buildTranscriptDisplayBlocks(items, "sess-2"); + const boundary = blocks.find((b) => b.kind === "session-boundary"); + assert.ok(boundary, "boundary present"); + assert.equal( + boundary.labelState, + "current", + "labelState=current when newest session matches live id", + ); + assert.equal(boundary.sessionId, "sess-2", "boundary sessionId is sess-2"); +}); + +test("buildTranscriptDisplayBlocks_newestNoLive_labelStateMostRecent", () => { + const items = [ + sessionItem("old", "sess-1", "2026-07-08T00:00:01.000Z"), + sessionItem("new", "sess-2", "2026-07-08T00:00:02.000Z"), + ]; + // latestLiveSessionId is null → newest session is "most-recent" + const blocksNoLive = buildTranscriptDisplayBlocks(items, null); + const boundaryNoLive = blocksNoLive.find( + (b) => b.kind === "session-boundary", + ); + assert.equal( + boundaryNoLive.labelState, + "most-recent", + "labelState=most-recent when no live id (archived-only view)", + ); + + // latestLiveSessionId is a DIFFERENT session → newest is still "most-recent" + const blocksDiffLive = buildTranscriptDisplayBlocks(items, "sess-other"); + const boundaryDiff = blocksDiffLive.find( + (b) => b.kind === "session-boundary", + ); + assert.equal( + boundaryDiff.labelState, + "most-recent", + "labelState=most-recent when live id differs from newest visible", + ); +}); + +// ── Three sessions — two boundaries ──────────────────────────────────────────── + +test("buildTranscriptDisplayBlocks_threeSessions_twoBoundaries", () => { + const items = [ + sessionItem("a", "sess-1", "2026-07-08T00:00:01.000Z"), + sessionItem("b", "sess-2", "2026-07-08T00:00:02.000Z"), + sessionItem("c", "sess-3", "2026-07-08T00:00:03.000Z"), + ]; + const blocks = buildTranscriptDisplayBlocks(items); + const boundaryBlocks = blocks.filter((b) => b.kind === "session-boundary"); + assert.equal(boundaryBlocks.length, 2, "two boundaries for three sessions"); + // With no live id: newest (sess-3) boundary = "most-recent"; older = "earlier". + const newestBoundary = boundaryBlocks[boundaryBlocks.length - 1]; + const olderBoundary = boundaryBlocks[0]; + assert.equal( + newestBoundary.labelState, + "most-recent", + "newest boundary is most-recent when latestLiveSessionId is null", + ); + assert.equal( + olderBoundary.labelState, + "earlier", + "older boundary is earlier when latestLiveSessionId is null", + ); +}); + +// ── Null sessionId items stay in the current run ─────────────────────────────── + +test("buildTranscriptDisplayBlocks_nullSessionId_staysInCurrentRun", () => { + // An item with null sessionId should not start a new run. + const items = [ + sessionItem("a", "sess-1", "2026-07-08T00:00:01.000Z"), + // null sessionId — stays in sess-1 run + { ...sessionItem("b", null, "2026-07-08T00:00:02.000Z"), sessionId: null }, + sessionItem("c", "sess-1", "2026-07-08T00:00:03.000Z"), + ]; + const blocks = buildTranscriptDisplayBlocks(items); + const boundaryBlocks = blocks.filter((b) => b.kind === "session-boundary"); + assert.equal( + boundaryBlocks.length, + 0, + "no boundary when null-sessionId items are present within a single session", + ); +}); + +// ── flattenDisplayBlocks skips session-boundary blocks ──────────────────────── + +test("flattenDisplayBlocks_skipsSessionBoundaryBlocks", () => { + const items = [ + sessionItem("a", "sess-1", "2026-07-08T00:00:01.000Z"), + sessionItem("b", "sess-2", "2026-07-08T00:00:02.000Z"), + ]; + const blocks = buildTranscriptDisplayBlocks(items); + assert.ok( + blocks.some((b) => b.kind === "session-boundary"), + "test setup: boundary must be present", + ); + const flat = flattenDisplayBlocks(blocks); + const ids = flat.map((i) => i.id); + assert.ok(ids.includes("a"), "item a is in flattened output"); + assert.ok(ids.includes("b"), "item b is in flattened output"); + assert.equal( + flat.filter((i) => i.kind === "session-boundary").length, + 0, + "session-boundary items are excluded from flatten", + ); +}); + +// ── isObserverEventAfter — latest-live ordering ────────────────────────────── + +test("isObserverEventAfter returns true when candidate has later timestamp", () => { + const stored = { timestamp: "2026-07-08T00:00:01.000Z", seq: 5 }; + const candidate = { timestamp: "2026-07-08T00:00:02.000Z", seq: 1 }; + assert.ok(isObserverEventAfter(candidate, stored)); +}); + +test("isObserverEventAfter returns false when candidate has earlier timestamp", () => { + const stored = { timestamp: "2026-07-08T00:00:02.000Z", seq: 5 }; + const candidate = { timestamp: "2026-07-08T00:00:01.000Z", seq: 10 }; + assert.ok(!isObserverEventAfter(candidate, stored)); +}); + +test("isObserverEventAfter returns true for same timestamp, higher seq — session B advances over session A", () => { + // This is the tiebreak case: timestamp equal, seq tiebreak must mirror + // compareObserverEvents so latest-live never drifts from transcript order. + const stored = { timestamp: "2026-07-08T00:00:01.000Z", seq: 3 }; + const candidate = { timestamp: "2026-07-08T00:00:01.000Z", seq: 7 }; + assert.ok(isObserverEventAfter(candidate, stored)); +}); + +test("isObserverEventAfter returns false for same timestamp, same seq", () => { + const stored = { timestamp: "2026-07-08T00:00:01.000Z", seq: 3 }; + const candidate = { timestamp: "2026-07-08T00:00:01.000Z", seq: 3 }; + assert.ok(!isObserverEventAfter(candidate, stored)); +}); + +test("isObserverEventAfter returns false for same timestamp, lower seq", () => { + const stored = { timestamp: "2026-07-08T00:00:01.000Z", seq: 7 }; + const candidate = { timestamp: "2026-07-08T00:00:01.000Z", seq: 3 }; + assert.ok(!isObserverEventAfter(candidate, stored)); +}); diff --git a/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.ts b/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.ts index efbf1c5c2..ced54d19c 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.ts +++ b/desktop/src/features/agents/ui/agentSessionTranscriptGrouping.ts @@ -15,7 +15,26 @@ export type TranscriptTurnSegment = export type TranscriptDisplayBlock = | { kind: "single"; item: TranscriptItem } - | { kind: "turn"; turnId: string; segments: TranscriptTurnSegment[] }; + | { kind: "turn"; turnId: string; segments: TranscriptTurnSegment[] } + | { + /** + * Session boundary divider injected between consecutive session runs. + * `sessionId` is the id of the session that FOLLOWS the divider (the + * newer session in reading order). + * + * `labelState` encodes three distinct states: + * - `"current"` — newest-visible session AND matches the live relay + * session id. Agent is actively running this session. + * - `"most-recent"` — newest-visible session but no live session match + * (archived-only view or session ended). This is the + * most recently observed session — not current context. + * - `"earlier"` — an older session, not newest-visible. + */ + kind: "session-boundary"; + sessionId: string; + sessionStartTimestamp: string; + labelState: "current" | "most-recent" | "earlier"; + }; export type TranscriptToolRunChildSegment = | { kind: "item"; item: TranscriptItem } @@ -356,6 +375,34 @@ function getRenderClass(item: TranscriptItem) { return item.renderClass ?? descriptor.renderClass; } +/** + * Split a flat, time-ordered array of TranscriptItems into contiguous session + * runs. Items with a null sessionId are attributed to the most recently seen + * session (or a synthetic "unknown" run if no session has appeared yet). + * + * A new run begins whenever the sessionId changes to a distinct non-null value. + */ +function splitIntoSessionRuns( + items: TranscriptItem[], +): Array<{ sessionId: string; items: TranscriptItem[] }> { + const runs: Array<{ sessionId: string; items: TranscriptItem[] }> = []; + let currentRun: { sessionId: string; items: TranscriptItem[] } | null = null; + + for (const item of items) { + const sid: string = item.sessionId ?? currentRun?.sessionId ?? "unknown"; + if ( + !currentRun || + (item.sessionId && item.sessionId !== currentRun.sessionId) + ) { + currentRun = { sessionId: sid, items: [] }; + runs.push(currentRun); + } + currentRun.items.push(item); + } + + return runs; +} + /** * Build presentation-only display blocks from normalized transcript items. * Raw observer order is preserved in the source items; this only reorders @@ -365,9 +412,85 @@ function getRenderClass(item: TranscriptItem) { * turnId=null. They are injected into the prompt segment of the first turn * that follows them in stream order — placing System prompt between the user * message bubble and the Prompt context sections in the rendered output. + * + * When items span multiple sessions (archived history + live session), a + * `session-boundary` block is injected between consecutive session runs. + * The newest-visible run is labeled distinctly when it equals + * `latestLiveSessionId`, signalling "current session" vs archived history. + * No boundary is emitted when only one session run is present. + * + * @param latestLiveSessionId - The session ID that is currently live on the + * relay (from `observerRelayStore`). Used to distinguish "current session" + * from "most recent observed session". Pass `null` when the relay is idle. */ export function buildTranscriptDisplayBlocks( items: TranscriptItem[], + latestLiveSessionId: string | null = null, +): TranscriptDisplayBlock[] { + const sessionRuns = splitIntoSessionRuns(items); + + // Fast path: single session (or zero) — no boundary blocks needed. + if (sessionRuns.length <= 1) { + return buildBlocksForRun( + sessionRuns[0]?.items ?? [], + /* isNewestRun */ true, + sessionRuns[0]?.sessionId ?? null, + latestLiveSessionId, + /* emitBoundary */ false, + ); + } + + // Multi-session: build blocks per run and interleave session-boundary blocks. + const allBlocks: TranscriptDisplayBlock[] = []; + for (let i = 0; i < sessionRuns.length; i++) { + const run = sessionRuns[i]; + const isNewestRun = i === sessionRuns.length - 1; + + // Inject boundary before this run's blocks (not before the oldest run). + if (i > 0) { + const firstItem = run.items.find((it) => it.timestamp); + const sessionStartTimestamp = + firstItem?.timestamp ?? new Date(0).toISOString(); + const labelState: "current" | "most-recent" | "earlier" = + isNewestRun && + latestLiveSessionId !== null && + run.sessionId === latestLiveSessionId + ? "current" + : isNewestRun + ? "most-recent" + : "earlier"; + allBlocks.push({ + kind: "session-boundary", + sessionId: run.sessionId, + sessionStartTimestamp, + labelState, + }); + } + + const runBlocks = buildBlocksForRun( + run.items, + isNewestRun, + run.sessionId, + latestLiveSessionId, + /* emitBoundary */ false, + ); + allBlocks.push(...runBlocks); + } + + return allBlocks; +} + +/** + * Build display blocks for a single session run's items. + * Internal helper — extracted so `buildTranscriptDisplayBlocks` can call it + * once per run without code duplication. + */ +function buildBlocksForRun( + items: TranscriptItem[], + _isNewestRun: boolean, + _sessionId: string | null, + _latestLiveSessionId: string | null, + _emitBoundary: boolean, ): TranscriptDisplayBlock[] { const blocks: TranscriptDisplayBlock[] = []; const turnBuckets = new Map(); @@ -469,6 +592,11 @@ export function flattenDisplayBlocks( continue; } + // session-boundary blocks carry no items — skip. + if (block.kind === "session-boundary") { + continue; + } + for (const segment of block.segments) { if (segment.kind === "item") { result.push(segment.item); diff --git a/desktop/src/features/agents/ui/useObserverEvents.ts b/desktop/src/features/agents/ui/useObserverEvents.ts index 73603a029..e5248e7bb 100644 --- a/desktop/src/features/agents/ui/useObserverEvents.ts +++ b/desktop/src/features/agents/ui/useObserverEvents.ts @@ -9,10 +9,14 @@ import { } from "@/features/agents/observerRelayStore"; import { listSaveSubscriptions, - readArchivedEvents, + readArchivedObserverEventsForChannel, + readUnindexedObserverRows, + indexObserverChannelId, } from "@/shared/api/tauriArchive"; +import { decryptObserverEvent } from "@/shared/api/tauriObserver"; import { useIdentityQuery } from "@/shared/api/hooks"; import type { TranscriptItem } from "./agentSessionTypes"; +import type { RelayEvent } from "@/shared/api/types"; // Stable subscribe reference shared by all useSyncExternalStore hooks. // subscribeAgentObserverStore already has a fixed identity, so this thin @@ -55,16 +59,25 @@ export function useAgentTranscript( const ARCHIVED_EVENTS_PAGE_SIZE = 50; /** - * Load-older-on-scroll for archived observer frames. + * Load-older-on-scroll for archived observer frames, scoped to a single channel. * - * Checks whether an `owner_p` save subscription exists for the current - * identity. If one does, exposes `fetchOlderArchived` and `hasOlderArchived` - * for wiring into a sentinel-based scroll loader. + * Reads from `observer_channel_index` (via `readArchivedObserverEventsForChannel`) + * so only frames attributable to this channel are loaded — cross-channel + * contamination is impossible. Frames with null/decrypt-failed channelId are + * excluded at the Rust level (Will's (a) ruling). * - * Degrades cleanly when no subscription exists (returns `hasOlderArchived: - * false` without making any archive calls). + * On first mount, runs a one-shot idempotent backfill: decrypts all + * not-yet-indexed `owner_p` kind 24200 rows and writes their (id, channelId) + * pairs into the index, so existing archived history is available immediately + * without requiring the user to scroll through every page. + * + * Degrades cleanly when no `owner_p` subscription exists or when `channelId` + * is null (returns `hasOlderArchived: false` without making any archive calls). */ -export function useLoadArchivedObserverEvents(enabled: boolean) { +export function useLoadArchivedObserverEvents( + enabled: boolean, + channelId: string | null, +) { const identityQuery = useIdentityQuery(); const identityPubkey = identityQuery.data?.pubkey ?? null; @@ -74,6 +87,22 @@ export function useLoadArchivedObserverEvents(enabled: boolean) { ); const [hasOlderArchived, setHasOlderArchived] = React.useState(true); const isFetchingRef = React.useRef(false); + // Backfill state: "pending" → "running" → "done". + // fetchOlderArchived awaits backfillPromiseRef before reading the index so + // the first scroll-trigger never races the write path and incorrectly marks + // the channel exhausted before backfill has completed. + const backfillStatusRef = React.useRef<"pending" | "running" | "done">( + "pending", + ); + const backfillPromiseRef = React.useRef | null>(null); + const backfillResolveRef = React.useRef<(() => void) | null>(null); + // Expose a promise that resolves when backfill is done. Created eagerly so + // fetchOlderArchived can await it before the effect that starts backfill fires. + if (!backfillPromiseRef.current) { + backfillPromiseRef.current = new Promise((resolve) => { + backfillResolveRef.current = resolve; + }); + } // Compound keyset cursor: tracks both `created_at` and `id` of the oldest // event seen so far. Mirrors the SQL `ORDER BY created_at DESC, id DESC` so // same-second siblings are never skipped at a page boundary. @@ -98,12 +127,18 @@ export function useLoadArchivedObserverEvents(enabled: boolean) { setHasSubscription(hasSub); if (!hasSub) { setHasOlderArchived(false); + // No subscription → backfill will never run; resolve the promise + // immediately so fetchOlderArchived doesn't await indefinitely. + backfillStatusRef.current = "done"; + backfillResolveRef.current?.(); } }) .catch(() => { if (!cancelled) { setHasSubscription(false); setHasOlderArchived(false); + backfillStatusRef.current = "done"; + backfillResolveRef.current?.(); } }); return () => { @@ -111,22 +146,112 @@ export function useLoadArchivedObserverEvents(enabled: boolean) { }; }, [enabled, identityPubkey]); + // One-shot idempotent backfill: attempt to decrypt all not-yet-processed + // owner_p kind 24200 rows and write their (id, channelId?) into + // observer_channel_index. A status row is written for EVERY processed event — + // null/failed channelId rows get channel_id=null, so re-runs skip them. + // Runs once per mount when the subscription is confirmed; gated by + // backfillStatusRef so fetchOlderArchived can await completion. + React.useEffect(() => { + if ( + !enabled || + !hasSubscription || + backfillStatusRef.current !== "pending" + ) { + return; + } + backfillStatusRef.current = "running"; + const promise = (async () => { + try { + const rows = await readUnindexedObserverRows(); + + const toIndex: Array<{ + eventId: string; + channelId: string | null; + createdAt: number; + }> = []; + + for (const row of rows) { + let parsed: RelayEvent; + try { + parsed = JSON.parse(row.rawJson) as RelayEvent; + } catch { + // Malformed JSON: write a null status row so we skip on re-run. + toIndex.push({ + eventId: row.id, + channelId: null, + createdAt: row.createdAt, + }); + continue; + } + try { + const decoded = (await decryptObserverEvent(parsed)) as { + channelId?: string | null; + }; + // Write a status row for every event — non-null channelId is + // attributable; null/undefined channelId writes channel_id=null so + // the frame is marked processed and excluded from scoped views. + toIndex.push({ + eventId: row.id, + channelId: decoded?.channelId ?? null, + createdAt: row.createdAt, + }); + } catch { + // Decrypt failure → write null status row (processed, unscoped). + toIndex.push({ + eventId: row.id, + channelId: null, + createdAt: row.createdAt, + }); + } + } + + if (toIndex.length > 0) { + await indexObserverChannelId(toIndex); + } + } catch (error) { + console.error( + "[useLoadArchivedObserverEvents] backfill failed:", + error, + ); + } finally { + backfillStatusRef.current = "done"; + backfillResolveRef.current?.(); + } + })(); + backfillPromiseRef.current = promise; + }, [enabled, hasSubscription]); + const fetchOlderArchived = React.useCallback(async () => { if ( !enabled || !identityPubkey || !hasSubscription || + !channelId || isFetchingRef.current || !hasOlderArchived ) { return; } + // Await backfill completion before reading the channel index. This + // guarantees the index is populated before the first paginated read, so + // a scroll-trigger that fires before backfill writes can't return 0 rows + // and falsely mark the channel exhausted. + if (backfillPromiseRef.current) { + await backfillPromiseRef.current; + } + + // Re-check after awaiting: hasOlderArchived might have been set false + // while we were waiting (e.g. subscription check failed). + if (!hasOlderArchived) { + return; + } + isFetchingRef.current = true; try { const before = cursorRef.current ?? undefined; - const events = await readArchivedEvents("owner_p", identityPubkey, { - kinds: [24200], + const events = await readArchivedObserverEventsForChannel(channelId, { before: before ?? null, limit: ARCHIVED_EVENTS_PAGE_SIZE, }); @@ -143,7 +268,7 @@ export function useLoadArchivedObserverEvents(enabled: boolean) { await ingestArchivedObserverEvents(events); } - // A short page means the archive is exhausted. + // A short page means the archive is exhausted for this channel. if (events.length < ARCHIVED_EVENTS_PAGE_SIZE) { setHasOlderArchived(false); } @@ -152,7 +277,7 @@ export function useLoadArchivedObserverEvents(enabled: boolean) { } finally { isFetchingRef.current = false; } - }, [enabled, identityPubkey, hasSubscription, hasOlderArchived]); + }, [enabled, identityPubkey, hasSubscription, channelId, hasOlderArchived]); return { fetchOlderArchived, hasOlderArchived }; } diff --git a/desktop/src/features/channels/ui/AgentSessionThreadPanel.tsx b/desktop/src/features/channels/ui/AgentSessionThreadPanel.tsx index f0d2669b0..0de5fd272 100644 --- a/desktop/src/features/channels/ui/AgentSessionThreadPanel.tsx +++ b/desktop/src/features/channels/ui/AgentSessionThreadPanel.tsx @@ -131,7 +131,7 @@ export function AgentSessionThreadPanel({ : `Last updated ${new Date(latestActivityAt).toLocaleString()}`; const { fetchOlderArchived, hasOlderArchived } = - useLoadArchivedObserverEvents(isLive); + useLoadArchivedObserverEvents(isLive, sessionChannelId ?? null); useLoadOlderOnScroll({ fetchOlder: fetchOlderArchived, diff --git a/desktop/src/shared/api/tauriArchive.ts b/desktop/src/shared/api/tauriArchive.ts index 4c927373c..5dcc1013a 100644 --- a/desktop/src/shared/api/tauriArchive.ts +++ b/desktop/src/shared/api/tauriArchive.ts @@ -220,6 +220,92 @@ export async function archiveEvents( }); } +/** + * Read a paginated page of archived kind 24200 (observer) events for a + * specific channel, using the `observer_channel_index`. + * + * Only returns frames whose `channelId` was successfully decrypted and + * matched this channel. Frames with null/decrypt-failed channelId are + * excluded (Will's (a) ruling). Compound cursor + short-page exhaustion + * signal work identically to `readArchivedEvents`. + */ +export async function readArchivedObserverEventsForChannel( + channelId: string, + opts?: { + before?: { createdAt: number; id: string } | null; + limit?: number; + }, +): Promise { + const rawRows = await invokeTauri( + "read_archived_observer_events_for_channel", + { + channelId, + beforeCreatedAt: opts?.before?.createdAt ?? null, + beforeId: opts?.before?.id ?? null, + limit: opts?.limit ?? null, + }, + ); + return rawRows + .map((raw) => { + try { + return JSON.parse(raw) as import("@/shared/api/types").RelayEvent; + } catch { + console.warn( + "[tauriArchive] failed to parse archived observer raw_json:", + raw, + ); + return null; + } + }) + .filter((e): e is import("@/shared/api/types").RelayEvent => e !== null); +} + +/** + * Index one or more archived observer frames by channelId. + * + * `channelId` is nullable: pass `null` for frames that are unscoped, + * malformed, or whose payload could not be decrypted. Null rows are written + * as a processed-state marker so re-runs skip them; they are never returned + * by channel-scoped reads (which filter `channel_id = ?`). + * + * Idempotent — already-indexed frames are silently skipped. + */ +export async function indexObserverChannelId( + entries: Array<{ + eventId: string; + channelId: string | null; + createdAt: number; + }>, +): Promise { + if (entries.length === 0) return; + await invokeTauri("index_observer_channel_id", { + entries: entries.map((e) => ({ + event_id: e.eventId, + channel_id: e.channelId, + created_at: e.createdAt, + })), + }); +} + +/** + * Return all `owner_p` kind 24200 archived event rows not yet indexed. + * + * Used by the one-shot backfill driver. Returns raw Nostr event JSON plus + * event id and created_at for each row so the caller can decrypt and index. + */ +export async function readUnindexedObserverRows(): Promise< + Array<{ id: string; rawJson: string; createdAt: number }> +> { + const rows = await invokeTauri< + Array<{ id: string; raw_json: string; created_at: number }> + >("read_unindexed_observer_rows"); + return rows.map((r) => ({ + id: r.id, + rawJson: r.raw_json, + createdAt: r.created_at, + })); +} + /** * Read a paginated page of archived raw events for a scope. *