From df0a0861776f358e8c365e18a429900044b12989 Mon Sep 17 00:00:00 2001 From: Will Pfleger Date: Thu, 23 Jul 2026 15:35:19 -0400 Subject: [PATCH 1/4] fix(cli): install rustls crypto provider to unbreak WSS publishes in release builds (#2590) Signed-off-by: Will Pfleger Co-authored-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 --- Cargo.lock | 1 + crates/buzz-cli/Cargo.toml | 7 +++++++ crates/buzz-cli/src/lib.rs | 12 ++++++++++++ 3 files changed, 20 insertions(+) diff --git a/Cargo.lock b/Cargo.lock index b5b2605a4..32152e41a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -899,6 +899,7 @@ dependencies = [ "nostr", "rand 0.10.1", "reqwest 0.13.4", + "rustls", "serde", "serde_json", "sha2 0.11.0", diff --git a/crates/buzz-cli/Cargo.toml b/crates/buzz-cli/Cargo.toml index a12260b52..1476e60bf 100644 --- a/crates/buzz-cli/Cargo.toml +++ b/crates/buzz-cli/Cargo.toml @@ -76,6 +76,13 @@ dirs = "6" # WebSocket client — ephemeral event publish (kind:20001 is WS-only on the relay) buzz-ws-client = { path = "../buzz-ws-client" } +# Explicit rustls dep with ring provider — required to install the process-level +# CryptoProvider at startup. Without this the standalone `buzz` binary panics when +# a multi-package release build (buzz-acp + buzz-dev-mcp + buzz-cli in one cargo +# invocation) unifies both ring and aws-lc-rs features, leaving rustls unable to +# auto-select a provider. See crates/buzz-acp/Cargo.toml for the same dependency. +rustls = { version = "0.23", default-features = false, features = ["ring", "std"] } + # Random number generation — full jitter for exponential backoff in with_retry rand = { workspace = true } diff --git a/crates/buzz-cli/src/lib.rs b/crates/buzz-cli/src/lib.rs index d5c6b6f9a..0f8caa416 100644 --- a/crates/buzz-cli/src/lib.rs +++ b/crates/buzz-cli/src/lib.rs @@ -25,6 +25,18 @@ where I: IntoIterator, S: Into + Clone, { + // Install ring as the process-level rustls CryptoProvider. Required because the + // release workflow builds all binaries in one cargo invocation, which unifies + // features across the workspace and enables *both* ring (from buzz-acp/buzz-dev-mcp) + // and aws-lc-rs (from reqwest's rustls feature via hyper-rustls). With both on, + // rustls cannot auto-select a provider, and any code that reaches + // ClientConfig::builder() — specifically the WSS path in publish_ephemeral_event + // used by `agents draft-create`, `agents draft-update`, and `users set-presence` + // — panics at rustls crypto/mod.rs. The `let _ =` swallow is intentional: when + // buzz-dev-mcp delegates to run_from_args, it has already installed ring; the + // double-install returns Err and is harmless. + let _ = rustls::crypto::ring::default_provider().install_default(); + let cli = match Cli::try_parse_from(args) { Ok(cli) => cli, Err(e) => { From 8cb05028be84829501a4f87f8db4ee7034fb786f Mon Sep 17 00:00:00 2001 From: Will Pfleger Date: Thu, 23 Jul 2026 15:40:46 -0400 Subject: [PATCH 2/4] fix(observer): eager archive hydration on panel open + 200-frame pages (#2574) Signed-off-by: Will Pfleger Co-authored-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 --- .../ingestArchivedObserverEvents.test.mjs | 205 +++- .../features/agents/ui/archivePagingState.ts | 67 +- .../ui/useLoadArchivedObserverEvents.test.mjs | 936 ++++++++++++++++++ .../features/agents/ui/useObserverEvents.ts | 107 +- 4 files changed, 1293 insertions(+), 22 deletions(-) create mode 100644 desktop/src/features/agents/ui/useLoadArchivedObserverEvents.test.mjs diff --git a/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs b/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs index eb38da9a8..c44c2a88e 100644 --- a/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs +++ b/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs @@ -354,7 +354,7 @@ describe("load-older cursor advance logic", () => { it("test_short_page_signals_archive_exhausted", () => { // A page with fewer events than the limit signals end-of-archive. - const PAGE_SIZE = 50; + const PAGE_SIZE = 200; const page = Array.from({ length: 30 }, (_, i) => ({ created_at: 1000 - i, })); @@ -367,8 +367,8 @@ describe("load-older cursor advance logic", () => { }); it("test_full_page_signals_more_archive_available", () => { - const PAGE_SIZE = 50; - const page = Array.from({ length: 50 }, (_, i) => ({ + const PAGE_SIZE = 200; + const page = Array.from({ length: 200 }, (_, i) => ({ created_at: 1000 - i, })); const exhausted = page.length < PAGE_SIZE; @@ -394,6 +394,7 @@ describe("load-older cursor advance logic", () => { import { createArchivePagingState, applyChannelReset, + runHydrationLoop, } from "@/features/agents/ui/archivePagingState.ts"; describe("archive paging state reset on channel change", () => { @@ -410,6 +411,12 @@ describe("archive paging state reset on channel change", () => { "backfillPromise is eagerly initialized", ); assert.equal(ps.cursor, null, "cursor starts null"); + assert.equal( + ps.initialHydrationDone, + false, + "initialHydrationDone starts false", + ); + assert.equal(ps.activeChannelId, null, "activeChannelId starts null"); }); it("test_channel_switch_resets_cursor_exhaustion_and_fetch_lock", () => { @@ -419,11 +426,13 @@ describe("archive paging state reset on channel change", () => { ps.cursor = { createdAt: 1000, id: "event-a5" }; ps.hasOlderArchived = false; // channel A exhausted ps.isFetching = true; // mid-flight request (edge case) + ps.initialHydrationDone = true; // hydration ran for channel A + ps.activeChannelId = "chan-a"; ps.backfillStatus = "done"; // backfill ran once already const originalPromise = ps.backfillPromise; // must survive reset // Channel switch — this is what the useEffect([channelId]) calls. - applyChannelReset(ps); + applyChannelReset(ps, "chan-b"); assert.equal(ps.cursor, null, "cursor resets to null on channel switch"); assert.equal( @@ -436,6 +445,16 @@ describe("archive paging state reset on channel change", () => { false, "isFetching resets to false on channel switch", ); + assert.equal( + ps.initialHydrationDone, + false, + "initialHydrationDone resets to false on channel switch so the new channel hydrates", + ); + assert.equal( + ps.activeChannelId, + "chan-b", + "activeChannelId updates to new channel on switch", + ); // Backfill state must NOT be touched — it is identity-level and should // survive channel switches so the backfill only runs once per identity mount. @@ -459,10 +478,11 @@ describe("archive paging state reset on channel change", () => { it("test_multiple_channel_switches_each_start_fresh", () => { const ps = createArchivePagingState(); - // Switch to channel A: exhaust it. + // Switch to channel A: exhaust it and complete hydration. ps.cursor = { createdAt: 500, id: "a-oldest" }; ps.hasOlderArchived = false; - applyChannelReset(ps); + ps.initialHydrationDone = true; + applyChannelReset(ps, "chan-b"); assert.equal(ps.cursor, null, "switch A→B: cursor reset"); assert.equal( @@ -470,11 +490,22 @@ describe("archive paging state reset on channel change", () => { true, "switch A→B: hasOlderArchived reset", ); + assert.equal( + ps.initialHydrationDone, + false, + "switch A→B: initialHydrationDone reset", + ); + assert.equal( + ps.activeChannelId, + "chan-b", + "switch A→B: activeChannelId updated", + ); - // Simulate channel B also being paged. + // Simulate channel B also being paged and hydrated. ps.cursor = { createdAt: 200, id: "b-oldest" }; ps.hasOlderArchived = false; - applyChannelReset(ps); + ps.initialHydrationDone = true; + applyChannelReset(ps, "chan-c"); assert.equal(ps.cursor, null, "switch B→C: cursor reset again"); assert.equal( @@ -482,6 +513,164 @@ describe("archive paging state reset on channel change", () => { true, "switch B→C: hasOlderArchived reset again", ); + assert.equal( + ps.initialHydrationDone, + false, + "switch B→C: initialHydrationDone reset again", + ); + assert.equal( + ps.activeChannelId, + "chan-c", + "switch B→C: activeChannelId updated", + ); + }); +}); + +// ── Eager initial hydration loop logic ─────────────────────────────────────── +// +// The initial hydration loop in useLoadArchivedObserverEvents calls +// fetchOlderArchived up to INITIAL_HYDRATION_BUDGET_PAGES times. The loop must: +// - Stop at budget (10 pages) even if more archive exists. +// - Stop early when the archive is exhausted (ps.hasOlderArchived → false). +// - Respect channel-switch cancellation (signal.cancelled). +// +// These tests call the PRODUCTION runHydrationLoop from archivePagingState.ts +// with mock fetchOnePage functions — so they fail if the production loop logic +// is deleted or misrouted, not just if a reimplemented copy breaks. + +describe("eager initial hydration loop control flow (production runHydrationLoop)", () => { + const BUDGET = 10; // mirrors INITIAL_HYDRATION_BUDGET_PAGES + + it("test_hydration_stops_at_budget_when_archive_never_exhausted", async () => { + const ps = createArchivePagingState(); + applyChannelReset(ps, "chan-1"); + let fetchCount = 0; + const fetchOnePage = async () => { + fetchCount++; + // archive remains non-empty — ps.hasOlderArchived stays true + }; + const signal = { cancelled: false }; + await runHydrationLoop(ps, fetchOnePage, BUDGET, signal); + assert.equal( + fetchCount, + BUDGET, + `production runHydrationLoop must stop after exactly ${BUDGET} pages (budget limit)`, + ); + }); + + it("test_hydration_stops_early_when_archive_exhausted", async () => { + const ps = createArchivePagingState(); + applyChannelReset(ps, "chan-1"); + let fetchCount = 0; + const fetchOnePage = async () => { + fetchCount++; + if (fetchCount >= 3) { + ps.hasOlderArchived = false; // mock: archive exhausted on page 3 + } + }; + const signal = { cancelled: false }; + await runHydrationLoop(ps, fetchOnePage, BUDGET, signal); + assert.equal( + fetchCount, + 3, + "production runHydrationLoop must stop as soon as ps.hasOlderArchived is false (before budget)", + ); + }); + + it("test_hydration_respects_cancellation_on_channel_switch", async () => { + const ps = createArchivePagingState(); + applyChannelReset(ps, "chan-1"); + const signal = { cancelled: false }; + let fetchCount = 0; + const fetchOnePage = async () => { + fetchCount++; + if (fetchCount >= 2) { + signal.cancelled = true; // mock: channel switched mid-loop + } + }; + await runHydrationLoop(ps, fetchOnePage, BUDGET, signal); + assert.equal( + fetchCount, + 2, + "production runHydrationLoop must stop when signal.cancelled is true (channel switch)", + ); + }); + + it("test_hydration_zero_iterations_when_already_exhausted", async () => { + // If ps.hasOlderArchived is already false before the loop starts (e.g. + // channel A was exhausted and reset did not run yet), the loop must not + // call fetchOnePage at all. + const ps = createArchivePagingState(); + applyChannelReset(ps, "chan-1"); + ps.hasOlderArchived = false; // already exhausted + let fetchCount = 0; + const fetchOnePage = async () => { + fetchCount++; + }; + const signal = { cancelled: false }; + await runHydrationLoop(ps, fetchOnePage, BUDGET, signal); + assert.equal( + fetchCount, + 0, + "must not fetch when archive is already exhausted", + ); + }); + + // Regression: stale React closure — the original fetchOlderArchived captured + // hasOlderArchived from React state, which was false (exhausted) for channel A + // while ps.hasOlderArchived had already been reset to true by applyChannelReset. + // The production loop uses ps.hasOlderArchived (the ref) to guard iterations, + // so the switch-then-hydrate path must call fetchOnePage for the new channel. + it("test_exhausted_channel_A_then_switch_to_B_hydrates_B", async () => { + const ps = createArchivePagingState(); + + // Simulate exhausting channel A. + applyChannelReset(ps, "chan-a"); + ps.hasOlderArchived = false; // channel A exhausted + + // Switch to channel B — resets hasOlderArchived and activeChannelId. + applyChannelReset(ps, "chan-b"); + + // ps.hasOlderArchived is now true; the loop should call fetchOnePage. + let fetchCount = 0; + const fetchOnePage = async () => { + fetchCount++; + ps.hasOlderArchived = false; // B also exhausted after 1 page + }; + const signal = { cancelled: false }; + await runHydrationLoop(ps, fetchOnePage, BUDGET, signal); + + assert.equal( + fetchCount, + 1, + "after switch from exhausted channel A to B, runHydrationLoop must call fetchOnePage for B (not skip due to stale exhaustion)", + ); + }); + + // Regression: stale cursor write — an in-flight read from channel A should + // not write A's cursor into ps.cursor after the switch to B. The activeChannelId + // token (set by applyChannelReset) is what lets fetchOlderArchived detect and + // discard the stale result. This test verifies that applyChannelReset correctly + // advances the token so a pre-switch requestChannelId !== ps.activeChannelId. + it("test_activeChannelId_token_detects_stale_read_from_prior_channel", () => { + const ps = createArchivePagingState(); + applyChannelReset(ps, "chan-a"); + const requestChannelId = ps.activeChannelId; // captured at request start = "chan-a" + + // Simulate channel switch before the Tauri read resolves. + applyChannelReset(ps, "chan-b"); + + // The in-flight A read checks requestChannelId !== ps.activeChannelId. + assert.notEqual( + requestChannelId, + ps.activeChannelId, + "requestChannelId from channel A must not match activeChannelId after switching to B — stale read must be discarded", + ); + assert.equal( + ps.activeChannelId, + "chan-b", + "activeChannelId must reflect the current channel after switch", + ); }); }); diff --git a/desktop/src/features/agents/ui/archivePagingState.ts b/desktop/src/features/agents/ui/archivePagingState.ts index 95574d311..eaaec7529 100644 --- a/desktop/src/features/agents/ui/archivePagingState.ts +++ b/desktop/src/features/agents/ui/archivePagingState.ts @@ -28,6 +28,20 @@ export interface ArchivePagingState { * Mirrors SQL ORDER BY created_at DESC, id DESC so same-second siblings are * never skipped at a page boundary. */ cursor: { createdAt: number; id: string } | null; + /** True once the initial eager-hydration pass for the current channel has + * completed (budget reached or archive exhausted). Resets on channel change + * so switching channels triggers a fresh hydration pass. */ + initialHydrationDone: boolean; + /** The channelId that this paging state is currently scoped to. + * Kept for diagnostics only; NOT used as a generation token (see + * resetGeneration). Channel equality is not a unique request identifier: + * A→B→A makes old-A channel checks pass again. */ + activeChannelId: string | null; + /** Monotonically increasing counter incremented by applyChannelReset. + * Each fetch snapshots this value at request start and checks it again + * after every async boundary — a mismatch means a channel switch occurred + * mid-flight (even A→B→A), and results are discarded. */ + resetGeneration: number; } /** @@ -43,6 +57,9 @@ export function createArchivePagingState(): ArchivePagingState { backfillPromise: null, backfillResolve: null, cursor: null, + initialHydrationDone: false, + activeChannelId: null, + resetGeneration: 0, }; state.backfillPromise = new Promise((resolve) => { state.backfillResolve = resolve; @@ -53,15 +70,57 @@ export function createArchivePagingState(): ArchivePagingState { /** * Reset per-channel paging state when the viewed channel changes. * - * Only cursor, exhaustion flag, and fetch lock are channel-scoped. Backfill - * state is identity-level (the index covers ALL channels and needs to run only - * once per identity mount), so it is intentionally NOT touched here. + * Only cursor, exhaustion flag, fetch lock, channel label, and generation token + * are channel-scoped. Backfill state is identity-level (the index covers ALL + * channels and needs to run only once per identity mount), so it is + * intentionally NOT touched here. + * + * `resetGeneration` is incremented on every call. In-flight fetches snapshot + * the generation at start and recheck it after every async boundary — a + * mismatch (including A→B→A) means the request is stale, so results are + * discarded. Channel ID is retained for diagnostics only; it does NOT serve + * as the generation token. * * Called by the useEffect([channelId]) in useLoadArchivedObserverEvents. * Exported so tests can verify the reset semantics without a React runtime. */ -export function applyChannelReset(state: ArchivePagingState): void { +export function applyChannelReset( + state: ArchivePagingState, + newChannelId: string | null, +): void { state.cursor = null; state.isFetching = false; state.hasOlderArchived = true; + state.initialHydrationDone = false; + state.activeChannelId = newChannelId; + state.resetGeneration += 1; +} + +/** + * Run the eager initial-hydration paging loop. + * + * Calls `fetchOnePage()` up to `budget` times. Stops early when: + * - `ps.hasOlderArchived` is false (archive exhausted for this channel), OR + * - `signal.cancelled` is true (channel switched away mid-loop). + * + * `fetchOnePage` is the per-page read unit: it must respect `ps.isFetching` + * (lock), await backfill, perform the Tauri read, ingest results, advance + * the cursor, and set `ps.hasOlderArchived = false` when the page is short. + * The hook wires the real implementation; tests supply a mock. + * + * Exported so tests can call the production loop logic directly — passing a + * mock `fetchOnePage` — without reimplementing the control flow. + */ +export async function runHydrationLoop( + ps: ArchivePagingState, + fetchOnePage: () => Promise, + budget: number, + signal: { cancelled: boolean }, +): Promise { + for (let page = 0; page < budget; page++) { + if (signal.cancelled || !ps.hasOlderArchived) { + break; + } + await fetchOnePage(); + } } diff --git a/desktop/src/features/agents/ui/useLoadArchivedObserverEvents.test.mjs b/desktop/src/features/agents/ui/useLoadArchivedObserverEvents.test.mjs new file mode 100644 index 000000000..c3ca1f678 --- /dev/null +++ b/desktop/src/features/agents/ui/useLoadArchivedObserverEvents.test.mjs @@ -0,0 +1,936 @@ +/** + * Mounted-hook lifecycle and race regression tests for + * useLoadArchivedObserverEvents. + * + * These tests mount the REAL production hook (including its useEffect wiring, + * fetchOlderArchived closure, runHydrationLoop call, and resetGeneration token + * checks) against a mocked Tauri IPC bridge and a real QueryClientProvider. + * They fail if any of the following is removed from the production hook: + * - the hydration effect + * - the resetGeneration checks (post-backfill, post-Tauri-read, post-ingest) + * - the generation-aware isFetching clear in finally + * - the post-backfill isFetching recheck before lock acquisition + * + * Four regressions: + * (a) exhausted-A → switch-to-B: B must read from a null cursor and ingest + * its rows. GREEN at dfb2d0385 (stale-closure was fixed in round 1). + * Fails at the pre-round-1 head where ps.hasOlderArchived was read from + * React state, not the ref. + * (b) deferred-I/O race (A→B): A is in flight (decrypt deferred), switch to B, + * resolve A's 1-row short ingest — A must NOT mark B exhausted or steal B's + * fetch lock, and B's eager loop must continue past page 1. Fails at + * dfb2d0385 (post-ingest token missing). Lock theft is asserted by calling + * fetchOlderArchived concurrently while B holds the lock: with correct + * protection the concurrent call is rejected (read count doesn't jump); + * without it, A's stale finally clears the lock and the concurrent call + * would start a duplicate read. + * (c) concurrent fetches during pending backfill: two callers both suspend on + * the backfill promise, both resume after it resolves — only one must + * acquire the lock and issue the Tauri read. Fails at 4c92a018d (no + * post-backfill isFetching recheck). + * (d) A→B→A: old-A's in-flight decrypt completes after the user returns to A. + * Old-A's generation no longer matches (each reset increments the counter), + * so it must not mark fresh-A exhausted or steal fresh-A's lock. Fails + * when resetGeneration is replaced with a channel-string equality check. + * + * ── DOM shim ───────────────────────────────────────────────────────────────── + * react-dom/client requires a minimal DOM; node has none. We install the same + * minimal shim used by MessageComposerDraftImagePersist.test.mjs. + * + * ── Tauri IPC mock ─────────────────────────────────────────────────────────── + * @tauri-apps/api/core calls window.__TAURI_INTERNALS__.invoke(cmd, args). + * We install a per-test mock at globalThis.__TAURI_INTERNALS__.invoke so every + * listSaveSubscriptions / readArchivedObserverEventsForChannel / readUnindexed / + * indexObserverChannelId call is intercepted by command name without patching + * module internals. + */ + +import assert from "node:assert/strict"; +import { describe, it, beforeEach } from "node:test"; + +// ── Minimal DOM shim (matches MessageComposerDraftImagePersist.test.mjs) ────── + +function installDOMShim() { + class MinimalEventTarget { + constructor() { + this._listeners = {}; + } + addEventListener(type, fn) { + if (!this._listeners[type]) this._listeners[type] = []; + this._listeners[type].push(fn); + } + removeEventListener(type, fn) { + if (this._listeners[type]) { + this._listeners[type] = this._listeners[type].filter((f) => f !== fn); + } + } + dispatchEvent(e) { + for (const fn of this._listeners[e.type] ?? []) fn(e); + return true; + } + } + + class MinimalNode extends MinimalEventTarget { + constructor(tagName) { + super(); + this.tagName = tagName; + this.children = []; + this.childNodes = []; + this.style = {}; + this.nodeType = 1; + this.parentNode = null; + } + get ownerDocument() { + return globalThis.document; + } + get firstChild() { + return this.children[0] ?? null; + } + get lastChild() { + return this.children[this.children.length - 1] ?? null; + } + get nextSibling() { + return null; + } + get nodeValue() { + return null; + } + appendChild(child) { + this.children.push(child); + this.childNodes.push(child); + child.parentNode = this; + return child; + } + removeChild(child) { + this.children = this.children.filter((c) => c !== child); + this.childNodes = this.childNodes.filter((c) => c !== child); + return child; + } + insertBefore(newNode, refNode) { + if (!refNode) return this.appendChild(newNode); + const i = this.children.indexOf(refNode); + if (i < 0) return this.appendChild(newNode); + this.children.splice(i, 0, newNode); + this.childNodes.splice(i, 0, newNode); + newNode.parentNode = this; + return newNode; + } + contains(node) { + if (!node) return false; + return this === node || this.children.some((c) => c?.contains?.(node)); + } + } + + class MinimalDocument extends MinimalEventTarget { + constructor() { + super(); + this.nodeType = 9; + } + createElement(tagName) { + return new MinimalNode(tagName); + } + createTextNode(value) { + const n = new MinimalNode("#text"); + n.nodeValue = value; + n.nodeType = 3; + return n; + } + createComment(value) { + const n = new MinimalNode("#comment"); + n.nodeValue = value; + n.nodeType = 8; + return n; + } + get body() { + if (!this._body) this._body = this.createElement("body"); + return this._body; + } + get activeElement() { + return null; + } + contains(node) { + return node != null; + } + } + + globalThis.document = new MinimalDocument(); + globalThis.HTMLIFrameElement = MinimalNode; + globalThis.HTMLElement = MinimalNode; + globalThis.IS_REACT_ACT_ENVIRONMENT = true; + process.env.IS_REACT_ACT_ENVIRONMENT = "true"; + + if (typeof globalThis.window === "undefined") { + Object.defineProperty(globalThis, "window", { + value: globalThis, + configurable: true, + }); + } + if (!Object.getOwnPropertyDescriptor(globalThis, "navigator")?.value) { + Object.defineProperty(globalThis, "navigator", { + value: { userAgent: "node" }, + configurable: true, + }); + } + globalThis.MutationObserver = class { + observe() {} + disconnect() {} + takeRecords() { + return []; + } + }; + globalThis.requestAnimationFrame = (fn) => setTimeout(fn, 0); +} + +installDOMShim(); + +// ── Tauri IPC interceptor ───────────────────────────────────────────────────── +// +// @tauri-apps/api/core calls window.__TAURI_INTERNALS__.invoke(cmd, args). +// Install a stub now (before any module that imports tauriArchive is loaded) +// so listSaveSubscriptions, readArchivedObserverEventsForChannel, etc. can be +// controlled per-test by replacing ipcHandlers. + +/** @type {Map Promise>} */ +const ipcHandlers = new Map(); + +globalThis.__TAURI_INTERNALS__ = { + invoke: (cmd, args) => { + const handler = ipcHandlers.get(cmd); + if (handler) return handler(args); + return Promise.reject(new Error(`unmocked Tauri command: ${cmd}`)); + }, + transformCallback: (_cb) => { + const id = Math.random(); + return id; + }, +}; + +function setIpcHandler(cmd, fn) { + ipcHandlers.set(cmd, fn); +} +function clearIpcHandlers() { + ipcHandlers.clear(); +} + +// ── Production imports (after shim, after IPC stub) ─────────────────────────── + +import React from "react"; +import { createRoot } from "react-dom/client"; +import { act } from "react"; +import { QueryClient, QueryClientProvider } from "@tanstack/react-query"; + +import { useLoadArchivedObserverEvents } from "@/features/agents/ui/useObserverEvents.ts"; +import { + resetAgentObserverStore, + _testRegisterKnownAgents, + _testGetArchivedChannelEvents, +} from "@/features/agents/observerRelayStore.ts"; + +// ── Constants ───────────────────────────────────────────────────────────────── + +const AGENT_PUBKEY = "a".repeat(64); +const IDENTITY_PUBKEY = "c".repeat(64); +const SUB_ID = "test-hook-sub"; + +// ── Tauri wire-shape helpers ────────────────────────────────────────────────── + +/** Returns a list_save_subscriptions response with one owner_p subscription. */ +function makeOwnerPSubResponse() { + return [ + { + identity_pubkey: IDENTITY_PUBKEY, + relay_url: "wss://test", + scope_type: "owner_p", + scope_value: IDENTITY_PUBKEY, + kinds: "[24200]", + created_at: 1000, + }, + ]; +} + +/** Returns a raw archived observer event row for readArchivedObserverEventsForChannel. */ +function makeArchivedRow(seq, channelId = "chan-1") { + return { + id: `ev${String(seq).padStart(63, "0")}`, + pubkey: AGENT_PUBKEY, + created_at: 1000 + seq, + kind: 24200, + tags: [ + ["p", IDENTITY_PUBKEY], + ["agent", AGENT_PUBKEY], + ["frame", "telemetry"], + ], + content: JSON.stringify({ + seq, + timestamp: new Date(1_000_000 + seq * 1000).toISOString(), + channelId, + kind: "telemetry", + sessionId: "sess-1", + turnId: "turn-1", + payload: { method: "session/update", params: {} }, + }), + sig: "s".repeat(128), + }; +} + +// ── React mounting helpers ──────────────────────────────────────────────────── + +/** + * Mount useLoadArchivedObserverEvents in a real React tree with a QueryClient + * pre-seeded with the identity. Returns { unmount, render(channelId), + * getFetchOlderArchived() }. + * + * getFetchOlderArchived() returns the latest fetchOlderArchived function from + * the hook's return value, captured on each render. Tests can call it directly + * to probe lock behaviour without going through the hydration loop. + */ +function mountHook(_initialChannelId, queryClient) { + // Capture the latest hook return values so tests can call fetchOlderArchived. + const hookReturnRef = { current: null }; + + function HarnessComponent({ channelId }) { + const result = useLoadArchivedObserverEvents(true, channelId); + hookReturnRef.current = result; + return null; + } + + const container = document.createElement("div"); + const root = createRoot(container); + + const render = async (channelId) => { + await act(async () => { + root.render( + React.createElement( + QueryClientProvider, + { client: queryClient }, + React.createElement(HarnessComponent, { channelId }), + ), + ); + }); + }; + + return { + render, + getFetchOlderArchived: () => hookReturnRef.current?.fetchOlderArchived, + unmount: async () => { + await act(async () => { + root.unmount(); + }); + }, + }; +} + +/** Make a QueryClient pre-seeded with identity so useIdentityQuery resolves. */ +function makeQueryClient() { + const qc = new QueryClient({ defaultOptions: { queries: { retry: false } } }); + qc.setQueryData(["identity"], { pubkey: IDENTITY_PUBKEY }); + return qc; +} + +// ── Settle helper ───────────────────────────────────────────────────────────── +// +// Flushes microtasks + a few macrotask ticks so async effects can settle. +// Uses act() so React commits state updates from effects. + +async function settle(iterations = 3) { + for (let i = 0; i < iterations; i++) { + await act(async () => { + await new Promise((r) => setTimeout(r, 5)); + }); + } +} + +// ── Tests ───────────────────────────────────────────────────────────────────── + +describe("useLoadArchivedObserverEvents — mounted hook lifecycle regressions", () => { + beforeEach(() => { + resetAgentObserverStore(); + clearIpcHandlers(); + _testRegisterKnownAgents(SUB_ID, [AGENT_PUBKEY]); + }); + + /** + * Regression (a): exhausted channel A → switch to channel B. + * + * The original stale-closure bug (pre-round-1): fetchOlderArchived captured + * hasOlderArchived from React state (false for exhausted A). After the switch, + * React state was still false while ps.hasOlderArchived had been reset to true. + * The hydration loop called fetchOlderArchived up to 10 times, each returning + * immediately at the !ps.hasOlderArchived guard — B never got a read. + * + * After the fix (reading ps.hasOlderArchived from the ref), B gets at least + * one read from a null cursor and its rows are ingested. + * + * PROVENANCE: this test is GREEN at dfb2d0385 (the stale-closure was already + * fixed in round 1). It would be red at the pre-round-1 head (884ed9ba2) + * where hasOlderArchived was captured from React state in the closure. + */ + it("test_exhausted_channel_A_switch_to_B_hook_reads_B_from_null_cursor", async () => { + // Channel A: 1 page of 1 row (short page → exhausted immediately). + const aRows = [makeArchivedRow(1, "chan-a")]; + // Channel B: 1 page of 1 row. + const bRows = [makeArchivedRow(2, "chan-b")]; + + const aCalls = []; + const bCalls = []; + + setIpcHandler("list_save_subscriptions", async () => + makeOwnerPSubResponse(), + ); + setIpcHandler("read_unindexed_observer_rows", async () => []); + setIpcHandler("index_observer_channel_id", async () => null); + setIpcHandler("read_archived_observer_events_for_channel", async (args) => { + if (args.channelId === "chan-a") { + aCalls.push({ cursor: args.beforeCreatedAt ?? null }); + return aRows.map((r) => JSON.stringify(r)); + } + if (args.channelId === "chan-b") { + bCalls.push({ cursor: args.beforeCreatedAt ?? null }); + return bRows.map((r) => JSON.stringify(r)); + } + return []; + }); + // decrypt_observer_event is called inside ingestArchivedObserverEvents. + // invokeTauri passes { eventJson: JSON.stringify(rawRelayEvent) }. + // The row.content is the JSON-encoded ObserverEvent — return it parsed. + setIpcHandler("decrypt_observer_event", async (args) => { + try { + const event = JSON.parse(args.eventJson); + return JSON.parse(event.content); + } catch { + return { kind: "telemetry", channelId: null }; + } + }); + + const qc = makeQueryClient(); + const { render, unmount } = mountHook("chan-a", qc); + + // Mount on chan-a and let hydration settle. + await render("chan-a"); + await settle(10); + + // A must have been read (at least one call, from null cursor). + assert.ok(aCalls.length >= 1, `expected A reads, got ${aCalls.length}`); + assert.equal(aCalls[0].cursor, null, "A first read must use null cursor"); + + // Switch to chan-b. + await render("chan-b"); + await settle(10); + + // B must have been read from a null cursor (fresh channel, no inherited cursor). + assert.ok(bCalls.length >= 1, `expected B reads, got ${bCalls.length}`); + assert.equal( + bCalls[0].cursor, + null, + "B first read must use null cursor (not A's cursor)", + ); + + // B's rows must have been ingested into the archive store. + const bArchived = _testGetArchivedChannelEvents(AGENT_PUBKEY, "chan-b"); + assert.ok( + bArchived.length >= 1, + `B's rows must be ingested — found ${bArchived.length} (exhausted-A stale closure bug would leave this 0)`, + ); + + await unmount(); + }); + + /** + * Regression (b): deferred-I/O race — A→B: A's stale writes after ingest AND + * lock theft via stale finally. + * + * Two protections are under test independently: + * + * 1. Post-ingest exhaustion write (post-ingest token check): + * The bug at dfb2d0385: fetchOlderArchived for A checked the token BEFORE + * ingestArchivedObserverEvents but NOT after. A's deferred decrypt resumed + * after the channel switch, wrote ps.hasOlderArchived=false (short-page + * exhaustion) for B's paging state, and B's eager loop stopped after 1 page. + * Removing the post-ingest token check causes bCallCount==1. + * + * 2. Lock theft via generation-gated finally: + * If the generation guard in finally is removed, stale A's finally runs + * ps.isFetching=false while B holds the lock. We probe this directly: + * after resolving stale A (while B's first read is in flight), we call + * fetchOlderArchived() on B ourselves. With the correct guard, B holds + * the lock and the concurrent call returns immediately (bCallCount stays + * at 1 for now). Without it (stale A stole the lock), the concurrent call + * acquires the lock and starts an extra Tauri read (bCallCount jumps to 2 + * prematurely, with concurrent in-flight reads — the assertion catches this + * because it fires BEFORE B's deferred first read resolves). + * + * Precise race sequence: + * 1. Mount on chan-a. A's Tauri read returns 1 row. Cursor set. Ingest starts. + * 2. A's decrypt is DEFERRED (aDecryptDeferred). + * 3. Switch to channel B. + * 4. B's hydration loop starts. B's first Tauri read is ALSO DEFERRED. + * 5. Resolve A's decrypt → A finishes ingest, hits post-ingest check. + * - At dfb2d0385 (no post-ingest check): writes ps.hasOlderArchived=false. + * Also, if finally is unguarded, ps.isFetching=false (lock stolen). + * - After fix: both writes discarded (generation mismatch). + * 6. LOCK PROBE (while B's first read is still deferred): + * call fetchOlderArchived() directly. Must return without starting a new + * Tauri read (B holds the lock; bCallCount must still be 1). + * 7. Resolve B's first Tauri read → B ingests 200 rows. + * 8. B loop continues: bCallCount >= 2. + * + * VERIFIED: removing the post-ingest token check causes bCallCount<2 (step 8). + * Removing just the finally generation guard causes bCallCount>=2 but the lock + * probe at step 6 catches the theft: bCallCount jumps to 2 before B's deferred + * first read resolves (duplicate concurrent read started while B is in flight). + * + * RED at dfb2d0385 (bCallCount==1). GREEN at current head. + */ + it("test_deferred_A_ingest_cannot_exhaust_B_or_steal_B_lock", async () => { + let resolveADecrypt; + const aDecryptDeferred = new Promise((resolve) => { + resolveADecrypt = resolve; + }); + + let resolveBFirstRead; + const bFirstReadDeferred = new Promise((resolve) => { + resolveBFirstRead = resolve; + }); + + const PAGE_SIZE = 200; + const makeBPage = (offset) => + Array.from({ length: PAGE_SIZE }, (_, i) => + makeArchivedRow(offset + i, "chan-b"), + ); + + let bCallCount = 0; + let aDecryptStarted = false; + let bFirstReadHeld = false; + + setIpcHandler("list_save_subscriptions", async () => + makeOwnerPSubResponse(), + ); + setIpcHandler("read_unindexed_observer_rows", async () => []); + setIpcHandler("index_observer_channel_id", async () => null); + setIpcHandler("read_archived_observer_events_for_channel", async (args) => { + if (args.channelId === "chan-a") { + return [JSON.stringify(makeArchivedRow(1, "chan-a"))]; // 1 row = short page + } + if (args.channelId === "chan-b") { + bCallCount++; + if (bCallCount === 1 && !bFirstReadHeld) { + // Defer B's first Tauri read until we explicitly release it. + bFirstReadHeld = true; + await bFirstReadDeferred; + return makeBPage(100).map((r) => JSON.stringify(r)); // full page + } + if (bCallCount <= 4) + return makeBPage(bCallCount * 100).map((r) => JSON.stringify(r)); + return [JSON.stringify(makeArchivedRow(9999, "chan-b"))]; // short = exhaust + } + return []; + }); + setIpcHandler("decrypt_observer_event", async (args) => { + try { + const event = JSON.parse(args.eventJson); + const parsed = JSON.parse(event.content); + if (parsed.channelId === "chan-a" && !aDecryptStarted) { + aDecryptStarted = true; + await aDecryptDeferred; // block A's decrypt + } + return parsed; + } catch { + return { kind: "telemetry", channelId: null }; + } + }); + + const qc = makeQueryClient(); + const { render, getFetchOlderArchived, unmount } = mountHook("chan-a", qc); + + // Step 1-2: Mount on chan-a. A's Tauri read completes (1 row), cursor set, + // ingest starts and blocks at A's decrypt. + await render("chan-a"); + await act(async () => { + await new Promise((r) => setTimeout(r, 20)); + }); + + // Step 3: Switch to chan-b while A's decrypt/ingest is blocked. + await render("chan-b"); + + // Step 4: B calls its first Tauri read and blocks. + await act(async () => { + await new Promise((r) => setTimeout(r, 20)); + }); + + // B's first read must have been attempted (bCallCount >= 1). + assert.ok( + bCallCount >= 1, + `B must have started its first read before the lock probe — got bCallCount=${bCallCount}`, + ); + const bCallCountBeforeAResolve = bCallCount; + + // Step 5: Resolve A's decrypt. At dfb2d0385 this writes + // ps.hasOlderArchived=false (if no post-ingest check) and/or clears + // ps.isFetching (if finally is unguarded). + resolveADecrypt(); + await act(async () => { + await new Promise((r) => setTimeout(r, 10)); + }); + + // Step 6: LOCK PROBE — B still holds the lock (B's first read is deferred). + // Call fetchOlderArchived directly. With correct protection, B's lock is + // intact and this call returns immediately without a Tauri read — bCallCount + // must NOT increase. If A stole the lock, this call acquires it and starts + // a new Tauri read before B's deferred first read resolves (bCallCount jumps). + const fetchFn = getFetchOlderArchived(); + if (fetchFn) { + await act(async () => { + await fetchFn(); + }); + } + + assert.equal( + bCallCount, + bCallCountBeforeAResolve, + `Lock probe must not start a new B read while B holds the lock (bCallCount=${bCallCount}, expected ${bCallCountBeforeAResolve}). If A's stale finally cleared the lock, this concurrent call would start a duplicate read.`, + ); + + // Step 7: Now resolve B's first Tauri read. + resolveBFirstRead(); + + // Let B's loop run. + await settle(10); + + // Step 8: B must have made at least 2 Tauri reads. + // At dfb2d0385: A corrupted ps.hasOlderArchived=false before B's first page + // resolved, so after B's first page, the loop checks and exits. bCallCount==1. + // After fix: A's write was discarded, ps.hasOlderArchived is still true, + // B continues to page 2+. + assert.ok( + bCallCount >= 2, + `B must read at least 2 pages — got ${bCallCount}. Post-ingest token missing at dfb2d0385 let A corrupt B's exhaustion state (bCallCount==1).`, + ); + + // A's row must NOT appear in B's channel archive. + const bArchived = _testGetArchivedChannelEvents(AGENT_PUBKEY, "chan-b"); + for (const evt of bArchived) { + assert.equal( + evt.channelId, + "chan-b", + "B's archive must only contain B-channel events", + ); + } + + await unmount(); + }); + + /** + * Regression (c): concurrent fetches during pending backfill — only one + * same-generation call may acquire the lock after backfill resolves. + * + * The gap at 4c92a018d: `ps.isFetching` was checked only at the top of + * fetchOlderArchived, BEFORE `await ps.backfillPromise`. Two callers + * (eager hydration loop + a concurrent scroll trigger) could both observe + * isFetching=false, both suspend on the same pending backfill promise, then + * both resume and proceed past the post-backfill guard (which only checked + * generation + exhaustion). Both would set isFetching=true and issue + * readArchivedObserverEventsForChannel from the SAME cursor — duplicating + * the first page. + * + * Fix: after `await ps.backfillPromise`, recheck `ps.isFetching` immediately + * before acquiring the lock. The second caller finds the lock already taken + * and returns without reading. + * + * The test defers the ARCHIVE response (not just backfill) so we can count + * reads while the winner's first response is still in flight. If both + * callers acquired the lock, archiveCallCount will be 2 before the first + * deferred response is released. With the fix, archiveCallCount is 1. + * + * Precise race sequence: + * 1. Mount on chan-a. Backfill is DEFERRED (backfillDeferred). + * Archive responses are also deferred until resolveArchive() is called. + * 2. Hydration loop call #1 suspends on backfill. + * 3. Inject manual call #2 — it also sees isFetching=false and suspends on + * backfill. + * 4. Resolve backfill. Both calls resume and race to acquire the lock. + * - Without fix: both pass the post-backfill guard, both set + * isFetching=true, both issue the Tauri read — archiveCallCount == 2. + * - With fix: one passes, sets isFetching=true; the other sees the lock + * taken and returns — archiveCallCount == 1. + * 5. Assert archiveCallCount == 1 before releasing archive response. + * (Archive is still deferred, so loop hasn't advanced past page 1 yet — + * any count > 1 is purely from the concurrent race, not loop progress.) + * + * VERIFIED: removing the post-backfill `ps.isFetching` recheck causes + * archiveCallCount == 2 at step 5. GREEN at current head. + */ + it("test_two_concurrent_fetches_during_backfill_only_one_proceeds", async () => { + let resolveBackfill; + const backfillDeferred = new Promise((resolve) => { + resolveBackfill = resolve; + }); + + // Defer ALL archive responses until we release them. This way, if two + // callers both acquire the lock, archiveCallCount jumps to 2 before we + // release the response — and we can catch it unambiguously. + let resolveArchive; + const archiveDeferred = new Promise((resolve) => { + resolveArchive = resolve; + }); + + let archiveCallCount = 0; + + setIpcHandler("list_save_subscriptions", async () => + makeOwnerPSubResponse(), + ); + // Defer readUnindexedObserverRows to simulate a pending backfill. + setIpcHandler("read_unindexed_observer_rows", async () => { + await backfillDeferred; + return []; + }); + setIpcHandler("index_observer_channel_id", async () => null); + setIpcHandler("read_archived_observer_events_for_channel", async (args) => { + if (args.channelId === "chan-a") { + archiveCallCount++; + // Hold this response until we explicitly release it so we can + // count concurrent reads before any result is returned. + await archiveDeferred; + return Array.from({ length: 200 }, (_, i) => + JSON.stringify(makeArchivedRow(i, "chan-a")), + ); + } + return []; + }); + setIpcHandler("decrypt_observer_event", async (args) => { + try { + const event = JSON.parse(args.eventJson); + return JSON.parse(event.content); + } catch { + return { kind: "telemetry", channelId: null }; + } + }); + + const qc = makeQueryClient(); + const { render, getFetchOlderArchived, unmount } = mountHook("chan-a", qc); + + // Step 1-2: Mount on chan-a. Backfill is in flight (deferred). + // Hydration loop call #1 enters fetchOlderArchived and suspends on backfill. + await render("chan-a"); + await act(async () => { + await new Promise((r) => setTimeout(r, 20)); + }); + + // Step 3: Inject call #2 while backfill is still pending and call #1 is + // suspended. Both observe isFetching=false here. + const fetchFn = getFetchOlderArchived(); + let call2Promise; + if (fetchFn) { + // Do NOT await yet — let it run concurrently with call #1. + call2Promise = fetchFn(); + } + + // Let call #2 reach its backfill await before we resolve backfill. + await act(async () => { + await new Promise((r) => setTimeout(r, 5)); + }); + + // Step 4: Resolve backfill. Both calls resume and race to acquire lock. + resolveBackfill(); + + // Yield to let both calls advance past the post-backfill guard and issue + // their Tauri reads (or be blocked by the lock recheck). + await act(async () => { + await new Promise((r) => setTimeout(r, 10)); + }); + + // Step 5: Archive responses are still deferred — loop cannot have advanced + // past page 1. Any archiveCallCount > 1 here is purely from concurrent + // reads racing through the backfill await without a lock recheck. + assert.equal( + archiveCallCount, + 1, + `Exactly one archive read must start after backfill resolves — got ${archiveCallCount}. Without the post-backfill isFetching recheck, both concurrent callers acquire the lock and both issue a Tauri read (archiveCallCount == 2).`, + ); + + // Release archive responses so the lock holder can finish and the test + // can unmount cleanly. + resolveArchive(); + + // Wait for call #2 to settle as well. + if (call2Promise) { + await act(async () => { + await call2Promise; + }); + } + + await settle(5); + await unmount(); + }); + + /** + * Regression (d): A→B→A — old-A's stale in-flight request must not corrupt + * fresh-A's paging state when the user returns to channel A. + * + * The gap in the channel-string equality check (af45fbf05): using + * `requestChannelId === ps.activeChannelId` as the ownership test fails when + * the user navigates A→B→A. Old-A's request sees `ps.activeChannelId === "A"` + * after the second return to A — so every check passes. Old-A can mark fresh-A + * exhausted and its finally releases fresh-A's lock. + * + * With resetGeneration each applyChannelReset() call increments the counter, + * so A(gen=1) → B(gen=2) → A(gen=3): old-A snapshotted gen=1, which never + * equals gen=3, so all writes are discarded regardless of channel name. + * + * Race sequence: + * 1. Mount on chan-a (gen=1). A's Tauri read = 1 row (short page). Ingest + * starts. A's decrypt is DEFERRED. + * 2. Switch to chan-b (gen=2). B's paging starts. + * 3. Switch BACK to chan-a (gen=3). Fresh-A's hydration starts. Fresh-A's + * first Tauri read is DEFERRED (freshAReadDeferred). + * 4. Resolve old-A's decrypt. Old-A hits short-page branch. + * - Without generation: old-A sees ps.activeChannelId==="chan-a", writes + * ps.hasOlderArchived=false and clears ps.isFetching. + * - With generation: gen=1 !== gen=3, writes discarded. + * 5. LOCK PROBE: call fetchOlderArchived directly. Fresh-A holds the lock + * (its deferred read is in flight). With correct protection the probe + * returns immediately (freshACallCount unchanged). Without it (old-A's + * finally stole the lock), the probe starts a duplicate read. + * 6. Resolve fresh-A's first read (full page). Fresh-A loop continues. + * 7. Assert freshACallCount >= 2 (fresh-A ran past page 1). + * + * VERIFIED: this test is RED when activeChannelId string equality replaces + * resetGeneration (old-A passes every check, marks fresh-A exhausted at step 4). + * GREEN at current head. + */ + it("test_A_B_A_old_request_cannot_corrupt_fresh_A_state", async () => { + let resolveOldADecrypt; + const oldADecryptDeferred = new Promise((resolve) => { + resolveOldADecrypt = resolve; + }); + + let resolveFreshAFirstRead; + const freshAFirstReadDeferred = new Promise((resolve) => { + resolveFreshAFirstRead = resolve; + }); + + const PAGE_SIZE = 200; + const makePage = (channelId, offset) => + Array.from({ length: PAGE_SIZE }, (_, i) => + makeArchivedRow(offset + i, channelId), + ); + + let oldADecryptStarted = false; + // Track calls per channel / phase. We only care about chan-a reads on fresh-A. + let freshACallCount = 0; + // After we switch back to A (gen=3), track reads for that phase. + let onFreshA = false; + + setIpcHandler("list_save_subscriptions", async () => + makeOwnerPSubResponse(), + ); + setIpcHandler("read_unindexed_observer_rows", async () => []); + setIpcHandler("index_observer_channel_id", async () => null); + setIpcHandler("read_archived_observer_events_for_channel", async (args) => { + if (args.channelId === "chan-a") { + if (!onFreshA) { + // Old-A's read: 1 row = short page. + return [JSON.stringify(makeArchivedRow(1, "chan-a"))]; + } + // Fresh-A's reads. + freshACallCount++; + if (freshACallCount === 1) { + // Defer fresh-A's first read. + await freshAFirstReadDeferred; + return makePage("chan-a", 200).map((r) => JSON.stringify(r)); // full + } + if (freshACallCount <= 4) + return makePage("chan-a", freshACallCount * 200).map((r) => + JSON.stringify(r), + ); + return [JSON.stringify(makeArchivedRow(9999, "chan-a"))]; // exhaust + } + if (args.channelId === "chan-b") { + // B gets one short page (we don't care about B's progress here). + return [JSON.stringify(makeArchivedRow(50, "chan-b"))]; + } + return []; + }); + setIpcHandler("decrypt_observer_event", async (args) => { + try { + const event = JSON.parse(args.eventJson); + const parsed = JSON.parse(event.content); + if (parsed.channelId === "chan-a" && !oldADecryptStarted && !onFreshA) { + oldADecryptStarted = true; + await oldADecryptDeferred; // block OLD A's decrypt + } + return parsed; + } catch { + return { kind: "telemetry", channelId: null }; + } + }); + + const qc = makeQueryClient(); + const { render, getFetchOlderArchived, unmount } = mountHook("chan-a", qc); + + // Step 1: Mount on chan-a (gen=1). Old-A's Tauri read returns 1 row. + // Ingest starts, decrypt blocks. + await render("chan-a"); + await act(async () => { + await new Promise((r) => setTimeout(r, 20)); + }); + + // Step 2: Switch to chan-b (gen=2). + await render("chan-b"); + await act(async () => { + await new Promise((r) => setTimeout(r, 10)); + }); + + // Step 3: Switch back to chan-a (gen=3). Fresh-A hydration starts. + onFreshA = true; + await render("chan-a"); + await act(async () => { + await new Promise((r) => setTimeout(r, 20)); + }); + + // Fresh-A must have started its first (deferred) read. + assert.ok( + freshACallCount >= 1, + `fresh-A must have started its first read — got freshACallCount=${freshACallCount}`, + ); + const freshACountBeforeOldResolve = freshACallCount; + + // Step 4: Resolve old-A's decrypt. Without generation check, old-A's + // post-ingest branch writes ps.hasOlderArchived=false (marks fresh-A + // exhausted) and its finally clears ps.isFetching (steals fresh-A's lock). + resolveOldADecrypt(); + await act(async () => { + await new Promise((r) => setTimeout(r, 10)); + }); + + // Step 5: LOCK PROBE — fresh-A holds the lock (first read deferred). + // A concurrent fetchOlderArchived call must be rejected (lock held). + // If old-A stole the lock, this probe starts a duplicate read (freshACallCount + // jumps before the deferred first read resolves). + const fetchFn = getFetchOlderArchived(); + if (fetchFn) { + await act(async () => { + await fetchFn(); + }); + } + + assert.equal( + freshACallCount, + freshACountBeforeOldResolve, + `Lock probe must not start a new fresh-A read while fresh-A holds the lock (freshACallCount=${freshACallCount}, expected ${freshACountBeforeOldResolve}). Old-A's stale finally stole the lock (A→B→A channel-string equality bug).`, + ); + + // Step 6: Resolve fresh-A's first read (full page). Loop continues. + resolveFreshAFirstRead(); + await settle(10); + + // Step 7: Fresh-A must have made at least 2 reads (loop continued past page 1). + // Without generation check, old-A wrote ps.hasOlderArchived=false before + // fresh-A's first page resolved, causing the loop to exit — freshACallCount==1. + assert.ok( + freshACallCount >= 2, + `fresh-A must read at least 2 pages — got ${freshACallCount}. Old-A's stale writes (A→B→A) would leave freshACallCount==1.`, + ); + + await unmount(); + }); +}); diff --git a/desktop/src/features/agents/ui/useObserverEvents.ts b/desktop/src/features/agents/ui/useObserverEvents.ts index 8f4ecc103..0c44a64d5 100644 --- a/desktop/src/features/agents/ui/useObserverEvents.ts +++ b/desktop/src/features/agents/ui/useObserverEvents.ts @@ -21,6 +21,7 @@ import type { RelayEvent } from "@/shared/api/types"; import { createArchivePagingState, applyChannelReset, + runHydrationLoop, } from "./archivePagingState"; export type { ArchivePagingState } from "./archivePagingState"; @@ -88,7 +89,12 @@ export function useArchivedChannelEvents( return React.useSyncExternalStore(subscribeToStore, getSnapshot); } -const ARCHIVED_EVENTS_PAGE_SIZE = 50; +const ARCHIVED_EVENTS_PAGE_SIZE = 200; + +// Number of pages to load eagerly on panel open (before any scroll). Each page +// is ARCHIVED_EVENTS_PAGE_SIZE frames; 10 pages = 2000 frames, which covers +// agent turns that emit hundreds of frames (e.g. a full code-review turn ~900). +const INITIAL_HYDRATION_BUDGET_PAGES = 10; /** * Load-older-on-scroll for archived observer frames, scoped to a single channel. @@ -138,10 +144,13 @@ export function useLoadArchivedObserverEvents( // Reset per-channel paging state when channelId changes. Backfill state is // identity-level (not per-channel) and must NOT be reset here — the backfill // index covers all channels and only needs to run once per identity mount. - // Only the cursor, exhaustion flag, and fetching lock are channel-scoped. + // Only the cursor, exhaustion flag, fetching lock, channel label, and + // resetGeneration are channel-scoped. resetGeneration is incremented by + // applyChannelReset so in-flight reads from any prior reset (including + // A→B→A) detect staleness and discard their results. // biome-ignore lint/correctness/useExhaustiveDependencies: channelId is the intentional reset key; ps is a stable ref excluded from deps by convention; setHasOlderArchived is a stable React state setter React.useEffect(() => { - applyChannelReset(ps); + applyChannelReset(ps, channelId); setHasOlderArchived(true); }, [channelId]); @@ -259,7 +268,7 @@ export function useLoadArchivedObserverEvents( ps.backfillPromise = promise; }, [enabled, hasSubscription]); - // biome-ignore lint/correctness/useExhaustiveDependencies: ps is a stable ref; ps.isFetching/ps.cursor/ps.backfillPromise/ps.hasOlderArchived are read via the stable ref object, not reactive values + // biome-ignore lint/correctness/useExhaustiveDependencies: ps is a stable ref; all per-page state (isFetching, cursor, hasOlderArchived, resetGeneration) is read from the ref rather than React state, so this callback is intentionally stable across exhaustion/channel changes const fetchOlderArchived = React.useCallback(async () => { if ( !enabled || @@ -267,11 +276,18 @@ export function useLoadArchivedObserverEvents( !hasSubscription || !channelId || ps.isFetching || - !hasOlderArchived + !ps.hasOlderArchived ) { return; } + // Snapshot the reset generation at the start of this request. Every + // shared-state write below rechecks requestGeneration === ps.resetGeneration + // first. A mismatch means at least one channel switch occurred while we + // were awaiting async I/O — even A→B→A is detected because each switch + // increments resetGeneration. Channel ID is kept only as the query input. + const requestGeneration = ps.resetGeneration; + // 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 @@ -280,12 +296,21 @@ export function useLoadArchivedObserverEvents( await ps.backfillPromise; } - // Re-check after awaiting: hasOlderArchived might have been set false - // while we were waiting (e.g. subscription check failed). - if (!hasOlderArchived) { + // Re-check after awaiting: generation may have advanced (channel switched), + // archive exhausted, or another concurrent caller may have acquired the + // fetch lock while we were suspended on backfill. All three must be + // re-evaluated because any of them could have changed mid-await. + if ( + !ps.hasOlderArchived || + requestGeneration !== ps.resetGeneration || + ps.isFetching + ) { return; } + // Acquire the fetch lock under this request's generation. The finally + // block only releases the lock if the generation still matches — so a stale + // in-flight request cannot clear the lock that belongs to a later reset. ps.isFetching = true; try { const before = ps.cursor ?? undefined; @@ -294,6 +319,13 @@ export function useLoadArchivedObserverEvents( limit: ARCHIVED_EVENTS_PAGE_SIZE, }); + // Discard result if the generation advanced while the Tauri read was in + // flight (channel switch, including A→B→A). The new channel will start + // its own read with a null cursor. + if (requestGeneration !== ps.resetGeneration) { + return; + } + if (events.length > 0) { // Cursor = the last row in newest-first order = the oldest event on // this page. Capture both created_at and id to mirror the compound @@ -306,6 +338,14 @@ export function useLoadArchivedObserverEvents( await ingestArchivedObserverEvents(events); } + // Re-check generation after ingestArchivedObserverEvents: ingestion + // decrypts each frame asynchronously and may take time. If a channel + // switch occurred during that await (including A→B→A), discard all + // remaining shared-state writes — exhaustion and React mirror. + if (requestGeneration !== ps.resetGeneration) { + return; + } + // A short page means the archive is exhausted for this channel. if (events.length < ARCHIVED_EVENTS_PAGE_SIZE) { setHasOlderArchived(false); @@ -314,9 +354,56 @@ export function useLoadArchivedObserverEvents( } catch (error) { console.error("[useLoadArchivedObserverEvents] fetch failed:", error); } finally { - ps.isFetching = false; + // Only release the fetch lock if this request still owns it. If the + // generation advanced (any channel switch including A→B→A), the new + // channel acquired its own lock — releasing here would steal it. + if (requestGeneration === ps.resetGeneration) { + ps.isFetching = false; + } } - }, [enabled, identityPubkey, hasSubscription, channelId, hasOlderArchived]); + }, [enabled, identityPubkey, hasSubscription, channelId]); + + // Eager initial hydration: on panel open (or channel switch), load archive + // pages automatically until the budget is reached or the channel is exhausted. + // This makes archived history visible immediately without any scrolling. + // + // Runs when: enabled + subscription confirmed + channelId resolved + + // hydration not yet done for this channel. Respects `applyChannelReset` + // (which resets initialHydrationDone) so channel switches trigger a fresh + // pass. Uses fetchOlderArchived's existing lock/cursor/backfill-await + // machinery — no parallel state machine. + // + // fetchOlderArchived is now stable (it no longer captures hasOlderArchived + // from React state — it reads ps.hasOlderArchived from the ref), so it is + // safe to call from this effect without coupling the hydration lifecycle to + // React state identity changes. + // biome-ignore lint/correctness/useExhaustiveDependencies: ps is a stable ref; initialHydrationDone is read from ps (not as a reactive dep) to avoid triggering re-runs; fetchOlderArchived is stable and intentionally omitted + React.useEffect(() => { + if ( + !enabled || + !identityPubkey || + !hasSubscription || + !channelId || + ps.initialHydrationDone + ) { + return; + } + + // Mark done immediately to prevent concurrent hydration loops. The loop + // runs asynchronously; the signal object handles mid-loop cancellation on + // channel switch (the cleanup fn sets signal.cancelled = true). + ps.initialHydrationDone = true; + const signal = { cancelled: false }; + void runHydrationLoop( + ps, + fetchOlderArchived, + INITIAL_HYDRATION_BUDGET_PAGES, + signal, + ); + return () => { + signal.cancelled = true; + }; + }, [enabled, identityPubkey, hasSubscription, channelId]); return { fetchOlderArchived, hasOlderArchived }; } From 80244f82318c85f931d1055e419456945c5eca99 Mon Sep 17 00:00:00 2001 From: Logan Johnson Date: Thu, 23 Jul 2026 12:45:00 -0700 Subject: [PATCH 3/4] Fix avatar upload lifecycle edge cases (#2277) Signed-off-by: Logan Johnson Signed-off-by: npub1dpf98sl35hm9k65t8h6cvh5n6knn2msugh5k5nmysxwnw5wlh7uqn2wgp4 <685253c3f1a5f65b6a8b3df5865e93d5a7356e1c45e96a4f64819d3751dfbfb8@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub13n66s06epmqf2kc3v373ez8hj65cuzyvxzjf93vwpervxqn2u7jq2qd9je <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> Co-authored-by: npub1dpf98sl35hm9k65t8h6cvh5n6knn2msugh5k5nmysxwnw5wlh7uqn2wgp4 <685253c3f1a5f65b6a8b3df5865e93d5a7356e1c45e96a4f64819d3751dfbfb8@sprout-oss.stage.blox.sqprod.co> Co-authored-by: npub13n66s06epmqf2kc3v373ez8hj65cuzyvxzjf93vwpervxqn2u7jq2qd9je <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> --- AGENTS.md | 5 +- desktop/src-tauri/src/commands/profile.rs | 135 +++++- desktop/src-tauri/src/lib.rs | 1 + desktop/src-tauri/src/relay.rs | 53 +-- desktop/src-tauri/src/relay/submit.rs | 60 +++ desktop/src/app/App.tsx | 3 + .../features/communities/useCommunityInit.ts | 22 +- .../onboarding/ui/CommunityOnboardingFlow.tsx | 34 +- .../profile/avatarPresentationStore.ts | 39 +- .../profile/avatarProfileSync.test.mjs | 303 ++++++++++++ .../src/features/profile/avatarProfileSync.ts | 262 ++++++++++ .../profile/profileCacheSync.test.mjs | 91 ++++ .../src/features/profile/profileCacheSync.ts | 91 ++++ .../profile/ui/ProfileAvatarEditor.tsx | 24 +- desktop/src/shared/api/tauriProfiles.ts | 13 + desktop/src/testing/e2eBridge.ts | 7 + desktop/tests/e2e/onboarding.spec.ts | 450 ++++++++++++++++-- 17 files changed, 1477 insertions(+), 116 deletions(-) create mode 100644 desktop/src-tauri/src/relay/submit.rs create mode 100644 desktop/src/features/profile/avatarProfileSync.test.mjs create mode 100644 desktop/src/features/profile/avatarProfileSync.ts create mode 100644 desktop/src/features/profile/profileCacheSync.test.mjs create mode 100644 desktop/src/features/profile/profileCacheSync.ts diff --git a/AGENTS.md b/AGENTS.md index c94ba3881..c9761d4ea 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -475,13 +475,16 @@ class instances, cached promises) survive across remounts. Every community-scope singleton needs a reset function wired into `resetCommunityState()` in `desktop/src/features/communities/useCommunityInit.ts`. -Current singletons that are reset on community switch: +Current singletons that are reset on relay boundary changes (same-relay +reconnects preserve pending avatar verification work): - `relayClient.disconnect()` — WebSocket teardown + promise rejection - `resetRateLimitGate()` — clears any active rate-limit window from the old relay - `clearAllDrafts()` — message draft cache - `resetAgentObserverStore()` — agent observer relay store - `resetActiveAgentTurnsStore()` — active agent turn timers - `resetAgentWorkingSignal()` — agent working indicator signal +- `resetAvatarProfileSync()` — pending verified-avatar profile writes +- `resetAvatarPresentations()` — avatar probes, previews, and Retry toasts - `resetSidebarRelayConnectionCardState()` — sidebar relay card dismiss state - `resetMediaCaches()` — proxy port and relay origin caches - `resetVideoPlayerState()` — video player singleton diff --git a/desktop/src-tauri/src/commands/profile.rs b/desktop/src-tauri/src/commands/profile.rs index ab36b45c5..ef67fac57 100644 --- a/desktop/src-tauri/src/commands/profile.rs +++ b/desktop/src-tauri/src/commands/profile.rs @@ -7,9 +7,13 @@ use tauri::State; use crate::{ app_state::AppState, events, + managed_agents::persona_events::monotonic_created_at, models::{ProfileInfo, SearchUsersResponse, UserNotesResponse, UsersBatchResponse}, nostr_convert, - relay::{query_relay, submit_event}, + relay::{ + query_relay, query_relay_at_with_keys, relay_http_base_url, submit_event, + submit_event_at_with_keys, + }, }; #[tauri::command] @@ -94,6 +98,85 @@ pub async fn update_profile( .unwrap_or_else(|| empty_profile_info(¤t_pubkey_hex_unwrap(&state)))) } +#[tauri::command] +pub async fn update_profile_at_relay( + relay_url: String, + expected_pubkey: String, + expected_avatar_url: Option, + avatar_url: String, + state: State<'_, AppState>, +) -> Result { + let signer = capture_expected_signer(&state, &expected_pubkey)?; + + let api_base_url = relay_http_base_url(&relay_url); + let filter = serde_json::json!({ + "kinds": [0], + "authors": [expected_pubkey], + "limit": 1 + }); + let prior_events = query_relay_at_with_keys( + &state, + &api_base_url, + std::slice::from_ref(&filter), + &signer, + None, + ) + .await?; + let prior_event = prior_events.first(); + let current: Value = prior_event + .and_then(|event| serde_json::from_str::(&event.content).ok()) + .unwrap_or(Value::Null); + let current_avatar_url = current + .get("picture") + .and_then(Value::as_str) + .map(str::to_string); + if normalized_avatar_url(current_avatar_url.as_deref()) + != normalized_avatar_url(expected_avatar_url.as_deref()) + { + return Err("profile avatar changed before deferred save".to_string()); + } + + let builder = build_deferred_profile_event(¤t, &avatar_url, prior_event)?; + submit_event_at_with_keys(builder, &state, &api_base_url, &signer).await?; + + let events = query_relay_at_with_keys(&state, &api_base_url, &[filter], &signer, None).await?; + Ok(events + .first() + .map(nostr_convert::profile_info_from_event) + .transpose()? + .unwrap_or_else(|| empty_profile_info(&expected_pubkey))) +} + +fn build_deferred_profile_event( + current: &Value, + avatar_url: &str, + prior_event: Option<&nostr::Event>, +) -> Result { + let display_name = current.get("display_name").and_then(Value::as_str); + let name = current.get("name").and_then(Value::as_str); + let about = current.get("about").and_then(Value::as_str); + let nip05 = current.get("nip05").and_then(Value::as_str); + + Ok( + events::build_profile(display_name, name, Some(avatar_url), about, nip05)? + .custom_created_at(monotonic_created_at( + prior_event.map(|event| event.created_at.as_secs() as i64), + )), + ) +} + +fn capture_expected_signer(state: &AppState, expected_pubkey: &str) -> Result { + let signer = state.signing_keys()?; + if signer.public_key().to_hex() != expected_pubkey { + return Err("profile identity changed before avatar save".to_string()); + } + Ok(signer) +} + +fn normalized_avatar_url(avatar_url: Option<&str>) -> Option<&str> { + avatar_url.map(str::trim).filter(|value| !value.is_empty()) +} + #[tauri::command] pub async fn get_user_profile( pubkey: Option, @@ -335,6 +418,56 @@ fn empty_profile_info(pubkey: &str) -> ProfileInfo { mod tests { use super::*; + #[test] + fn deferred_profile_signer_is_captured_and_rejects_wrong_identity() { + let state = crate::app_state::build_app_state(); + let original = state.signing_keys().expect("signable identity"); + let original_pubkey = original.public_key().to_hex(); + + let captured = capture_expected_signer(&state, &original_pubkey) + .expect("matching identity should be captured"); + *state.keys.lock().expect("lock keys") = nostr::Keys::generate(); + + assert_eq!(captured.public_key().to_hex(), original_pubkey); + assert_ne!( + state.keys.lock().expect("lock keys").public_key().to_hex(), + original_pubkey + ); + assert_eq!( + capture_expected_signer(&state, &original_pubkey).unwrap_err(), + "profile identity changed before avatar save" + ); + } + + #[test] + fn deferred_profile_event_is_strictly_newer_than_prior_head() { + let keys = nostr::Keys::generate(); + let prior_created_at = nostr::Timestamp::now().as_secs() + 60; + let prior_event = nostr::EventBuilder::new( + nostr::Kind::Metadata, + serde_json::json!({"display_name": "Larry"}).to_string(), + ) + .custom_created_at(nostr::Timestamp::from(prior_created_at)) + .sign_with_keys(&keys) + .expect("sign prior profile"); + + let builder = build_deferred_profile_event( + &serde_json::json!({"display_name": "Larry"}), + "https://example.com/avatar.png", + Some(&prior_event), + ) + .expect("build deferred profile"); + let event = builder + .sign_with_keys(&keys) + .expect("sign deferred profile"); + + assert_eq!(event.created_at.as_secs(), prior_created_at + 1); + assert_eq!( + serde_json::from_str::(&event.content).unwrap()["picture"], + "https://example.com/avatar.png" + ); + } + #[test] fn user_search_filter_requests_prefix_mode_for_typeahead() { // Every caller of `search_users` is a typeahead surface. Whole-word diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index 9017356ab..48e897b19 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -667,6 +667,7 @@ pub fn run() { persist_current_identity, get_profile, update_profile, + update_profile_at_relay, get_user_profile, get_users_batch, get_user_notes, diff --git a/desktop/src-tauri/src/relay.rs b/desktop/src-tauri/src/relay.rs index e6463f509..1c9ba0095 100644 --- a/desktop/src-tauri/src/relay.rs +++ b/desktop/src-tauri/src/relay.rs @@ -347,6 +347,7 @@ pub async fn query_relay_at_with_keys( keys: &Keys, auth_tag: Option<&str>, ) -> Result, String> { + crate::relay_admission::wait_for_rate_limit().await; let url = format!("{}/query", api_base_url); let body_bytes = serde_json::to_vec(filters).map_err(|e| format!("filter serialization failed: {e}"))?; @@ -530,56 +531,8 @@ pub struct AgentProfileInfo { // ── Signed-event submission ───────────────────────────────────────────────── -/// Response from `POST /events`. -#[derive(Debug, Deserialize, serde::Serialize)] -pub struct SubmitEventResponse { - pub event_id: String, - pub accepted: bool, - pub message: String, -} - -/// Build an `EventBuilder` from the events module, sign it with the user's keys, -/// and POST the signed event to `/events` with NIP-98 auth. -pub async fn submit_event( - builder: nostr::EventBuilder, - state: &AppState, -) -> Result { - crate::relay_admission::wait_for_rate_limit().await; - // All synchronous work (signing) must complete before any .await - // so the MutexGuard is dropped and the future remains Send. - let url = format!("{}/events", relay_api_base_url_with_override(state)); - let (auth_header, body_bytes) = { - let keys = state.signing_keys()?; - let event = builder - .sign_with_keys(&keys) - .map_err(|e| format!("failed to sign event: {e}"))?; - let body = event.as_json().into_bytes(); - let auth = build_nip98_auth_header_for_keys(&keys, &Method::POST, &url, &body)?; - (auth, body) - }; // keys dropped here - - let response = state - .http_client - .post(&url) - .header("Authorization", auth_header) - .header("Content-Type", "application/json") - .body(body_bytes) - .send() - .await - .map_err(|e| classify_request_error(&e))?; - - if !response.status().is_success() { - return Err(relay_error_message(response).await); - } - - let result: SubmitEventResponse = parse_json_response(response).await?; - - if !result.accepted { - return Err(format!("relay rejected event: {}", result.message)); - } - - Ok(result) -} +mod submit; +pub use submit::{submit_event, submit_event_at_with_keys, SubmitEventResponse}; /// POST an already-signed event to `/events` with NIP-98 auth. /// diff --git a/desktop/src-tauri/src/relay/submit.rs b/desktop/src-tauri/src/relay/submit.rs new file mode 100644 index 000000000..7fb3f9404 --- /dev/null +++ b/desktop/src-tauri/src/relay/submit.rs @@ -0,0 +1,60 @@ +use super::*; + +/// Response from `POST /events`. +#[derive(Debug, Deserialize, serde::Serialize)] +pub struct SubmitEventResponse { + pub event_id: String, + pub accepted: bool, + pub message: String, +} + +/// Sign with an explicit identity and POST the event to an explicit relay. +/// +/// The caller owns the signer lifetime. This is important for deferred work: +/// an in-process identity swap cannot retarget the event or its NIP-98 auth +/// after the caller has validated which identity the operation belongs to. +pub async fn submit_event_at_with_keys( + builder: nostr::EventBuilder, + state: &AppState, + api_base_url: &str, + keys: &nostr::Keys, +) -> Result { + crate::relay_admission::wait_for_rate_limit().await; + let url = format!("{}/events", api_base_url.trim_end_matches('/')); + let event = builder + .sign_with_keys(keys) + .map_err(|e| format!("failed to sign event: {e}"))?; + let body_bytes = event.as_json().into_bytes(); + let auth_header = build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes)?; + + let response = state + .http_client + .post(&url) + .header("Authorization", auth_header) + .header("Content-Type", "application/json") + .body(body_bytes) + .send() + .await + .map_err(|e| classify_request_error(&e))?; + + if !response.status().is_success() { + return Err(relay_error_message(response).await); + } + + let result: SubmitEventResponse = parse_json_response(response).await?; + if !result.accepted { + return Err(format!("relay rejected event: {}", result.message)); + } + + Ok(result) +} + +/// Build and submit an event to the currently active workspace relay. +pub async fn submit_event( + builder: nostr::EventBuilder, + state: &AppState, +) -> Result { + let api_base_url = relay_api_base_url_with_override(state); + let keys = state.signing_keys()?; + submit_event_at_with_keys(builder, state, &api_base_url, &keys).await +} diff --git a/desktop/src/app/App.tsx b/desktop/src/app/App.tsx index e89bf61fb..ed9c465bc 100644 --- a/desktop/src/app/App.tsx +++ b/desktop/src/app/App.tsx @@ -44,6 +44,7 @@ import { import { WelcomeSetup } from "@/features/communities/ui/WelcomeSetup"; import { CommunityApplyErrorScreen } from "@/features/communities/ui/CommunityApplyErrorScreen"; import { CommunityChangeOverlay } from "@/features/communities/ui/CommunityChangeOverlay"; +import { setAvatarProfileSyncQueryClient } from "@/features/profile/avatarProfileSync"; import { createBuzzQueryClient } from "@/shared/api/queryClient"; import { isSharedIdentity as isSharedIdentityCmd } from "@/shared/api/tauri"; import { getProfile } from "@/shared/api/tauriProfiles"; @@ -197,6 +198,8 @@ function CommunitySwitchGate() { function CommunityQueryProvider({ children }: { children: ReactNode }) { const [queryClient] = useState(createBuzzQueryClient); + useEffect(() => setAvatarProfileSyncQueryClient(queryClient), [queryClient]); + useEffect(() => { const e2eWindow = window as Window & { __BUZZ_E2E__?: unknown; diff --git a/desktop/src/features/communities/useCommunityInit.ts b/desktop/src/features/communities/useCommunityInit.ts index a50a8665c..d46c8d62a 100644 --- a/desktop/src/features/communities/useCommunityInit.ts +++ b/desktop/src/features/communities/useCommunityInit.ts @@ -19,6 +19,8 @@ import { } from "@/features/agents/activeAgentTurnsStore"; import { resetAgentWorkingSignal } from "@/features/agents/agentWorkingSignal"; import { resetAgentObserverStore } from "@/features/agents/observerRelayStore"; +import { resetAvatarPresentations } from "@/features/profile/avatarPresentationStore"; +import { resetAvatarProfileSync } from "@/features/profile/avatarProfileSync"; import { resetSidebarRelayConnectionCardState } from "@/features/sidebar/ui/useSidebarRelayConnectionCard"; import { clearMarkdownNodeCache } from "@/shared/ui/markdown/nodeCache"; import { resetVideoPlayerState } from "@/shared/ui/videoPlayerState"; @@ -33,13 +35,21 @@ import type { Community } from "./types"; * destroyed via effect cleanup and do not need entries here. * See AGENTS.md "Community Switching" for the full contract. */ -function resetCommunityState(): void { +function resetCommunityState({ + resetAvatarState, +}: { + resetAvatarState: boolean; +}): void { relayClient.disconnect(); resetRateLimitGate(); clearAllDrafts(); resetAgentObserverStore(); resetActiveAgentTurnsStore(); resetAgentWorkingSignal(); + if (resetAvatarState) { + resetAvatarProfileSync(); + resetAvatarPresentations(); + } resetSidebarRelayConnectionCardState(); resetMediaCaches(); resetVideoPlayerState(); @@ -84,6 +94,10 @@ export function useCommunityInit( // Track the previously-applied community ID so we can save its turn state // before resetting when the user switches to a different community. const prevCommunityIdRef = useRef(null); + // Deferred avatar work owns the relay captured when it was queued. A + // same-relay reconnect during onboarding must not cancel that work, while an + // actual relay boundary must clear both the queue and its presentation probe. + const appliedRelayUrlRef = useRef(null); // biome-ignore lint/correctness/useExhaustiveDependencies: we intentionally depend on specific properties (id/relayUrl/token/reposDir) — depending on the whole object would trigger resets on name-only changes useEffect(() => { @@ -145,9 +159,13 @@ export function useCommunityInit( // store under the outgoing community ID and delete its snapshot. prevCommunityIdRef.current = null; } - resetCommunityState(); + resetCommunityState({ + resetAvatarState: + appliedRelayUrlRef.current !== activeCommunity.relayUrl, + }); } hasInitializedRef.current = true; + appliedRelayUrlRef.current = activeCommunity.relayUrl; // Apply community config to the Tauri backend. // diff --git a/desktop/src/features/onboarding/ui/CommunityOnboardingFlow.tsx b/desktop/src/features/onboarding/ui/CommunityOnboardingFlow.tsx index 10c194ce6..4b729ab33 100644 --- a/desktop/src/features/onboarding/ui/CommunityOnboardingFlow.tsx +++ b/desktop/src/features/onboarding/ui/CommunityOnboardingFlow.tsx @@ -14,6 +14,7 @@ import { WELCOME_SURFACE_READY_EVENT, } from "@/features/onboarding/welcome"; import { useAvatarPresentation } from "@/features/profile/avatarPresentationStore"; +import { registerAvatarWhenReady } from "@/features/profile/avatarProfileSync"; import { profileQueryKey } from "@/features/profile/hooks"; import { ProfileAvatar } from "@/features/profile/ui/ProfileAvatar"; import { @@ -152,6 +153,7 @@ export function CommunityOnboardingFlow({ const systemColorScheme = useSystemColorScheme(); const [displayName, setDisplayName] = React.useState(""); const [avatarUrl, setAvatarUrl] = React.useState(""); + const avatarPresentation = useAvatarPresentation(avatarUrl); const [isUploadingAvatar, setIsUploadingAvatar] = React.useState(false); const [isAvatarEditorOpen, setIsAvatarEditorOpen] = React.useState(false); const [starterPersonas, setStarterPersonas] = React.useState( @@ -380,10 +382,34 @@ export function CommunityOnboardingFlow({ if (!displayName.trim()) return; setIsPending(true); try { - await updateProfile({ - displayName: displayName.trim(), - avatarUrl: avatarUrl.trim() || undefined, - }); + const candidateAvatarUrl = avatarUrl.trim(); + const presentationState = avatarPresentation?.state; + const shouldSaveCandidate = + candidateAvatarUrl.length > 0 && + presentationState !== "failed" && + presentationState !== "pending"; + + const deferredAvatar = + candidateAvatarUrl && presentationState && presentationState !== "ready" + ? registerAvatarWhenReady({ + avatarUrl: candidateAvatarUrl, + relayUrl: transaction.relayUrl, + }) + : null; + + try { + const profile = await updateProfile({ + displayName: displayName.trim(), + avatarUrl: shouldSaveCandidate ? candidateAvatarUrl : undefined, + }); + deferredAvatar?.release({ + expectedPubkey: profile.pubkey, + expectedAvatarUrl: profile.avatarUrl, + }); + } catch (error) { + deferredAvatar?.cancel(); + throw error; + } update({ stage: "team-intro", error: undefined }); } catch (error) { if (isRelayMembershipDeniedError(error)) { diff --git a/desktop/src/features/profile/avatarPresentationStore.ts b/desktop/src/features/profile/avatarPresentationStore.ts index 769e7b27b..cfdad9724 100644 --- a/desktop/src/features/profile/avatarPresentationStore.ts +++ b/desktop/src/features/profile/avatarPresentationStore.ts @@ -125,7 +125,10 @@ async function verifyPresentation( export function beginAvatarPresentation(remoteUrl: string, image: Blob): void { const existing = presentations.get(remoteUrl); - if (existing) releaseLocalPreview(existing); + if (existing) { + toast.dismiss(toastId(remoteUrl)); + releaseLocalPreview(existing); + } const localPreviewUrl = URL.createObjectURL(image); const entry: AvatarPresentationEntry = { @@ -139,6 +142,27 @@ export function beginAvatarPresentation(remoteUrl: string, image: Blob): void { void verifyPresentation(entry); } +export function disposeAvatarPresentation(remoteUrl: string): void { + const entry = presentations.get(remoteUrl); + if (!entry) return; + + entry.generation = nextGeneration++; + toast.dismiss(toastId(remoteUrl)); + releaseLocalPreview(entry); + presentations.delete(remoteUrl); + emitChange(); +} + +export function resetAvatarPresentations(): void { + for (const entry of presentations.values()) { + entry.generation = nextGeneration++; + toast.dismiss(toastId(entry.remoteUrl)); + releaseLocalPreview(entry); + } + presentations.clear(); + emitChange(); +} + export function retryAvatarPresentation(remoteUrl: string): void { const entry = presentations.get(remoteUrl); if (entry?.snapshot.state !== "failed") return; @@ -172,3 +196,16 @@ export function useAvatarPresentation( () => null, ); } + +export function useAvatarSelection( + avatarUrl: string, + onUrlChange: (avatarUrl: string) => void, +): (avatarUrl: string) => void { + return React.useCallback( + (nextAvatarUrl: string) => { + if (avatarUrl !== nextAvatarUrl) disposeAvatarPresentation(avatarUrl); + onUrlChange(nextAvatarUrl); + }, + [avatarUrl, onUrlChange], + ); +} diff --git a/desktop/src/features/profile/avatarProfileSync.test.mjs b/desktop/src/features/profile/avatarProfileSync.test.mjs new file mode 100644 index 000000000..b814208bc --- /dev/null +++ b/desktop/src/features/profile/avatarProfileSync.test.mjs @@ -0,0 +1,303 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; + +import { + createAvatarProfileSync, + createProfileCacheRefreshQueue, +} from "./avatarProfileSync.ts"; + +const INPUT = { + avatarUrl: "https://old-relay.example/avatar.png", + relayUrl: "wss://old-relay.example", + expectedPubkey: "pubkey", + expectedAvatarUrl: null, +}; + +function createHarness({ + initialState = "pending", + saveProfile, + getActivePubkey = async () => INPUT.expectedPubkey, + scheduleRetry, +} = {}) { + let presentation = { displayUrl: INPUT.avatarUrl, state: initialState }; + let listener = () => {}; + let unsubscribeCount = 0; + const saves = []; + const refreshed = []; + const sync = createAvatarProfileSync({ + getPresentation: () => presentation, + subscribe: (nextListener) => { + listener = nextListener; + return () => { + unsubscribeCount += 1; + }; + }, + saveProfile: + saveProfile ?? + (async (input) => { + saves.push(input); + return { avatarUrl: input.avatarUrl, pubkey: input.expectedPubkey }; + }), + getActivePubkey, + refreshCaches: async (profile, input) => { + refreshed.push({ profile, input }); + }, + scheduleRetry, + }); + + return { + get unsubscribeCount() { + return unsubscribeCount; + }, + listener: () => listener(), + refreshed, + saves, + removePresentation: () => { + presentation = null; + }, + setState: (state) => { + presentation = { + ...(presentation ?? { displayUrl: INPUT.avatarUrl }), + state, + }; + }, + sync, + }; +} + +async function flushPromises() { + await new Promise((resolve) => setImmediate(resolve)); +} + +test("saves at the captured relay and refreshes caches after verification", async () => { + const harness = createHarness(); + harness.sync.saveWhenReady(INPUT); + + harness.setState("ready"); + harness.listener(); + await flushPromises(); + + assert.deepEqual(harness.saves, [INPUT]); + assert.equal(harness.refreshed.length, 1); + assert.equal(harness.refreshed[0].input.relayUrl, INPUT.relayUrl); + assert.equal(harness.unsubscribeCount, 1); +}); + +test("community reset cancels a pending avatar save", async () => { + const harness = createHarness(); + harness.sync.saveWhenReady(INPUT); + + harness.sync.reset(); + harness.setState("ready"); + harness.listener(); + await flushPromises(); + + assert.deepEqual(harness.saves, []); + assert.equal(harness.unsubscribeCount, 1); +}); + +test("a reset sync accepts deferred work from the next community", async () => { + const harness = createHarness(); + harness.sync.reset(); + harness.setState("ready"); + const nextInput = { + ...INPUT, + relayUrl: "wss://next-relay.example", + }; + + harness.sync.saveWhenReady(nextInput); + await flushPromises(); + + assert.deepEqual(harness.saves, [nextInput]); + assert.equal(harness.refreshed.length, 1); +}); + +test("skips cache refresh when the active identity changes during save", async () => { + let resolveSave; + const savePromise = new Promise((resolve) => { + resolveSave = resolve; + }); + let activePubkey = INPUT.expectedPubkey; + const harness = createHarness({ + initialState: "ready", + saveProfile: () => savePromise, + getActivePubkey: async () => activePubkey, + }); + harness.sync.saveWhenReady(INPUT); + + activePubkey = "replacement-pubkey"; + resolveSave({ avatarUrl: INPUT.avatarUrl, pubkey: INPUT.expectedPubkey }); + await flushPromises(); + + assert.deepEqual(harness.refreshed, []); + assert.equal(harness.unsubscribeCount, 1); +}); + +test("retries a transient save and keeps the sync pending until success", async () => { + const scheduled = []; + let attempt = 0; + const harness = createHarness({ + initialState: "ready", + saveProfile: async (input) => { + attempt += 1; + if (attempt === 1) throw new Error("relay unreachable: network error"); + return { avatarUrl: input.avatarUrl, pubkey: input.expectedPubkey }; + }, + scheduleRetry: (callback, delayMs) => { + const retry = { callback, delayMs, cancelled: false }; + scheduled.push(retry); + return () => { + retry.cancelled = true; + }; + }, + }); + + harness.sync.saveWhenReady(INPUT); + await flushPromises(); + + assert.equal(attempt, 1); + assert.equal(harness.unsubscribeCount, 0); + assert.equal(scheduled[0].delayMs, 5_000); + + harness.listener(); + await flushPromises(); + assert.equal(attempt, 1, "store updates must not bypass retry backoff"); + + scheduled[0].callback(); + await flushPromises(); + + assert.equal(attempt, 2); + assert.equal(harness.refreshed.length, 1); + assert.equal(harness.unsubscribeCount, 1); +}); + +test("reset cancels a scheduled transient retry", async () => { + let scheduled; + let attempts = 0; + const harness = createHarness({ + initialState: "ready", + saveProfile: async () => { + attempts += 1; + throw new Error("relay rate-limited: retry in 5s"); + }, + scheduleRetry: (callback, delayMs) => { + scheduled = { callback, delayMs, cancelled: false }; + return () => { + scheduled.cancelled = true; + }; + }, + }); + + harness.sync.saveWhenReady(INPUT); + await flushPromises(); + harness.sync.reset(); + + assert.equal(scheduled.delayMs, 5_000); + assert.equal(scheduled.cancelled, true); + scheduled.callback(); + await flushPromises(); + assert.equal(attempts, 1); +}); + +test("registration preserves readiness until the initial profile write completes", async () => { + const harness = createHarness(); + const registration = harness.sync.registerWhenReady({ + avatarUrl: INPUT.avatarUrl, + relayUrl: INPUT.relayUrl, + }); + + harness.setState("ready"); + harness.listener(); + harness.removePresentation(); + assert.deepEqual(harness.saves, []); + + registration.release({ + expectedPubkey: INPUT.expectedPubkey, + expectedAvatarUrl: null, + }); + await flushPromises(); + + assert.deepEqual(harness.saves, [INPUT]); + assert.equal(harness.refreshed.length, 1); +}); + +test("cancelled registration cannot save after the initial profile write fails", async () => { + const harness = createHarness(); + const registration = harness.sync.registerWhenReady({ + avatarUrl: INPUT.avatarUrl, + relayUrl: INPUT.relayUrl, + }); + + harness.setState("ready"); + harness.listener(); + registration.cancel(); + registration.release({ + expectedPubkey: INPUT.expectedPubkey, + expectedAvatarUrl: null, + }); + await flushPromises(); + + assert.deepEqual(harness.saves, []); + assert.equal(harness.unsubscribeCount, 1); +}); + +test("cache refresh waits for a provider and flushes exactly once", async () => { + const refreshed = []; + const queue = createProfileCacheRefreshQueue( + async (client, profile, relayUrl) => { + refreshed.push({ client, profile, relayUrl }); + }, + ); + const profile = { avatarUrl: INPUT.avatarUrl, pubkey: INPUT.expectedPubkey }; + const client = {}; + + await queue.enqueue({ profile, input: INPUT }); + assert.deepEqual(refreshed, []); + + const detach = queue.setClient(client); + await flushPromises(); + assert.deepEqual(refreshed, [{ client, profile, relayUrl: INPUT.relayUrl }]); + + detach(); + const detachAgain = queue.setClient(client); + await flushPromises(); + assert.equal(refreshed.length, 1); + detachAgain(); +}); + +test("cache refresh reset discards work from the previous community", async () => { + const refreshed = []; + const queue = createProfileCacheRefreshQueue( + async (client, profile, relayUrl) => { + refreshed.push({ client, profile, relayUrl }); + }, + ); + + await queue.enqueue({ + profile: { avatarUrl: INPUT.avatarUrl, pubkey: INPUT.expectedPubkey }, + input: INPUT, + }); + queue.reset(); + queue.setClient({}); + await flushPromises(); + + assert.deepEqual(refreshed, []); +}); + +test("cache refresh follows only a successful save", async () => { + let rejectSave; + const savePromise = new Promise((_, reject) => { + rejectSave = reject; + }); + const harness = createHarness({ + initialState: "ready", + saveProfile: () => savePromise, + }); + harness.sync.saveWhenReady(INPUT); + + rejectSave(new Error("stale baseline")); + await flushPromises(); + + assert.deepEqual(harness.refreshed, []); + assert.equal(harness.unsubscribeCount, 1); +}); diff --git a/desktop/src/features/profile/avatarProfileSync.ts b/desktop/src/features/profile/avatarProfileSync.ts new file mode 100644 index 000000000..5bca6d8da --- /dev/null +++ b/desktop/src/features/profile/avatarProfileSync.ts @@ -0,0 +1,262 @@ +import type { QueryClient } from "@tanstack/react-query"; + +import { + getAvatarPresentation, + subscribeAvatarPresentations, + type AvatarPresentation, +} from "@/features/profile/avatarPresentationStore"; +import { refreshProfileCaches } from "@/features/profile/profileCacheSync"; +import { getIdentity } from "@/shared/api/tauriIdentity"; +import { updateProfileAtRelay } from "@/shared/api/tauriProfiles"; +import type { Profile } from "@/shared/api/types"; +import { isRelayUnreachableError } from "@/shared/lib/relayError"; + +const AVATAR_SAVE_RETRY_DELAYS_MS = [5_000, 30_000, 120_000] as const; + +type PendingAvatarSave = { + avatarUrl: string; + relayUrl: string; + expectedPubkey: string; + expectedAvatarUrl: string | null; +}; + +type DeferredAvatarSave = Pick; + +type AvatarSaveRegistration = { + cancel: () => void; + release: ( + input: Pick, + ) => void; +}; + +type AvatarProfileSyncDependencies = { + getPresentation: (avatarUrl: string) => AvatarPresentation | null; + subscribe: (listener: () => void) => () => void; + saveProfile: (input: PendingAvatarSave) => Promise; + getActivePubkey: () => Promise; + refreshCaches: (profile: Profile, input: PendingAvatarSave) => Promise; + scheduleRetry?: (callback: () => void, delayMs: number) => () => void; +}; + +function isRetryableAvatarSaveError(error: unknown): boolean { + const message = error instanceof Error ? error.message : String(error); + return ( + isRelayUnreachableError(error) || message.startsWith("relay rate-limited:") + ); +} + +export function createAvatarProfileSync( + dependencies: AvatarProfileSyncDependencies, +) { + const pendingSyncs = new Map void>(); + let generation = 0; + + const reset = () => { + generation += 1; + for (const stop of pendingSyncs.values()) stop(); + pendingSyncs.clear(); + }; + + const queueSave = (input: PendingAvatarSave, assumeReady = false): void => { + const syncKey = `${input.relayUrl}:${input.expectedPubkey}:${input.avatarUrl}`; + if (pendingSyncs.has(syncKey)) return; + + let isSaving = false; + let isReady = assumeReady; + let retryAttempt = 0; + let cancelRetry: (() => void) | null = null; + let unsubscribe = () => {}; + const queuedGeneration = generation; + const stop = () => { + cancelRetry?.(); + cancelRetry = null; + unsubscribe(); + pendingSyncs.delete(syncKey); + }; + const saveIfReady = () => { + if (generation !== queuedGeneration || cancelRetry !== null) return; + const presentation = dependencies.getPresentation(input.avatarUrl); + if (presentation?.state === "ready") isReady = true; + if (!presentation && !isReady) { + stop(); + return; + } + if (!isReady || isSaving) return; + + isSaving = true; + void dependencies + .saveProfile(input) + .then(async (profile) => { + if (generation !== queuedGeneration) return; + const activePubkey = await dependencies.getActivePubkey(); + if ( + generation !== queuedGeneration || + activePubkey?.toLowerCase() !== input.expectedPubkey.toLowerCase() + ) { + return; + } + await dependencies.refreshCaches(profile, input); + }) + .then(stop) + .catch((error: unknown) => { + if ( + generation !== queuedGeneration || + !isRetryableAvatarSaveError(error) + ) { + stop(); + return; + } + const delayMs = AVATAR_SAVE_RETRY_DELAYS_MS[retryAttempt]; + if (delayMs === undefined) { + stop(); + return; + } + retryAttempt += 1; + isSaving = false; + const scheduleRetry = + dependencies.scheduleRetry ?? + ((callback, delay) => { + const timeout = window.setTimeout(callback, delay); + return () => window.clearTimeout(timeout); + }); + cancelRetry = scheduleRetry(() => { + cancelRetry = null; + saveIfReady(); + }, delayMs); + }); + }; + + unsubscribe = dependencies.subscribe(saveIfReady); + pendingSyncs.set(syncKey, stop); + saveIfReady(); + }; + + const registerWhenReady = ( + input: DeferredAvatarSave, + ): AvatarSaveRegistration => { + const registrationKey = `registration:${input.relayUrl}:${input.avatarUrl}`; + if (pendingSyncs.has(registrationKey)) { + return { cancel: () => {}, release: () => {} }; + } + + let observedReady = false; + let active = true; + const queuedGeneration = generation; + const observe = () => { + if (generation !== queuedGeneration) return; + if (dependencies.getPresentation(input.avatarUrl)?.state === "ready") { + observedReady = true; + } + }; + const unsubscribe = dependencies.subscribe(observe); + const cancel = () => { + if (!active) return; + active = false; + unsubscribe(); + pendingSyncs.delete(registrationKey); + }; + pendingSyncs.set(registrationKey, cancel); + observe(); + + return { + cancel, + release: (completion) => { + if (!active || generation !== queuedGeneration) return; + cancel(); + queueSave({ ...input, ...completion }, observedReady); + }, + }; + }; + + return { registerWhenReady, reset, saveWhenReady: queueSave }; +} + +type ProfileCacheRefresh = { + profile: Profile; + input: PendingAvatarSave; +}; + +type ProfileCacheRefreshQueue = { + enqueue: (refresh: ProfileCacheRefresh) => Promise; + reset: () => void; + setClient: (client: QueryClient) => () => void; +}; + +export function createProfileCacheRefreshQueue( + refresh: ( + client: QueryClient, + profile: Profile, + relayUrl: string, + ) => Promise, +): ProfileCacheRefreshQueue { + let client: QueryClient | null = null; + const pending = new Map(); + const refreshKey = ({ profile, input }: ProfileCacheRefresh) => + `${input.relayUrl}:${profile.pubkey.toLowerCase()}`; + + const flush = (nextClient: QueryClient) => { + const queued = [...pending.values()]; + pending.clear(); + for (const item of queued) { + void refresh(nextClient, item.profile, item.input.relayUrl); + } + }; + + return { + enqueue: async (item) => { + if (client) { + await refresh(client, item.profile, item.input.relayUrl); + return; + } + pending.set(refreshKey(item), item); + }, + reset: () => pending.clear(), + setClient: (nextClient) => { + client = nextClient; + flush(nextClient); + return () => { + if (client === nextClient) client = null; + }; + }, + }; +} + +const profileCacheRefreshQueue = + createProfileCacheRefreshQueue(refreshProfileCaches); + +const avatarProfileSync = createAvatarProfileSync({ + getPresentation: getAvatarPresentation, + subscribe: subscribeAvatarPresentations, + saveProfile: updateProfileAtRelay, + getActivePubkey: async () => { + try { + return (await getIdentity()).pubkey; + } catch { + return null; + } + }, + refreshCaches: async (profile, input) => { + await profileCacheRefreshQueue.enqueue({ profile, input }); + }, +}); + +export function setAvatarProfileSyncQueryClient( + client: QueryClient, +): () => void { + return profileCacheRefreshQueue.setClient(client); +} + +export function registerAvatarWhenReady( + input: DeferredAvatarSave, +): AvatarSaveRegistration { + return avatarProfileSync.registerWhenReady(input); +} + +export function saveAvatarWhenReady(input: PendingAvatarSave): void { + avatarProfileSync.saveWhenReady(input); +} + +export function resetAvatarProfileSync(): void { + profileCacheRefreshQueue.reset(); + avatarProfileSync.reset(); +} diff --git a/desktop/src/features/profile/profileCacheSync.test.mjs b/desktop/src/features/profile/profileCacheSync.test.mjs new file mode 100644 index 000000000..6c423abad --- /dev/null +++ b/desktop/src/features/profile/profileCacheSync.test.mjs @@ -0,0 +1,91 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { QueryClient } from "@tanstack/react-query"; + +import { readSelfProfileCache } from "./lib/selfProfileStorage.ts"; +import { refreshProfileCaches } from "./profileCacheSync.ts"; + +function installBrowserStubs() { + const values = new Map(); + globalThis.window = { + dispatchEvent() {}, + localStorage: { + getItem: (key) => values.get(key) ?? null, + setItem: (key, value) => values.set(key, value), + removeItem: (key) => values.delete(key), + key: (index) => [...values.keys()][index] ?? null, + get length() { + return values.size; + }, + }, + }; + globalThis.CustomEvent = class CustomEvent {}; + globalThis.fetch = async () => ({ ok: false }); +} + +const PUBKEY = "abcdef"; +const RELAY_URL = "wss://relay.example"; +const PROFILE = { + pubkey: PUBKEY, + displayName: "Alice", + avatarUrl: "https://cdn.example/avatar.png", + about: "About Alice", + nip05Handle: null, + ownerPubkey: null, + hasProfileEvent: true, +}; + +test("successful deferred save synchronizes every profile cache", async () => { + installBrowserStubs(); + const queryClient = new QueryClient(); + queryClient.setQueryData(["profile"], { ...PROFILE, avatarUrl: null }); + queryClient.setQueryData(["user-profile", PUBKEY], { + ...PROFILE, + avatarUrl: null, + }); + queryClient.setQueryData(["users-batch-entry", PUBKEY], { + summary: { displayName: "Alice", avatarUrl: null }, + fetchedAt: Date.now(), + }); + queryClient.setQueryData(["users-batch", PUBKEY], { + profiles: { + [PUBKEY]: { + displayName: "Alice", + avatarUrl: null, + nip05Handle: null, + ownerPubkey: null, + }, + }, + missing: [], + }); + queryClient.setQueryData( + ["user-search", "alice", 8], + [{ pubkey: PUBKEY, displayName: "Alice", avatarUrl: null }], + ); + + await refreshProfileCaches(queryClient, PROFILE, RELAY_URL); + + assert.deepEqual(queryClient.getQueryData(["profile"]), PROFILE); + assert.deepEqual(queryClient.getQueryData(["user-profile", PUBKEY]), PROFILE); + assert.equal( + queryClient.getQueryData(["users-batch", PUBKEY]).profiles[PUBKEY] + .avatarUrl, + PROFILE.avatarUrl, + ); + assert.equal( + queryClient.getQueryData(["users-batch-entry", PUBKEY]), + undefined, + ); + assert.equal( + queryClient.getQueryState(["user-search", "alice", 8]).isInvalidated, + true, + ); + const persisted = readSelfProfileCache(RELAY_URL, PUBKEY); + assert.equal(persisted.displayName, PROFILE.displayName); + assert.equal(persisted.avatarUrl, PROFILE.avatarUrl); + assert.equal(persisted.about, PROFILE.about); + assert.equal(persisted.avatarDataUrl, null); + assert.equal(persisted.hasProfileEvent, true); + assert.ok(persisted.updatedAt > 0); +}); diff --git a/desktop/src/features/profile/profileCacheSync.ts b/desktop/src/features/profile/profileCacheSync.ts new file mode 100644 index 000000000..656096c97 --- /dev/null +++ b/desktop/src/features/profile/profileCacheSync.ts @@ -0,0 +1,91 @@ +import type { Query, QueryClient } from "@tanstack/react-query"; + +import { + evictUsersBatchEntries, + profileQueryKey, +} from "@/features/profile/hooks"; +import { + fetchAvatarDataUrl, + readSelfProfileCache, + resolveAvatarDataUrl, + writeSelfProfileCache, +} from "@/features/profile/lib/selfProfileStorage"; +import type { + Profile, + UserProfileSummary, + UsersBatchResponse, +} from "@/shared/api/types"; +import { getAvatarSnapshotUrl } from "@/shared/lib/animatedAvatar"; +import { rewriteRelayUrl } from "@/shared/lib/mediaUrl"; + +function queryContainsPubkey(query: Query, pubkey: string): boolean { + return query.queryKey.includes(pubkey); +} + +export async function refreshProfileCaches( + queryClient: QueryClient, + profile: Profile, + relayUrl: string, +): Promise { + const pubkey = profile.pubkey.toLowerCase(); + await queryClient.cancelQueries({ + predicate: (query) => + query.queryKey[0] === profileQueryKey[0] || + (query.queryKey[0] === "user-profile" && + queryContainsPubkey(query, pubkey)) || + (query.queryKey[0] === "users-batch" && + queryContainsPubkey(query, pubkey)), + }); + + queryClient.setQueryData(profileQueryKey, profile); + queryClient.setQueryData(["user-profile", pubkey], profile); + evictUsersBatchEntries(queryClient, [pubkey]); + queryClient.setQueriesData( + { + predicate: (query) => + query.queryKey[0] === "users-batch" && + queryContainsPubkey(query, pubkey), + }, + (current) => { + if (!current?.profiles[pubkey]) return current; + return { + ...current, + profiles: { + ...current.profiles, + [pubkey]: { + ...current.profiles[pubkey], + avatarUrl: profile.avatarUrl, + } satisfies UserProfileSummary, + }, + }; + }, + ); + // Search result pages also embed profile avatars, but their arbitrary query + // text/page shape makes a safe targeted rewrite brittle. Mark every search + // view stale; active searches refetch immediately and inactive ones refresh + // when next opened. + await queryClient.invalidateQueries({ queryKey: ["user-search"] }); + + const existing = readSelfProfileCache(relayUrl, profile.pubkey); + const baseCache = { + version: 1 as const, + displayName: profile.displayName, + avatarUrl: profile.avatarUrl, + about: profile.about, + avatarDataUrl: resolveAvatarDataUrl(profile.avatarUrl, null, existing), + updatedAt: Date.now(), + ...(profile.hasProfileEvent && { hasProfileEvent: true as const }), + }; + // Persist the canonical profile before attempting the optional image snapshot, + // so quitting during that fetch cannot leave the durable fallback stale. + writeSelfProfileCache(relayUrl, profile.pubkey, baseCache); + + const snapshotUrl = getAvatarSnapshotUrl(profile.avatarUrl); + if (!snapshotUrl) return; + const fetched = await fetchAvatarDataUrl(rewriteRelayUrl(snapshotUrl)); + if (!fetched) return; + writeSelfProfileCache(relayUrl, profile.pubkey, { + ...baseCache, + avatarDataUrl: fetched, + }); +} diff --git a/desktop/src/features/profile/ui/ProfileAvatarEditor.tsx b/desktop/src/features/profile/ui/ProfileAvatarEditor.tsx index 9e7901fff..d66385048 100644 --- a/desktop/src/features/profile/ui/ProfileAvatarEditor.tsx +++ b/desktop/src/features/profile/ui/ProfileAvatarEditor.tsx @@ -9,6 +9,7 @@ import { AnimatedAvatarCapture } from "@/features/profile/ui/AnimatedAvatarCaptu import { AvatarCustomColorPanel } from "@/features/profile/ui/AvatarCustomColorPanel"; import { ProfileAvatarUploadPreview } from "@/features/profile/ui/ProfileAvatarUploadPreview"; import { ProfileAvatarModeTabs } from "@/features/profile/ui/ProfileAvatarModeTabs"; +import { useAvatarSelection } from "@/features/profile/avatarPresentationStore"; import { useAvatarUpload } from "@/features/profile/useAvatarUpload"; import { cn } from "@/shared/lib/cn"; import { Button } from "@/shared/ui/button"; @@ -161,14 +162,15 @@ export function ProfileAvatarEditor({ }, [mode, onModeChange], ); + const setAvatar = useAvatarSelection(avatarUrl, onUrlChange); const handleUploadSuccess = React.useCallback( (uploadedUrl: string) => { setUrlDraft(""); onUploadedAvatarChange?.(uploadedUrl); - onUrlChange(uploadedUrl); + setAvatar(uploadedUrl); updateMode("image"); }, - [onUploadedAvatarChange, onUrlChange, updateMode], + [onUploadedAvatarChange, setAvatar, updateMode], ); const [isAnimatedApplyPending, setIsAnimatedApplyPending] = React.useState(false); @@ -192,14 +194,14 @@ export function ProfileAvatarEditor({ clearUploadError(); setUrlDraft(""); onUploadedAvatarChange?.(animatedUrl); - onUrlChange(animatedUrl); + setAvatar(animatedUrl); onAnimatedAvatarApply?.(animatedUrl); }, [ clearUploadError, onAnimatedAvatarApply, onUploadedAvatarChange, - onUrlChange, + setAvatar, ], ); // Done on the animated tab uploads the pending recording first, then @@ -347,14 +349,14 @@ export function ProfileAvatarEditor({ } onUploadedAvatarChange?.(null); - onUrlChange(nextAvatarUrl); + setAvatar(nextAvatarUrl); }, [ avatarUrl, customColorDraft, isCustomColorPickerOpen, onUploadedAvatarChange, - onUrlChange, selectedEmoji, + setAvatar, ]); const handleFiles = React.useCallback( @@ -379,14 +381,14 @@ export function ProfileAvatarEditor({ clearUploadError(); onUploadedAvatarChange?.(null); - onUrlChange(nextUrl); + setAvatar(nextUrl); hasUserEditedUrlDraftRef.current = false; updateMode("image"); }, [ clearUploadError, isInputDisabled, onUploadedAvatarChange, - onUrlChange, + setAvatar, updateMode, urlDraft, ]); @@ -396,10 +398,10 @@ export function ProfileAvatarEditor({ setUrlDraft(""); hasUserEditedUrlDraftRef.current = false; onUploadedAvatarChange?.(null); - onUrlChange(emojiAvatarDataUrl(emoji, color)); + setAvatar(emojiAvatarDataUrl(emoji, color)); onEmojiAvatarChange?.(); }, - [onEmojiAvatarChange, onUploadedAvatarChange, onUrlChange, selectedColor], + [onEmojiAvatarChange, onUploadedAvatarChange, selectedColor, setAvatar], ); const openCustomColorPicker = React.useCallback(() => { @@ -699,7 +701,7 @@ export function ProfileAvatarEditor({ hasUserEditedUrlDraftRef.current = true; setUrlDraft(event.target.value); onUploadedAvatarChange?.(null); - onUrlChange(event.target.value); + setAvatar(event.target.value); }} onFocus={() => { isUrlInputFocusedRef.current = true; diff --git a/desktop/src/shared/api/tauriProfiles.ts b/desktop/src/shared/api/tauriProfiles.ts index aab1fd7a1..c8e52f516 100644 --- a/desktop/src/shared/api/tauriProfiles.ts +++ b/desktop/src/shared/api/tauriProfiles.ts @@ -83,6 +83,19 @@ export async function updateProfile( return fromRawProfile(profile); } +export async function updateProfileAtRelay(input: { + relayUrl: string; + expectedPubkey: string; + expectedAvatarUrl: string | null; + avatarUrl: string; +}): Promise { + const profile = await invokeTauri( + "update_profile_at_relay", + input, + ); + return fromRawProfile(profile); +} + export async function getUserProfile(pubkey?: string): Promise { const profile = await invokeTauri("get_user_profile", { pubkey }); return fromRawProfile(profile); diff --git a/desktop/src/testing/e2eBridge.ts b/desktop/src/testing/e2eBridge.ts index 4ef5ee56e..af7fb4bc2 100644 --- a/desktop/src/testing/e2eBridge.ts +++ b/desktop/src/testing/e2eBridge.ts @@ -9326,6 +9326,13 @@ export function maybeInstallE2eTauriMocks() { payload as Parameters[0], activeConfig, ); + case "update_profile_at_relay": + return handleUpdateProfile( + { + avatarUrl: (payload as { avatarUrl: string }).avatarUrl, + }, + activeConfig, + ); case "get_user_profile": return handleGetUserProfile( (payload as Parameters[0]) ?? {}, diff --git a/desktop/tests/e2e/onboarding.spec.ts b/desktop/tests/e2e/onboarding.spec.ts index c3c716bec..33e441c0f 100644 --- a/desktop/tests/e2e/onboarding.spec.ts +++ b/desktop/tests/e2e/onboarding.spec.ts @@ -58,6 +58,8 @@ async function setRelayConnectionState( const HOME_SEEN_STORAGE_KEY_PREFIX = "buzz-home-feed-seen.v1:"; const COMMUNITY_ONBOARDING_TRANSACTION_STORAGE_KEY = "buzz-community-onboarding-transaction.v1"; +const ONE_PIXEL_PNG_BASE64 = + "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII="; const DEFAULT_MOCK_PUBKEY = "deadbeef".repeat(8); const BLANK_TYLER_IDENTITY = { ...TEST_IDENTITIES.tyler, @@ -89,6 +91,48 @@ async function seedOnboardingCompletion(page: Page, pubkey: string) { ); } +async function seedCommunityProfileStage(page: Page, id: string) { + await seedActiveIdentity(page, BLANK_TYLER_IDENTITY); + await page.addInitScript( + ({ pubkey, transactionId, transactionStorageKey }) => { + window.localStorage.setItem( + `buzz-machine-onboarding-complete.v2:${pubkey}`, + "true", + ); + const timestamp = new Date().toISOString(); + window.localStorage.setItem( + transactionStorageKey, + JSON.stringify({ + id: transactionId, + source: "first-community", + stage: "profile", + relayUrl: "wss://default.example.com", + communityName: "Default", + communityId: "e2e-default-community", + addedCommunity: true, + createdAt: timestamp, + updatedAt: timestamp, + }), + ); + }, + { + pubkey: BLANK_TYLER_IDENTITY.pubkey, + transactionId: id, + transactionStorageKey: COMMUNITY_ONBOARDING_TRANSACTION_STORAGE_KEY, + }, + ); +} + +async function uploadCommunityAvatar(page: Page, filename: string) { + await page.getByTestId("community-avatar-open").click(); + await page.getByTestId("community-avatar-input").setInputFiles({ + buffer: Buffer.from(ONE_PIXEL_PNG_BASE64, "base64"), + mimeType: "image/png", + name: filename, + }); + await page.getByTestId("community-avatar-done").click(); +} + async function readHomeSeenStorageKeys(page: Page) { return page.evaluate((prefix) => { return Object.keys(window.localStorage).filter((key) => @@ -412,6 +456,23 @@ async function invokeMockCommand( ); } +async function seedCurrentAvatar(page: Page, avatarUrl: string) { + await page.waitForFunction(() => { + const bridgeWindow = window as Window & { + __BUZZ_E2E_INVOKE_MOCK_COMMAND__?: unknown; + __TAURI_INTERNALS__?: { invoke?: unknown }; + }; + return ( + typeof bridgeWindow.__BUZZ_E2E_INVOKE_MOCK_COMMAND__ === "function" || + typeof bridgeWindow.__TAURI_INTERNALS__?.invoke === "function" + ); + }); + await invokeMockCommand(page, "update_profile", { avatarUrl }); + await page.evaluate(() => { + window.__BUZZ_E2E_COMMAND_PAYLOADS__ = []; + }); +} + async function getWelcomeChannelId(page: Page) { const channels = await getMockChannels(page); return ( @@ -1636,37 +1697,46 @@ test("connected first-community profile step offers equal-width Next and Back co .toBeNull(); }); -test("pending avatar stays navigable and exposes retry after propagation fails", async ({ +test("name-only community profile save preserves an existing avatar", async ({ page, }) => { - await seedActiveIdentity(page, BLANK_TYLER_IDENTITY); - await page.addInitScript( - ({ pubkey, transactionStorageKey }) => { - window.localStorage.setItem( - `buzz-machine-onboarding-complete.v2:${pubkey}`, - "true", - ); - const timestamp = new Date().toISOString(); - window.localStorage.setItem( - transactionStorageKey, - JSON.stringify({ - id: "txn-avatar-propagation", - source: "first-community", - stage: "profile", - relayUrl: "wss://default.example.com", - communityName: "Default", - communityId: "e2e-default-community", - addedCommunity: true, - createdAt: timestamp, - updatedAt: timestamp, - }), - ); - }, - { - pubkey: BLANK_TYLER_IDENTITY.pubkey, - transactionStorageKey: COMMUNITY_ONBOARDING_TRANSACTION_STORAGE_KEY, - }, + await seedCommunityProfileStage(page, "txn-avatar-preserve-existing"); + await installMockBridge(page, undefined, { + relayWsUrl: "wss://default.example.com", + skipOnboardingSeed: true, + }); + await page.goto("/"); + + const existingAvatarUrl = + "https://mock.relay/media/existing-community-avatar.png"; + await seedCurrentAvatar(page, existingAvatarUrl); + await page.getByTestId("community-profile-name-key").fill("Tyler"); + await page.getByTestId("community-profile-next").click(); + + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined]); + const profile = await invokeMockCommand<{ avatar_url: string | null }>( + page, + "get_profile", ); + expect(profile.avatar_url).toBe(existingAvatarUrl); +}); + +test("pending avatar stays navigable, clears failures, and retries", async ({ + page, +}) => { + await seedCommunityProfileStage(page, "txn-avatar-propagation"); const uploadedAvatarUrl = "https://mock.relay/media/pending-community-avatar.png"; @@ -1678,10 +1748,7 @@ test("pending avatar stays navigable and exposes retry after propagation fails", } await new Promise((resolve) => setTimeout(resolve, 500)); await route.fulfill({ - body: Buffer.from( - "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII=", - "base64", - ), + body: Buffer.from(ONE_PIXEL_PNG_BASE64, "base64"), contentType: "image/png", }); }); @@ -1708,16 +1775,7 @@ test("pending avatar stays navigable and exposes retry after propagation fails", await page.goto("/"); await page.getByTestId("community-profile-name-key").fill("Tyler"); - await page.getByTestId("community-avatar-open").click(); - await page.getByTestId("community-avatar-input").setInputFiles({ - buffer: Buffer.from( - "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII=", - "base64", - ), - mimeType: "image/png", - name: "pending-community-avatar.png", - }); - await page.getByTestId("community-avatar-done").click(); + await uploadCommunityAvatar(page, "pending-community-avatar.png"); const avatarImage = page.getByTestId("community-avatar-circle-image"); await expect(avatarImage).toHaveAttribute("src", /^blob:/); @@ -1773,16 +1831,316 @@ test("pending avatar stays navigable and exposes retry after propagation fails", page.getByText("Your default avatar is showing instead."), ).toBeVisible(); + await page.getByTestId("community-profile-next").click(); + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined]); + avatarReady = true; await page.getByRole("button", { name: "Retry" }).click(); + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined, uploadedAvatarUrl]); + await expect(page.getByText("Avatar couldn’t finish uploading")).toHaveCount( + 0, + ); +}); + +test("a pending avatar never becomes durable if propagation fails after onboarding unmounts", async ({ + page, +}) => { + await seedCommunityProfileStage(page, "txn-avatar-saved-before-failure"); + const uploadedAvatarUrl = + "https://mock.relay/media/saved-pending-community-avatar.png"; + let allowAvatarFailure = false; + await page.route(`${uploadedAvatarUrl}*`, async (route) => { + while (!allowAvatarFailure) { + await new Promise((resolve) => setTimeout(resolve, 50)); + } + await route.fulfill({ status: 404 }); + }); + await installMockBridge( + page, + { + uploadDescriptors: [ + { + filename: "saved-pending-community-avatar.png", + sha256: "f".repeat(64), + size: 128, + type: "image/png", + uploaded: 1_779_900_002, + url: uploadedAvatarUrl, + }, + ], + }, + { + relayWsUrl: "wss://default.example.com", + skipOnboardingSeed: true, + }, + ); + await page.goto("/"); + + await page.getByTestId("community-profile-name-key").fill("Tyler"); + await uploadCommunityAvatar(page, "saved-pending-community-avatar.png"); await expect( page.getByTestId("community-avatar-circle-upload-pending"), ).toBeVisible(); - await expect(avatarImage).toHaveAttribute("src", /^blob:/); - await expect(avatarImage).toHaveClass(/brightness-75/); + await page.getByTestId("community-profile-next").click(); + await page.getByTestId("community-team-intro-enter").click(); + await expect(page.getByTestId("community-onboarding-flow")).toHaveCount(0, { + timeout: 10_000, + }); + allowAvatarFailure = true; + await expect - .poll(() => avatarImage.getAttribute("src")) - .not.toMatch(/^blob:/); + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined]); + await expect( + page.getByText("Avatar couldn’t finish uploading"), + ).toBeVisible(); + const profile = await invokeMockCommand<{ avatar_url: string | null }>( + page, + "get_profile", + ); + expect(profile.avatar_url).toBeNull(); +}); + +test("a pending avatar becomes durable after onboarding unmounts once ready", async ({ + page, +}) => { + await seedCommunityProfileStage(page, "txn-avatar-ready-after-unmount"); + const uploadedAvatarUrl = + "https://mock.relay/media/ready-after-unmount-community-avatar.png"; + let allowAvatarReady = false; + await page.route(`${uploadedAvatarUrl}*`, async (route) => { + while (!allowAvatarReady) { + await new Promise((resolve) => setTimeout(resolve, 50)); + } + await route.fulfill({ + body: Buffer.from(ONE_PIXEL_PNG_BASE64, "base64"), + contentType: "image/png", + }); + }); + await installMockBridge( + page, + { + uploadDescriptors: [ + { + filename: "ready-after-unmount-community-avatar.png", + sha256: "b".repeat(64), + size: 128, + type: "image/png", + uploaded: 1_779_900_004, + url: uploadedAvatarUrl, + }, + ], + }, + { + relayWsUrl: "wss://default.example.com", + skipOnboardingSeed: true, + }, + ); + await page.goto("/"); + + await page.getByTestId("community-profile-name-key").fill("Tyler"); + await uploadCommunityAvatar(page, "ready-after-unmount-community-avatar.png"); + await page.getByTestId("community-profile-next").click(); + await page.getByTestId("community-team-intro-enter").click(); + await expect(page.getByTestId("community-onboarding-flow")).toHaveCount(0, { + timeout: 10_000, + }); + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined]); + + allowAvatarReady = true; + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined, uploadedAvatarUrl]); + const profile = await invokeMockCommand<{ avatar_url: string | null }>( + page, + "get_profile", + ); + expect(profile.avatar_url).toBe(uploadedAvatarUrl); +}); + +test("a failed pending replacement leaves the confirmed avatar untouched", async ({ + page, +}) => { + await seedCommunityProfileStage(page, "txn-avatar-restore-existing"); + const existingAvatarUrl = + "https://mock.relay/media/existing-community-avatar.png"; + const uploadedAvatarUrl = + "https://mock.relay/media/replacement-community-avatar.png"; + await page.route(`${uploadedAvatarUrl}*`, (route) => + route.fulfill({ status: 404 }), + ); + await installMockBridge( + page, + { + uploadDescriptors: [ + { + filename: "replacement-community-avatar.png", + sha256: "a".repeat(64), + size: 128, + type: "image/png", + uploaded: 1_779_900_003, + url: uploadedAvatarUrl, + }, + ], + }, + { + relayWsUrl: "wss://default.example.com", + skipOnboardingSeed: true, + }, + ); + await page.goto("/"); + await seedCurrentAvatar(page, existingAvatarUrl); + + await page.getByTestId("community-profile-name-key").fill("Tyler"); + await uploadCommunityAvatar(page, "replacement-community-avatar.png"); + await page.getByTestId("community-profile-next").click(); + + await expect + .poll(() => + page.evaluate(() => + (window.__BUZZ_E2E_COMMAND_PAYLOADS__ ?? []) + .filter( + ({ command }) => + command === "update_profile" || + command === "update_profile_at_relay", + ) + .map(({ payload }) => (payload as { avatarUrl?: string }).avatarUrl), + ), + ) + .toEqual([undefined]); + const profile = await invokeMockCommand<{ avatar_url: string | null }>( + page, + "get_profile", + ); + expect(profile.avatar_url).toBe(existingAvatarUrl); +}); + +test("replacing a pending upload disposes its verifier and local preview", async ({ + page, +}) => { + await seedCommunityProfileStage(page, "txn-avatar-replacement"); + await page.addInitScript(() => { + const testWindow = window as Window & { + __BUZZ_E2E_REVOKED_OBJECT_URLS__?: string[]; + }; + const revokedUrls: string[] = []; + const revokeObjectUrl = URL.revokeObjectURL.bind(URL); + testWindow.__BUZZ_E2E_REVOKED_OBJECT_URLS__ = revokedUrls; + URL.revokeObjectURL = (url) => { + revokedUrls.push(url); + revokeObjectUrl(url); + }; + }); + + const uploadedAvatarUrl = + "https://mock.relay/media/superseded-community-avatar.png"; + await page.route(`${uploadedAvatarUrl}*`, (route) => + route.fulfill({ status: 404 }), + ); + await installMockBridge( + page, + { + uploadDescriptors: [ + { + filename: "superseded-community-avatar.png", + sha256: "e".repeat(64), + size: 128, + type: "image/png", + uploaded: 1_779_900_001, + url: uploadedAvatarUrl, + }, + ], + }, + { + relayWsUrl: "wss://default.example.com", + skipOnboardingSeed: true, + }, + ); + await page.goto("/"); + + await page.getByTestId("community-profile-name-key").fill("Tyler"); + await uploadCommunityAvatar(page, "superseded-community-avatar.png"); + const supersededPreviewUrl = await page + .getByTestId("community-avatar-circle-image") + .getAttribute("src"); + expect(supersededPreviewUrl).toMatch(/^blob:/); + + await page.getByTestId("community-avatar-open").click(); + await page.getByRole("tab", { name: "Emoji" }).click(); + await selectFirstEmojiFromPicker(page); + await page.getByTestId("community-avatar-done").click(); + + await expect + .poll(() => + page.evaluate(() => { + const testWindow = window as Window & { + __BUZZ_E2E_REVOKED_OBJECT_URLS__?: string[]; + }; + return testWindow.__BUZZ_E2E_REVOKED_OBJECT_URLS__ ?? []; + }), + ) + .toContain(supersededPreviewUrl); + await page.waitForTimeout(6_000); + await expect(page.getByText("Avatar couldn’t finish uploading")).toHaveCount( + 0, + ); + await expect(page.getByRole("button", { name: "Retry" })).toHaveCount(0); }); test("membership denial on community profile save offers recovery", async ({ From daeaf7c33d5415199a33cbc3dab00244fad5c219 Mon Sep 17 00:00:00 2001 From: klopez4212 Date: Thu, 23 Jul 2026 20:50:29 +0100 Subject: [PATCH 4/4] Refine channel lifecycle settings (#2427) --- desktop/src/features/channels/hooks.ts | 13 +- .../channels/ui/ChannelManagementSheet.tsx | 308 ++++++++---------- .../ui/ChannelManagementSheetRows.tsx | 36 -- .../channels/ui/ChannelMembersBar.tsx | 2 +- .../ui/ChannelPermissionsSettings.tsx | 98 ++++++ .../channels/ui/ChannelTypePicker.tsx | 91 ++++++ .../channels/ui/ChannelTypeSettings.tsx | 161 +++++++++ .../features/channels/ui/channelFormStyles.ts | 5 + .../sidebar/lib/useCreateChannelForm.ts | 19 +- .../sidebar/ui/CreateChannelFormFields.tsx | 193 ++--------- desktop/tests/e2e/channel-controls.spec.ts | 281 +++++++++++++--- desktop/tests/e2e/channels.spec.ts | 139 +++++--- desktop/tests/e2e/integration.spec.ts | 31 +- .../welcome-agent-modal-screenshots.spec.ts | 6 +- 14 files changed, 876 insertions(+), 507 deletions(-) create mode 100644 desktop/src/features/channels/ui/ChannelPermissionsSettings.tsx create mode 100644 desktop/src/features/channels/ui/ChannelTypePicker.tsx create mode 100644 desktop/src/features/channels/ui/ChannelTypeSettings.tsx create mode 100644 desktop/src/features/channels/ui/channelFormStyles.ts diff --git a/desktop/src/features/channels/hooks.ts b/desktop/src/features/channels/hooks.ts index 121e2dbb1..9003d0f5a 100644 --- a/desktop/src/features/channels/hooks.ts +++ b/desktop/src/features/channels/hooks.ts @@ -348,13 +348,10 @@ export function useUpdateChannelMutation(channelId: string | null) { return updateChannel({ ...input, channelId }); }, + onMutate: () => ({ channelId }), onSuccess: (updatedChannel) => { - if (!channelId) { - return; - } - queryClient.setQueryData( - channelDetailQueryKey(channelId), + channelDetailQueryKey(updatedChannel.id), updatedChannel, ); queryClient.setQueryData(channelsQueryKey, (current = []) => @@ -365,7 +362,7 @@ export function useUpdateChannelMutation(channelId: string | null) { ), ); }, - onSettled: () => { + onSettled: (_data, _error, _variables, context) => { // refetchType "none": onSuccess already cached the relay-returned detail; // awaiting the full channel-list refetch kept the edit dialog stuck on // "Saving..." (same failure #1360 fixed for create). @@ -373,9 +370,9 @@ export function useUpdateChannelMutation(channelId: string | null) { queryKey: channelsQueryKey, refetchType: "none", }); - if (channelId) { + if (context?.channelId) { void queryClient.invalidateQueries({ - queryKey: channelDetailQueryKey(channelId), + queryKey: channelDetailQueryKey(context.channelId), refetchType: "none", }); } diff --git a/desktop/src/features/channels/ui/ChannelManagementSheet.tsx b/desktop/src/features/channels/ui/ChannelManagementSheet.tsx index 57c1f6c14..3c1727e00 100644 --- a/desktop/src/features/channels/ui/ChannelManagementSheet.tsx +++ b/desktop/src/features/channels/ui/ChannelManagementSheet.tsx @@ -27,8 +27,6 @@ import { useDeleteChannelMutation, useJoinChannelMutation, useLeaveChannelMutation, - useSetChannelPurposeMutation, - useSetChannelTopicMutation, useUnarchiveChannelMutation, useUpdateChannelMutation, } from "@/features/channels/hooks"; @@ -36,7 +34,6 @@ import { compareMembersByRole } from "@/features/channels/lib/memberUtils"; import { DEFAULT_EPHEMERAL_TTL_SECONDS, formatTtlDuration, - parseTtlDuration, } from "@/features/channels/lib/ephemeralChannel"; import type { Channel } from "@/shared/api/types"; import { cn } from "@/shared/lib/cn"; @@ -45,7 +42,6 @@ import { Button } from "@/shared/ui/button"; import { Dialog, DialogContent, - DialogDescription, DialogHeader, DialogTitle, } from "@/shared/ui/dialog"; @@ -68,6 +64,12 @@ import { PANEL_OVERLAY_CLASS, } from "@/shared/ui/OverlayPanelBackdrop"; import { ChannelCanvas } from "./ChannelCanvas"; +import { + CHANNEL_FORM_FIELD_CONTROL_CLASS, + CHANNEL_FORM_FIELD_SHELL_CLASS, +} from "./channelFormStyles"; +import { ChannelTypeSettings } from "./ChannelTypeSettings"; +import { ChannelPermissionsSettings } from "./ChannelPermissionsSettings"; import { ChannelHero, ChannelQuickAction, @@ -78,7 +80,6 @@ import { IngressRow, NarrativeField, NarrativeGroup, - ToggleRow, } from "./ChannelManagementSheetRows"; import { ChannelManagementModerationActions, @@ -118,13 +119,13 @@ export function ChannelManagementSheet({ const membersQuery = useChannelMembersQuery(channelId, open); const canvasQuery = useCanvasQuery(channelId, channelId !== null && open); const updateChannelDetailsMutation = useUpdateChannelMutation(channelId); - const setTopicMutation = useSetChannelTopicMutation(channelId); - const setPurposeMutation = useSetChannelPurposeMutation(channelId); const archiveChannelMutation = useArchiveChannelMutation(channelId); const unarchiveChannelMutation = useUnarchiveChannelMutation(channelId); const deleteChannelMutation = useDeleteChannelMutation(channelId); const joinChannelMutation = useJoinChannelMutation(channelId); const leaveChannelMutation = useLeaveChannelMutation(channelId); + const channelIdRef = React.useRef(channelId); + channelIdRef.current = channelId; const detail = detailsQuery.data ?? channel; const members = React.useMemo(() => { @@ -159,13 +160,17 @@ export function ChannelManagementSheet({ const [nameDraft, setNameDraft] = React.useState(""); const [descriptionDraft, setDescriptionDraft] = React.useState(""); - const [topicDraft, setTopicDraft] = React.useState(""); - const [purposeDraft, setPurposeDraft] = React.useState(""); const [isPrivateDraft, setIsPrivateDraft] = React.useState(false); const [isEphemeralDraft, setIsEphemeralDraft] = React.useState(false); - const [ttlDraft, setTtlDraft] = React.useState(""); + const [ttlSecondsDraft, setTtlSecondsDraft] = React.useState( + DEFAULT_EPHEMERAL_TTL_SECONDS, + ); const [isDeleteDialogOpen, setIsDeleteDialogOpen] = React.useState(false); const [isEditDialogOpen, setIsEditDialogOpen] = React.useState(false); + const [isConvertingVisibility, setIsConvertingVisibility] = + React.useState(false); + const [hasUserEditedChannelDraft, setHasUserEditedChannelDraft] = + React.useState(false); const [activeView, setActiveView] = React.useState<"summary" | "canvas">( "summary", ); @@ -194,13 +199,10 @@ export function ChannelManagementSheet({ setNameDraft(detail.name); setDescriptionDraft(detail.description); - setTopicDraft(detail.topic ?? ""); - setPurposeDraft(detail.purpose ?? ""); setIsPrivateDraft(detail.visibility === "private"); setIsEphemeralDraft(detail.ttlSeconds !== null); - setTtlDraft( - detail.ttlSeconds !== null ? formatTtlDuration(detail.ttlSeconds) : "", - ); + setTtlSecondsDraft(detail.ttlSeconds ?? DEFAULT_EPHEMERAL_TTL_SECONDS); + setHasUserEditedChannelDraft(false); setActiveView("summary"); }, [detail, open]); @@ -232,19 +234,13 @@ export function ChannelManagementSheet({ onOpenChange(next); } - // Parsed seconds for the ephemeral TTL field. `null` when the field is empty - // or malformed; the form blocks saving on a non-empty malformed value. - const parsedTtlSeconds = parseTtlDuration(ttlDraft); - const ttlInvalid = - isEphemeralDraft && ttlDraft.trim() !== "" && parsedTtlSeconds === null; - const currentVisibility = detail?.visibility ?? channel.visibility; const currentTtlSeconds = detail?.ttlSeconds ?? null; const nextVisibility: "open" | "private" = isPrivateDraft ? "private" : "open"; const nextTtlSeconds: number | null = isEphemeralDraft - ? (parsedTtlSeconds ?? DEFAULT_EPHEMERAL_TTL_SECONDS) + ? ttlSecondsDraft : null; const lifecycleDirty = nextVisibility !== currentVisibility || @@ -254,22 +250,11 @@ export function ChannelManagementSheet({ const nameDirty = nameDraft.trim() !== resolvedChannel.name.trim(); const descriptionDirty = descriptionDraft.trim() !== resolvedChannel.description.trim(); - const topicDirty = topicDraft.trim() !== (resolvedChannel.topic ?? "").trim(); - const purposeDirty = - purposeDraft.trim() !== (resolvedChannel.purpose ?? "").trim(); - const isSavingChannelEdits = - updateChannelDetailsMutation.isPending || - setTopicMutation.isPending || - setPurposeMutation.isPending; - const hasChannelEditChanges = - nameDirty || - descriptionDirty || - lifecycleDirty || - topicDirty || - purposeDirty; + const isSavingChannelEdits = updateChannelDetailsMutation.isPending; + const hasChannelEditChanges = nameDirty || descriptionDirty || lifecycleDirty; const canSaveChannelEdits = nameDraft.trim().length > 0 && - !ttlInvalid && + hasUserEditedChannelDraft && hasChannelEditChanges && !isSavingChannelEdits; const canvasContent = canvasQuery.data?.content?.trim() ?? ""; @@ -279,6 +264,18 @@ export function ChannelManagementSheet({ : undefined; const canOpenCanvas = hasCanvas || canEditNarrative; + function handleEditDialogOpenChange(next: boolean) { + if (!next) { + setNameDraft(resolvedChannel.name); + setDescriptionDraft(resolvedChannel.description); + setIsEphemeralDraft(currentTtlSeconds !== null); + setTtlSecondsDraft(currentTtlSeconds ?? DEFAULT_EPHEMERAL_TTL_SECONDS); + setHasUserEditedChannelDraft(false); + } + + setIsEditDialogOpen(next); + } + async function handleSaveChannelEdits() { try { if (nameDirty || descriptionDirty || lifecycleDirty) { @@ -294,20 +291,32 @@ export function ChannelManagementSheet({ }); } - if (topicDirty) { - await setTopicMutation.mutateAsync({ topic: topicDraft.trim() }); - } - - if (purposeDirty) { - await setPurposeMutation.mutateAsync({ purpose: purposeDraft.trim() }); - } - + setHasUserEditedChannelDraft(false); setIsEditDialogOpen(false); } catch { // React Query stores mutation errors; keep the dialog open and render them. } } + async function handleConvertVisibility(visibility: "open" | "private") { + if (visibility === currentVisibility) { + return; + } + setIsConvertingVisibility(true); + try { + const updatedChannel = await updateChannelDetailsMutation.mutateAsync({ + visibility, + }); + if (channelIdRef.current === updatedChannel.id) { + setIsPrivateDraft(visibility === "private"); + } + } catch { + // React Query stores mutation errors; keep the dialog open and render them. + } finally { + setIsConvertingVisibility(false); + } + } + return ( - + +
- Edit channel - - Update settings for{" "} - {resolvedChannel.name}. - + + Edit {currentVisibility === "private" ? "private" : "public"}{" "} + channel +
-
+
- setNameDraft(event.target.value)} - value={nameDraft} - /> +
+ { + setNameDraft(event.target.value); + setHasUserEditedChannelDraft(true); + }} + value={nameDraft} + /> +
-