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/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) => { 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/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 }; } 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} + /> +
-