From 94c533b37c25a5a7149ede77d95d8383d6af90bf Mon Sep 17 00:00:00 2001 From: cyberzero000 Date: Sat, 22 Aug 2026 11:37:23 -0700 Subject: [PATCH 1/2] fix(agents): track channel membership, not a stale profile copy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Desktop's `@mention` picker gated on `agent.channelIds.includes(channelId)` — the `channel_ids` array from the agent's kind:10100 profile. Nothing kept that array in sync with membership: `buzz-acp` only reads channels, and the relay never authors a 10100. So an agent invited to a new channel subscribed immediately and answered anything p-tagged, but never appeared in that channel's picker, and typed text was not consulted outside DMs. The agent was present, listening, and unmentionable until an operator republished the profile by hand. Mobile never had the bug because it resolves `@name` against relay membership. Fixed on both sides, either sufficient alone: - Desktop accepts relay membership in place of `channel_ids`. Membership is maintained by the relay and cannot drift. `respond_to` and `respondToAllowlist` are still enforced, so this widens discovery, not authority. - `buzz-acp` updates its profile when it observes a membership change it already handles. The update is a queued delta, not a rewrite from this harness's subscription set: that set is narrowed by `channels_override`, by rule matching, and by any channel whose startup subscribe failed, so publishing it wholesale would delete every channel this process happens not to serve, and two harnesses sharing a pubkey would flap the field against each other. Read-modify-write, so fields it does not own survive; a missing profile is left missing, because creating one belongs to deploy tooling. Deltas drain through a single ordered worker. kind:10100 is replaceable and read-modify-written, so two concurrent appliers would both start from the pre-change profile and the later write would drop the earlier channel, and two out-of-order appliers for the same channel would settle on the wrong answer. Queued deltas are coalesced into one publish because membership churn arrives in bursts and each replaceable write needs its own second. A failed publish keeps its deltas and retries with backoff rather than dropping them behind a single warn line. Signed-off-by: cyberzero000 --- crates/buzz-acp/src/lib.rs | 311 ++++++++++++++++++ .../lib/agentAutocompleteEligibility.test.mjs | 67 ++++ .../lib/agentAutocompleteEligibility.ts | 29 +- .../messages/lib/agentMentionRevalidation.ts | 7 + .../src/features/messages/lib/useMentions.ts | 24 +- desktop/src/shared/lib/pubkey.ts | 11 + 6 files changed, 432 insertions(+), 17 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 146214197a8..0ce2e2a4448 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -106,6 +106,298 @@ async fn publish_presence( Ok(()) } +/// Applies this harness's kind:10100 `channel_ids` changes one at a time. +/// +/// Membership deltas must be applied strictly in order and strictly one at a +/// time. kind:10100 is replaceable and read-modify-written, so two concurrent +/// appliers would both start from the pre-change profile and the later write +/// would drop the earlier channel; two out-of-order appliers for the *same* +/// channel would settle on the wrong answer (an add landing after its own +/// remove leaves the agent advertising a channel it left). A single worker +/// draining an ordered queue gives both properties. +struct ProfileChannelUpdater { + tx: tokio::sync::mpsc::UnboundedSender<(Uuid, bool)>, +} + +impl ProfileChannelUpdater { + fn spawn(rest_client: relay::RestClient, keys: nostr::Keys) -> Self { + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<(Uuid, bool)>(); + tokio::spawn(async move { + // The relay breaks a same-second replaceable tie by lowest event id + // and rolls the loser back while still reporting the write as + // accepted, so never sign two of these inside one second. + let mut last_published_secs = 0u64; + // Deltas from a publish that failed. The queue is the only copy, so + // they are carried into the next attempt rather than dropped: a + // peer's future timestamp or a backwards clock step would otherwise + // leave `channel_ids` wrong indefinitely behind a single warn line. + let mut deltas: Vec<(Uuid, bool)> = Vec::new(); + // Set while `deltas` holds work from a failed publish. Waiting only + // on `rx` would mean the retry never happens until the *next* + // membership change arrives — and membership churn is rare enough + // that "next change" can be never for the life of the process, + // leaving the agent absent from every `@mention` picker behind a + // single warn line. + let mut retry_delay: Option = None; + loop { + let received = match retry_delay { + None => match rx.recv().await { + Some(delta) => Some(delta), + None => break, + }, + Some(delay) => tokio::select! { + queued = rx.recv() => match queued { + Some(delta) => Some(delta), + None => break, + }, + _ = tokio::time::sleep(delay) => None, + }, + }; + // Coalesce everything already queued into a single publish. + // Membership churn arrives in bursts — a bulk invite queues one + // delta per channel within a few milliseconds — and each write + // needs its own second, so publishing them one at a time would + // demand more distinct seconds than the burst spans. Applying + // them together needs exactly one. + if let Some(first) = received { + deltas.push(first); + } + while let Ok(next) = rx.try_recv() { + deltas.push(next); + } + if deltas.is_empty() { + retry_delay = None; + continue; + } + match apply_profile_channel_delta(&rest_client, &keys, &deltas, last_published_secs) + .await + { + Ok(Some(published_secs)) => { + last_published_secs = published_secs; + deltas.clear(); + retry_delay = None; + } + Ok(None) => { + deltas.clear(); + retry_delay = None; + } + Err(e) => { + let next = retry_delay + .map(|delay| (delay * 2).min(MAX_PROFILE_RETRY_DELAY)) + .unwrap_or(BASE_PROFILE_RETRY_DELAY); + retry_delay = Some(next); + tracing::warn!( + "failed to republish agent profile, retrying in {next:?}: {e}" + ); + } + } + } + }); + Self { tx } + } + + /// Queue a membership change. Unbounded because dropping a delta would + /// leave the profile permanently stale, and membership churn is rare. + fn record(&self, channel_id: Uuid, joined: bool) { + if self.tx.send((channel_id, joined)).is_err() { + tracing::warn!("profile channel updater stopped; skipping republish"); + } + } +} + +/// How far ahead of local time a republish may stamp its `created_at`. +/// +/// Bounds how much peer clock skew we will out-bid before giving up loudly. +const MAX_PROFILE_PUBLISH_SKEW_SECS: u64 = 5; + +/// First retry delay after a failed profile republish. +const BASE_PROFILE_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); + +/// Ceiling for the profile-republish retry backoff. A stale `channel_ids` keeps +/// the agent out of `@mention` pickers, so this stays low enough to self-heal +/// within a work session. +const MAX_PROFILE_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(300); + +fn unix_secs_now() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0) +} + +/// Apply one membership change to this agent's self-declared `channel_ids`. +/// +/// `channel_ids` is a field desktop's @mention picker reads to decide whether +/// an agent is invocable in a channel, and nothing kept it in sync. An agent +/// invited to a new channel subscribed immediately but stayed absent from that +/// channel's picker until an operator republished the profile by hand. +/// +/// The update is a delta, not a rewrite from this harness's subscription set. +/// That set is narrowed by `channels_override`, by rule matching, and by any +/// channel whose startup subscribe failed — publishing it wholesale would +/// delete every channel this process happens not to serve, and two harnesses +/// sharing a pubkey with different overrides would flap the field against each +/// other. +/// +/// The profile is read-modify-written so fields this harness does not own +/// (`display_name`, `capabilities`, `channel_add_policy`) survive. A missing +/// profile is left missing: creating one is the deploy tooling's job, and a +/// partial profile invented here would read as a downgrade to clients. The +/// human-readable `channels` array is deliberately untouched — nothing gates on +/// it, and naming a just-joined channel would cost another round trip. +/// Returns the second the new profile was published in, or `None` when nothing +/// needed publishing. +async fn apply_profile_channel_delta( + rest_client: &relay::RestClient, + keys: &nostr::Keys, + deltas: &[(Uuid, bool)], + last_published_secs: u64, +) -> Result, relay::RelayError> { + use buzz_core::kind::KIND_AGENT_PROFILE; + use nostr::{EventBuilder, Kind}; + + let filter = nostr::Filter::new() + .kind(Kind::Custom(KIND_AGENT_PROFILE as u16)) + .author(keys.public_key()) + .limit(1); + let existing = rest_client.query(&[filter]).await?; + let existing_event = existing.as_array().and_then(|events| events.first()); + // Anyone may have published this profile — deploy tooling, `buzz channels + // set-add-policy`, another harness. Treat its `created_at` as a floor so + // our write cannot collide with theirs in the same second either. + let last_published_secs = last_published_secs.max( + existing_event + .and_then(|event| event.get("created_at")) + .and_then(serde_json::Value::as_u64) + .unwrap_or(0), + ); + let content = existing_event + .and_then(|event| event.get("content")) + .and_then(serde_json::Value::as_str) + .unwrap_or_default(); + let Ok(serde_json::Value::Object(mut profile)) = + serde_json::from_str::(content) + else { + tracing::debug!("no kind:10100 profile to update — skipping republish"); + return Ok(None); + }; + + let mut ids: Vec = profile + .get("channel_ids") + .and_then(serde_json::Value::as_array) + .map(|values| { + values + .iter() + .filter_map(serde_json::Value::as_str) + .map(str::to_string) + .collect() + }) + .unwrap_or_default(); + // Applied in queue order, so a join/leave/rejoin of one channel settles on + // the last observed state rather than on whichever delta happens to win. + let mut changed = false; + for (channel_id, joined) in deltas { + let target = channel_id.to_string(); + let had = ids.iter().any(|id| id == &target); + if *joined == had { + tracing::debug!(channel_id = %channel_id, "profile channel_ids already current"); + continue; + } + if *joined { + ids.push(target); + } else { + ids.retain(|id| id != &target); + } + changed = true; + } + if !changed { + return Ok(None); + } + ids.sort(); + ids.dedup(); + profile.insert( + "channel_ids".to_string(), + serde_json::Value::Array(ids.into_iter().map(serde_json::Value::String).collect()), + ); + + // Never sign two of these inside one second — see `ProfileChannelUpdater`. + // + // The floor is seeded from the stored profile's `created_at`, which any + // peer can set: deploy tooling, `buzz channels set-add-policy`, or another + // harness on a skewed clock. Sleeping until the wall clock passes it would + // park the single worker — and grow its unbounded queue — for the whole + // skew, and giving up on that sleep would silently sign a losing + // `created_at` with nothing to requeue the delta. Stamp the timestamp past + // the floor instead: it wins the LWW compare without waiting at all. + let now = unix_secs_now(); + let created_at = now.max(last_published_secs + 1); + let lead = created_at.saturating_sub(now); + // Relays reject timestamps too far in the future, so a wildly skewed peer + // is unwinnable. Fail loudly rather than publishing an event that cannot + // land — the worker logs it and leaves its own clock untouched. + // + // Only a peer can reach this. Our own writes cannot, because of the sleep + // below: it hands the lead back before returning, so each publish starts + // from `now` again instead of the stamps ratcheting a second higher every + // time and eventually tripping this guard on the harness's own traffic. + if lead > MAX_PROFILE_PUBLISH_SKEW_SECS { + return Err(relay::RelayError::Http(format!( + "stored agent profile is timestamped {lead}s ahead of local time; refusing to republish" + ))); + } + // Publishing is capped at one per second by the distinct-`created_at` + // requirement. Sleeping here rather than returning an error is what keeps + // that a rate limit instead of a data loss: the queue is the only copy of + // these deltas, and the worker has no requeue path. + if lead > 0 { + tokio::time::sleep(std::time::Duration::from_secs(lead)).await; + } + + let event = EventBuilder::new( + Kind::Custom(KIND_AGENT_PROFILE as u16), + serde_json::Value::Object(profile).to_string(), + ) + .tags([]) + .custom_created_at(nostr::Timestamp::from_secs(created_at)) + .sign_with_keys(keys) + .map_err(|e| relay::RelayError::Http(format!("profile sign error: {e}")))?; + // `submit_event` awaits the relay's response, unlike the socket publisher, + // which returns as soon as the command is queued. The next delta reads the + // profile back, so the write has to be durable before this returns. + let response = rest_client.submit_event(&event).await?; + // A 200 does not mean the profile changed. Two distinct failures hide + // behind one: + // + // - `accepted: false` — the relay refused the event outright. + // - `accepted: true` with a `duplicate:` message — the write was stored + // and then rolled back because a peer's copy won the replaceable + // compare (`Db::replace_addressable_event` reports `was_inserted = + // false`, and the ingest path maps that to accepted-plus-duplicate). + // + // The second is the one this design actually invites, and reading only + // `accepted` misses it: the worker would record a publish that never + // landed and clear the deltas, which are its only copy. + let accepted = response + .get("accepted") + .and_then(serde_json::Value::as_bool) + != Some(false); + let message = response + .get("message") + .and_then(serde_json::Value::as_str) + .unwrap_or_default(); + if !accepted || message.starts_with("duplicate:") { + let detail = if message.is_empty() { + "relay did not accept the profile update" + } else { + message + }; + return Err(relay::RelayError::Http(format!( + "profile channel_ids update did not land: {detail}" + ))); + } + Ok(Some(created_at)) +} + fn emit_runtime_lifecycle( observer: Option<&observer::ObserverHandle>, start_nonce: &str, @@ -2025,6 +2317,8 @@ async fn tokio_main() -> Result<()> { let presence_publisher = relay.event_publisher(); let presence_keys = config.keys.clone(); + let profile_channel_updater = + ProfileChannelUpdater::spawn(relay.rest_client(), config.keys.clone()); // Priority: BUZZ_AUTH_TAG (NIP-OA attestation) → --agent-owner flag. let startup_owner: Option = resolve_agent_owner(&config); @@ -2682,6 +2976,10 @@ async fn tokio_main() -> Result<()> { } membership_newest_ts.insert(ch, ts); + // `Some(joined)` once we know this agent's served + // channel set includes/excludes `ch`; `None` when + // the notification does not apply to us. + let mut profile_channel_delta: Option = None; if kind_u32 == KIND_MEMBER_ADDED_NOTIFICATION { // Clear removal tracking so sessions are not // stripped for a legitimately re-added channel. @@ -2689,18 +2987,21 @@ async fn tokio_main() -> Result<()> { if subscribed_channel_ids.contains(&ch) { tracing::debug!(channel_id = %ch, "membership notification: channel already subscribed"); + profile_channel_delta = Some(true); } else if let Some(filter) = config::resolve_dynamic_channel_filter(&config, ch, &rules) { tracing::info!(channel_id = %ch, "membership notification: subscribing to new channel"); if let Err(e) = relay.subscribe_channel_from(ch, filter, Some(ts)).await { tracing::warn!("failed to subscribe to new channel {ch}: {e}"); } else { subscribed_channel_ids.insert(ch); + profile_channel_delta = Some(true); } } else { tracing::debug!(channel_id = %ch, "membership notification: no matching rules — skipping"); } } else { subscribed_channel_ids.remove(&ch); + profile_channel_delta = Some(false); tracing::info!(channel_id = %ch, "membership notification: unsubscribing from channel"); if let Err(e) = relay.unsubscribe_channel(ch).await { tracing::warn!("failed to unsubscribe from channel {ch}: {e}"); @@ -2743,6 +3044,16 @@ async fn tokio_main() -> Result<()> { ); } } + // The set of channels this agent serves just + // changed. Refresh the self-declared copy in + // kind:10100 so discovery surfaces track it + // without an operator republishing by hand. + // Queued rather than applied here: the update + // is a relay round trip, and the deltas must + // land in the order they were observed. + if let Some(joined) = profile_channel_delta { + profile_channel_updater.record(ch, joined); + } continue; } diff --git a/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs b/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs index 21880eca2ff..4900d38c1c6 100644 --- a/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs +++ b/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs @@ -538,3 +538,70 @@ test("coalesceAgentAutocompleteCandidates: leaves non-agents alone", () => { assert.deepEqual(coalesce([first, second]), [first, second]); }); + +test("relayAgentCanRespondInChannel: relay membership stands in for a stale channelIds", () => { + // The agent was invited to "fresh" but its kind:10100 still lists only + // "general" — nothing republishes that array when membership changes. + const agent = { + pubkey: PUB_A, + respondTo: "anyone", + respondToAllowlist: [], + channelIds: ["general"], + }; + + assert.equal( + relayAgentCanRespondInChannel(agent, "fresh", CURRENT_PUBKEY), + false, + ); + assert.equal( + relayAgentCanRespondInChannel( + agent, + "fresh", + CURRENT_PUBKEY, + new Set([PUB_A]), + ), + true, + ); +}); + +test("relayAgentCanRespondInChannel: membership does not bypass an allowlist", () => { + const agent = { + pubkey: PUB_A, + respondTo: "allowlist", + respondToAllowlist: [OTHER_OWNER_PUBKEY], + channelIds: ["general"], + }; + + assert.equal( + relayAgentCanRespondInChannel( + agent, + "fresh", + CURRENT_PUBKEY, + new Set([PUB_A]), + ), + false, + ); +}); + +test("getMentionableAgentPubkeys: channel scope accepts a member with stale channelIds", () => { + const agent = { + pubkey: PUB_A, + respondTo: "anyone", + respondToAllowlist: [], + channelIds: ["general"], + }; + + assert.deepEqual( + [ + ...getMentionableAgentPubkeys({ + channelMemberPubkeys: new Set([PUB_A]), + currentPubkey: CURRENT_PUBKEY, + eligibilityScope: { type: "channel", channelId: "fresh" }, + managedAgentPubkeys: [], + relayAgents: [agent], + sharedChannelIds: new Set(["fresh"]), + }), + ], + [PUB_A], + ); +}); diff --git a/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts b/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts index 4e1c787f92e..c227b07a37a 100644 --- a/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts +++ b/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts @@ -9,6 +9,14 @@ export function getSharedChannelIds(channels: readonly Channel[] | undefined) { ); } +/** + * `sharesChannel` lets a caller assert the shared channel from authoritative + * state — relay membership — instead of the agent's self-declared + * `channelIds`. That array is a kind:10100 field nothing republishes when + * membership changes, so an agent invited to a new channel stays absent from + * its picker until someone reissues the profile by hand. Membership is + * maintained by the relay and cannot drift. `respondTo` is still enforced. + */ export function relayAgentIsSharedWithUser( agent: Pick< RelayAgent, @@ -16,6 +24,7 @@ export function relayAgentIsSharedWithUser( >, sharedChannelIds: ReadonlySet, currentPubkey?: string | null, + sharesChannel = false, ) { const normalizedCurrentPubkey = currentPubkey ? normalizePubkey(currentPubkey) @@ -37,21 +46,30 @@ export function relayAgentIsSharedWithUser( return ( agent.respondTo === "anyone" && - agent.channelIds.some((channelId) => sharedChannelIds.has(channelId)) + (sharesChannel || + agent.channelIds.some((channelId) => sharedChannelIds.has(channelId))) ); } export function relayAgentCanRespondInChannel( agent: Pick< RelayAgent, - "channelIds" | "ownerPubkey" | "respondTo" | "respondToAllowlist" + "channelIds" | "ownerPubkey" | "pubkey" | "respondTo" | "respondToAllowlist" >, channelId: string, currentPubkey?: string | null, + channelMemberPubkeys?: ReadonlySet, ) { + const isChannelMember = + channelMemberPubkeys?.has(normalizePubkey(agent.pubkey)) === true; return ( - agent.channelIds.includes(channelId) && - relayAgentIsSharedWithUser(agent, new Set([channelId]), currentPubkey) + (isChannelMember || agent.channelIds.includes(channelId)) && + relayAgentIsSharedWithUser( + agent, + new Set([channelId]), + currentPubkey, + isChannelMember, + ) ); } @@ -61,12 +79,14 @@ export type AgentEligibilityScope = | { type: "managed-only" }; export function getMentionableAgentPubkeys({ + channelMemberPubkeys, currentPubkey, eligibilityScope, managedAgentPubkeys, relayAgents, sharedChannelIds, }: { + channelMemberPubkeys?: ReadonlySet; currentPubkey?: string | null; eligibilityScope: AgentEligibilityScope; managedAgentPubkeys: Iterable; @@ -87,6 +107,7 @@ export function getMentionableAgentPubkeys({ agent, eligibilityScope.channelId, currentPubkey, + channelMemberPubkeys, ); if (isAllowed) { pubkeys.add(normalizePubkey(agent.pubkey)); diff --git a/desktop/src/features/messages/lib/agentMentionRevalidation.ts b/desktop/src/features/messages/lib/agentMentionRevalidation.ts index 37f7ce9d4e3..0298424fa3c 100644 --- a/desktop/src/features/messages/lib/agentMentionRevalidation.ts +++ b/desktop/src/features/messages/lib/agentMentionRevalidation.ts @@ -17,6 +17,7 @@ type DirectoryResult = { export async function revalidateAgentMentionPubkeys({ pubkeys, agentPubkeys, + channelMemberPubkeys, currentPubkey, eligibilityScope, sharedChannelIds, @@ -25,6 +26,7 @@ export async function revalidateAgentMentionPubkeys({ }: { pubkeys: readonly string[]; agentPubkeys: ReadonlySet; + channelMemberPubkeys?: ReadonlySet; currentPubkey: string | null; eligibilityScope: AgentEligibilityScope; sharedChannelIds: ReadonlySet; @@ -51,6 +53,7 @@ export async function revalidateAgentMentionPubkeys({ managedResult.data.map((agent) => normalizePubkey(agent.pubkey)), ); const mentionablePubkeys = getMentionableAgentPubkeys({ + channelMemberPubkeys, currentPubkey, eligibilityScope, managedAgentPubkeys: managedPubkeys, @@ -76,6 +79,7 @@ export async function revalidateAgentMentionPubkeys({ export function useAgentMentionRevalidation({ agentPubkeys, + channelMemberPubkeys, getSelectedAgentPubkeys, currentPubkey, eligibilityScope, @@ -83,6 +87,7 @@ export function useAgentMentionRevalidation({ refetchManagedAgents, }: { agentPubkeys: ReadonlySet; + channelMemberPubkeys?: ReadonlySet; getSelectedAgentPubkeys: () => ReadonlySet; currentPubkey: string | null; eligibilityScope: AgentEligibilityScope; @@ -94,6 +99,7 @@ export function useAgentMentionRevalidation({ revalidateAgentMentionPubkeys({ pubkeys, agentPubkeys: new Set([...agentPubkeys, ...getSelectedAgentPubkeys()]), + channelMemberPubkeys, currentPubkey, eligibilityScope, sharedChannelIds, @@ -108,6 +114,7 @@ export function useAgentMentionRevalidation({ }), [ agentPubkeys, + channelMemberPubkeys, currentPubkey, eligibilityScope, getSelectedAgentPubkeys, diff --git a/desktop/src/features/messages/lib/useMentions.ts b/desktop/src/features/messages/lib/useMentions.ts index 30b7a1f48a5..21cb3dbe990 100644 --- a/desktop/src/features/messages/lib/useMentions.ts +++ b/desktop/src/features/messages/lib/useMentions.ts @@ -34,7 +34,7 @@ import type { AutocompleteEdit } from "./useRichTextEditor"; import type { ChannelMember, ChannelType } from "@/shared/api/types"; import type { UserProfileLookup } from "@/features/profile/lib/identity"; import { detectPrefixQuery } from "@/shared/lib/detectPrefixQuery"; -import { normalizePubkey } from "@/shared/lib/pubkey"; +import { normalizePubkey, normalizePubkeySet } from "@/shared/lib/pubkey"; import { channelMemberPubkeySet } from "@/shared/lib/rosterDerivations"; import { trimMapToSize } from "@/shared/lib/trimMapToSize"; import { flushMentionDebounce } from "./flushMentionDebounce"; @@ -152,12 +152,7 @@ export function useMentions( [managedAgentsQuery.data], ); const managedAgentPubkeys = React.useMemo( - () => - new Set( - (managedAgentsQuery.data ?? []).map((agent) => - normalizePubkey(agent.pubkey), - ), - ), + () => normalizePubkeySet(managedAgentsQuery.data), [managedAgentsQuery.data], ); const relayAgentNamesByPubkey = React.useMemo( @@ -177,9 +172,16 @@ export function useMentions( const mentionChannelId = isAgentMentionChannelType(options?.channelType) ? channelId : null; + // Identity-cached (shared with the timeline's roster derivations) — the + // Set is built once per distinct roster instead of per consumer. + const memberPubkeys = React.useMemo( + () => (members ? channelMemberPubkeySet(members) : new Set()), + [members], + ); const mentionableAgentPubkeys = React.useMemo( () => getMentionableAgentPubkeys({ + channelMemberPubkeys: memberPubkeys, currentPubkey, eligibilityScope: mentionChannelId ? { type: "channel", channelId: mentionChannelId } @@ -191,6 +193,7 @@ export function useMentions( [ currentPubkey, managedAgentPubkeys, + memberPubkeys, mentionChannelId, relayAgentsQuery.data, sharedChannelIds, @@ -225,12 +228,6 @@ export function useMentions( () => new Set(activePersonas.map((persona) => persona.id)), [activePersonas], ); - // Identity-cached (shared with the timeline's roster derivations) — the - // Set is built once per distinct roster instead of per consumer. - const memberPubkeys = React.useMemo( - () => (members ? channelMemberPubkeySet(members) : new Set()), - [members], - ); const agentIdentityPubkeys = React.useMemo( () => getAgentIdentityPubkeys({ @@ -809,6 +806,7 @@ export function useMentions( ).current; const revalidateMentionPubkeys = useAgentMentionRevalidation({ agentPubkeys: agentIdentityPubkeys, + channelMemberPubkeys: memberPubkeys, getSelectedAgentPubkeys, currentPubkey, eligibilityScope: mentionChannelId diff --git a/desktop/src/shared/lib/pubkey.ts b/desktop/src/shared/lib/pubkey.ts index 6dcc48749a3..938edfd87af 100644 --- a/desktop/src/shared/lib/pubkey.ts +++ b/desktop/src/shared/lib/pubkey.ts @@ -8,6 +8,17 @@ export function normalizePubkey(pubkey: string): string { return pubkey.trim().toLowerCase(); } +/** + * Set of normalised pubkeys from anything carrying one — members, agents, + * profiles. Membership tests against a raw list are the usual place a + * case difference turns into a silent miss. + */ +export function normalizePubkeySet( + items: readonly { pubkey: string }[] | undefined, +): Set { + return new Set((items ?? []).map((item) => normalizePubkey(item.pubkey))); +} + /** * The ONE canonical compact display form for a pubkey: `abcd1234…wxyz`. * From 514a099aa0c8772661aab15a36d6fb608afd95e5 Mon Sep 17 00:00:00 2001 From: cyberzero000 Date: Sun, 23 Aug 2026 12:15:51 -0700 Subject: [PATCH 2/2] fix(agents): correct the membership-tracking premise and its failure modes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review of the first commit found the desktop premise wrong and three real defects in the harness path. `channel_ids` was never stale on desktop. `list_relay_agents` replaces the field with relay membership read from relay-signed kind:39002 before the frontend sees it (`relay_directory.rs`), so the kind:10100 array never reaches the picker. What the desktop change actually closes is a refresh-clock gap: `useRelayAgentsQuery` polls every 5 minutes with `refetchOnWindowFocus: false` while the channel roster updates live, so an agent invited to the channel you are looking at is missing from the picker until the next tick. Send-time was never affected — `revalidate_relay_agents` queries `#d: [channel_id]` fresh. Comments and test names now say that instead. The consumer that does read kind:10100 `channel_ids` is mobile's picker, for non-member relay agents (`agentIsSharedWithUser` in `mention_candidates.dart`). That is what the harness half repairs. Three fixes to that half: - `ids.sort()` is gone. `channel_ids` is paired with the profile's human-readable `channels` array by index in `deriveProfileChannels`, so sorting relabelled rows that were correct before. Adds append, removes delete in place, and the dedup preserves order. - The skew guard was 5s against a relay that accepts ±900s. A peer stamped 6-900s ahead is publishable, but the guard refused and retried forever at the 300s backoff ceiling, leaving `channel_ids` stale for the life of the process. Raised to the relay's own window, with the wait it can cost bounded separately: only our own previous publish is worth sleeping through (one second), and a peer's future stamp is out-bid rather than slept through, so the single worker is never parked for minutes. - `accepted` failed open. `!= Some(false)` read a missing field as success, and `submit_event` returns `Value::Null` for an empty 2xx body — the worker would clear deltas it never landed, and the queue is their only copy. Now fails closed and also matches the bare `duplicate` message the relay emits alongside `duplicate:`, as `buzz-cli`'s write path does. A burst that nets to no change no longer costs a replaceable write: `changed` compares the final array against the original rather than counting mutations. The desktop index-pairing is fixed too, not just stopped from getting worse. The backend already returns `channel_ids` in channel-id order against whatever order the profile's `channels` names are in, so the pairing was already wrong: a row could read `#general` and open `#random`. Both sides now resolve against the community channel list instead of by position. Tests: 11 for the extracted `apply_channel_id_deltas`, `profile_publish_stamp`, and `profile_write_landed`, where the first commit had none for 290 lines of queue, LWW-timestamp and response-parsing logic. 3 for `deriveProfileChannels`; 2 of those fail against the index-paired version. cargo test -p buzz-acp --lib: 812 passed, 0 failed. cargo clippy -p buzz-acp --all-targets -D warnings: clean. cargo fmt: clean. desktop: 5403 passed, 0 failed. tsc --noEmit and biome check: clean. Still open, deliberately: a delta arriving during the retry backoff cancels it and republishes immediately, so a persistent failure under churn retries hot. And member-added with no matching rule records no delta while member-removed always does, so `channel_ids` can still omit a channel the agent is a relay member of. Signed-off-by: cyberzero000 --- crates/buzz-acp/src/lib.rs | 338 ++++++++++++++---- .../lib/agentAutocompleteEligibility.test.mjs | 8 +- .../lib/agentAutocompleteEligibility.ts | 17 +- .../profile/ui/UserProfilePanelUtils.test.mjs | 69 ++++ .../profile/ui/UserProfilePanelUtils.ts | 25 +- 5 files changed, 364 insertions(+), 93 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 0ce2e2a4448..1082de0aae3 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -207,8 +207,22 @@ impl ProfileChannelUpdater { /// How far ahead of local time a republish may stamp its `created_at`. /// -/// Bounds how much peer clock skew we will out-bid before giving up loudly. -const MAX_PROFILE_PUBLISH_SKEW_SECS: u64 = 5; +/// Matches the relay's own ingest window (`MAX_TIMESTAMP_DRIFT_SECS`, ±900s): +/// a stamp inside it lands, so refusing one is pure loss — the worker would +/// retry it forever at the backoff ceiling and leave `channel_ids` stale for +/// the life of the process. Past the window the relay rejects the event +/// outright, so a peer skewed further than this is unwinnable and failing +/// loudly is the only honest outcome. +const MAX_PROFILE_PUBLISH_SKEW_SECS: u64 = 900; + +/// Longest wall-clock wait before signing a republish. +/// +/// Only this harness's own previous publish can require a wait, and it can only +/// ever be the single second the distinct-`created_at` rule costs. A peer's +/// future timestamp is out-bid by stamping past it, never by sleeping through +/// it: sleeping would park the single worker — and grow its unbounded queue — +/// for the whole skew, up to `MAX_PROFILE_PUBLISH_SKEW_SECS`. +const MAX_PROFILE_PUBLISH_SLEEP: std::time::Duration = std::time::Duration::from_secs(1); /// First retry delay after a failed profile republish. const BASE_PROFILE_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); @@ -225,12 +239,224 @@ fn unix_secs_now() -> u64 { .unwrap_or(0) } -/// Apply one membership change to this agent's self-declared `channel_ids`. +/// Apply queued membership deltas to a profile's `channel_ids`, in queue order. +/// +/// Order matters: a join/leave/rejoin of one channel settles on the last +/// observed state rather than on whichever delta happens to win. Returns +/// whether the array actually ended up different — an add for a channel already +/// listed, a remove for one never listed, or a burst that nets to nothing must +/// not cost a replaceable write. +/// +/// Adds append and removes delete in place; the array is never sorted. +/// `channel_ids` is paired with the profile's human-readable `channels` array by +/// index in at least one client, so a wholesale reorder relabels rows that were +/// correct before. +fn apply_channel_id_deltas(ids: &mut Vec, deltas: &[(Uuid, bool)]) -> bool { + let before = ids.clone(); + for (channel_id, joined) in deltas { + let target = channel_id.to_string(); + let had = ids.iter().any(|id| id == &target); + if *joined == had { + tracing::debug!(channel_id = %channel_id, "profile channel_ids already current"); + continue; + } + if *joined { + ids.push(target); + } else { + ids.retain(|id| id != &target); + } + } + // Duplicates can only arrive from the stored profile — the `had` check never + // adds one. Order-preserving, for the same reason the array is not sorted. + let mut seen = std::collections::HashSet::new(); + ids.retain(|id| seen.insert(id.clone())); + *ids != before +} + +/// `created_at` for the next republish, and how far ahead of `now` it lands. +/// +/// Never stamps two of these in the same second: the relay breaks a same-second +/// replaceable tie by lowest event id and rolls the loser back while still +/// reporting the write as accepted. `floor` is the newest `created_at` known to +/// be in play — this harness's last publish, or the stored profile's, whichever +/// is later. Stamping past the floor wins the LWW compare without waiting for +/// the wall clock to catch up to a peer's skew. +fn profile_publish_stamp(now: u64, floor: u64) -> (u64, u64) { + let created_at = now.max(floor.saturating_add(1)); + (created_at, created_at.saturating_sub(now)) +} + +/// Whether a `POST /events` response means the republish actually landed. +/// +/// A 2xx does not mean the profile changed. Three distinct failures hide behind +/// one: +/// +/// - `accepted: false` — the relay refused the event outright. +/// - `accepted: true` with a `duplicate` message — the write was stored and then +/// rolled back because a peer's copy won the replaceable compare +/// (`Db::replace_addressable_event` reports `was_inserted = false`, and the +/// ingest path maps that to accepted-plus-duplicate). +/// - no readable `accepted` field — an empty or unexpected body, which +/// `RestClient::submit_event` surfaces as `Value::Null`. +/// +/// All three fail closed. The queue is the worker's only copy of these deltas, +/// so reading an unreadable response as success would drop them silently. +fn profile_write_landed(response: &serde_json::Value) -> Result<(), String> { + let accepted = response + .get("accepted") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false); + let message = response + .get("message") + .and_then(serde_json::Value::as_str) + .unwrap_or_default(); + // The relay emits both bare `duplicate` and `duplicate: `; the CLI's + // write path matches the same two forms. + let superseded = message == "duplicate" || message.starts_with("duplicate:"); + if accepted && !superseded { + return Ok(()); + } + Err(if message.is_empty() { + "relay did not accept the profile update".to_string() + } else { + message.to_string() + }) +} + +#[cfg(test)] +mod profile_channel_tests { + use super::*; + + fn channel(byte: u8) -> Uuid { + Uuid::from_bytes([byte; 16]) + } + + #[test] + fn add_appends_instead_of_sorting() { + let first = channel(0xaa).to_string(); + let mut ids = vec![first.clone()]; + + assert!(apply_channel_id_deltas(&mut ids, &[(channel(0x11), true)])); + + // 0x11… sorts before 0xaa… but must stay second: the profile's + // index-paired `channels` names still point at their own ids. + assert_eq!(ids, vec![first, channel(0x11).to_string()]); + } + + #[test] + fn remove_drops_only_its_own_channel() { + let kept = channel(0xaa).to_string(); + let mut ids = vec![kept.clone(), channel(0xbb).to_string()]; + + assert!(apply_channel_id_deltas(&mut ids, &[(channel(0xbb), false)])); + + assert_eq!(ids, vec![kept]); + } + + #[test] + fn already_current_deltas_cost_no_write() { + let mut ids = vec![channel(0xaa).to_string()]; + + // Add of a listed channel, remove of an unlisted one. + assert!(!apply_channel_id_deltas( + &mut ids, + &[(channel(0xaa), true), (channel(0xbb), false)] + )); + assert_eq!(ids, vec![channel(0xaa).to_string()]); + } + + #[test] + fn a_burst_that_nets_to_nothing_costs_no_write() { + let mut ids = Vec::new(); + + // Retried deltas from a failed publish can carry an add and its own + // later remove. Applied in order they settle on "not a member", which + // is what the stored profile already says. + assert!(!apply_channel_id_deltas( + &mut ids, + &[(channel(0x22), true), (channel(0x22), false)] + )); + assert!(ids.is_empty()); + } + + #[test] + fn later_delta_wins_for_the_same_channel() { + let mut ids = Vec::new(); + + assert!(apply_channel_id_deltas( + &mut ids, + &[(channel(0x33), false), (channel(0x33), true)] + )); + + assert_eq!(ids, vec![channel(0x33).to_string()]); + } + + #[test] + fn stamp_never_reuses_the_floor_second() { + // Same second as the floor: step past it rather than tie and lose. + assert_eq!(profile_publish_stamp(100, 100), (101, 1)); + // Nothing in play: stamp now, wait for nothing. + assert_eq!(profile_publish_stamp(100, 0), (100, 0)); + } + + #[test] + fn peer_skew_is_out_bid_not_slept_through() { + let (created_at, lead) = profile_publish_stamp(1_000, 1_800); + + assert_eq!(created_at, 1_801); + assert_eq!(lead, 801); + // Inside the relay's ingest window, so this publishes; the worker waits + // at most MAX_PROFILE_PUBLISH_SLEEP regardless of the lead. + assert!(lead <= MAX_PROFILE_PUBLISH_SKEW_SECS); + } + + #[test] + fn unwinnable_skew_is_reported_not_published() { + let (_, lead) = profile_publish_stamp(1_000, 1_000 + MAX_PROFILE_PUBLISH_SKEW_SECS); + + assert!(lead > MAX_PROFILE_PUBLISH_SKEW_SECS); + } + + #[test] + fn only_an_explicit_accept_counts_as_landed() { + assert!(profile_write_landed(&serde_json::json!({"accepted": true})).is_ok()); + assert!( + profile_write_landed(&serde_json::json!({"accepted": true, "message": "ok"})).is_ok() + ); + } + + #[test] + fn a_lost_replaceable_compare_is_not_landed() { + // `accepted: true` plus a duplicate message: stored, then rolled back. + for message in ["duplicate", "duplicate:", "duplicate: superseded"] { + assert!( + profile_write_landed(&serde_json::json!({"accepted": true, "message": message})) + .is_err(), + "{message} must not read as landed" + ); + } + } + + #[test] + fn an_unreadable_response_fails_closed() { + // The queue is the only copy of the deltas, so anything we cannot read + // as an accept has to be retried rather than cleared. + assert!(profile_write_landed(&serde_json::Value::Null).is_err()); + assert!(profile_write_landed(&serde_json::json!({})).is_err()); + assert!(profile_write_landed(&serde_json::json!({"accepted": "yes"})).is_err()); + assert!(profile_write_landed(&serde_json::json!({"accepted": false})).is_err()); + } +} + +/// Apply observed membership changes to this agent's self-declared `channel_ids`. /// -/// `channel_ids` is a field desktop's @mention picker reads to decide whether -/// an agent is invocable in a channel, and nothing kept it in sync. An agent -/// invited to a new channel subscribed immediately but stayed absent from that -/// channel's picker until an operator republished the profile by hand. +/// kind:10100's `channel_ids` is what mobile's mention picker consults to decide +/// whether a *non-member* relay agent is mentionable +/// (`agentIsSharedWithUser` in `mention_candidates.dart`), and nothing kept it +/// in sync: `buzz-acp` only read channels, and the relay never authors a 10100. +/// Desktop reads relay membership instead — its Tauri directory replaces this +/// field before the frontend sees it (`relay_directory.rs`) — so desktop is not +/// what this repairs. /// /// The update is a delta, not a rewrite from this harness's subscription set. /// That set is narrowed by `channels_override`, by rule matching, and by any @@ -243,8 +469,11 @@ fn unix_secs_now() -> u64 { /// (`display_name`, `capabilities`, `channel_add_policy`) survive. A missing /// profile is left missing: creating one is the deploy tooling's job, and a /// partial profile invented here would read as a downgrade to clients. The -/// human-readable `channels` array is deliberately untouched — nothing gates on -/// it, and naming a just-joined channel would cost another round trip. +/// human-readable `channels` array is left alone, because naming a just-joined +/// channel would cost another round trip — which is also why `channel_ids` +/// keeps its existing order instead of being sorted: clients that pair the two +/// arrays by index would otherwise relabel rows that were correct before. +/// /// Returns the second the new profile was published in, or `None` when nothing /// needed publishing. async fn apply_profile_channel_delta( @@ -293,28 +522,9 @@ async fn apply_profile_channel_delta( .collect() }) .unwrap_or_default(); - // Applied in queue order, so a join/leave/rejoin of one channel settles on - // the last observed state rather than on whichever delta happens to win. - let mut changed = false; - for (channel_id, joined) in deltas { - let target = channel_id.to_string(); - let had = ids.iter().any(|id| id == &target); - if *joined == had { - tracing::debug!(channel_id = %channel_id, "profile channel_ids already current"); - continue; - } - if *joined { - ids.push(target); - } else { - ids.retain(|id| id != &target); - } - changed = true; - } - if !changed { + if !apply_channel_id_deltas(&mut ids, deltas) { return Ok(None); } - ids.sort(); - ids.dedup(); profile.insert( "channel_ids".to_string(), serde_json::Value::Array(ids.into_iter().map(serde_json::Value::String).collect()), @@ -324,33 +534,30 @@ async fn apply_profile_channel_delta( // // The floor is seeded from the stored profile's `created_at`, which any // peer can set: deploy tooling, `buzz channels set-add-policy`, or another - // harness on a skewed clock. Sleeping until the wall clock passes it would - // park the single worker — and grow its unbounded queue — for the whole - // skew, and giving up on that sleep would silently sign a losing - // `created_at` with nothing to requeue the delta. Stamp the timestamp past - // the floor instead: it wins the LWW compare without waiting at all. + // harness on a skewed clock. let now = unix_secs_now(); - let created_at = now.max(last_published_secs + 1); - let lead = created_at.saturating_sub(now); - // Relays reject timestamps too far in the future, so a wildly skewed peer - // is unwinnable. Fail loudly rather than publishing an event that cannot - // land — the worker logs it and leaves its own clock untouched. - // - // Only a peer can reach this. Our own writes cannot, because of the sleep - // below: it hands the lead back before returning, so each publish starts - // from `now` again instead of the stamps ratcheting a second higher every - // time and eventually tripping this guard on the harness's own traffic. + let (created_at, lead) = profile_publish_stamp(now, last_published_secs); + // Past the relay's ingest window the event cannot land at all, so a peer + // skewed further than this is unwinnable. Fail loudly: the worker logs it, + // keeps the deltas, and retries with backoff. if lead > MAX_PROFILE_PUBLISH_SKEW_SECS { return Err(relay::RelayError::Http(format!( "stored agent profile is timestamped {lead}s ahead of local time; refusing to republish" ))); } - // Publishing is capped at one per second by the distinct-`created_at` - // requirement. Sleeping here rather than returning an error is what keeps - // that a rate limit instead of a data loss: the queue is the only copy of - // these deltas, and the worker has no requeue path. + // Two different leads land here and only one is worth waiting out: + // + // - Our own previous publish, one second back. Sleeping through it is what + // keeps the one-per-second cap a rate limit instead of data loss — the + // queue is the only copy of these deltas and the worker has no requeue + // path — and it hands the second back before returning, so stamps never + // ratchet away from the wall clock. + // - A peer's future timestamp, up to the whole ingest window. That one is + // out-bid by the stamp above, not slept through: parking the single + // worker for minutes would stall every later delta behind it. if lead > 0 { - tokio::time::sleep(std::time::Duration::from_secs(lead)).await; + tokio::time::sleep(std::time::Duration::from_secs(lead).min(MAX_PROFILE_PUBLISH_SLEEP)) + .await; } let event = EventBuilder::new( @@ -365,36 +572,9 @@ async fn apply_profile_channel_delta( // which returns as soon as the command is queued. The next delta reads the // profile back, so the write has to be durable before this returns. let response = rest_client.submit_event(&event).await?; - // A 200 does not mean the profile changed. Two distinct failures hide - // behind one: - // - // - `accepted: false` — the relay refused the event outright. - // - `accepted: true` with a `duplicate:` message — the write was stored - // and then rolled back because a peer's copy won the replaceable - // compare (`Db::replace_addressable_event` reports `was_inserted = - // false`, and the ingest path maps that to accepted-plus-duplicate). - // - // The second is the one this design actually invites, and reading only - // `accepted` misses it: the worker would record a publish that never - // landed and clear the deltas, which are its only copy. - let accepted = response - .get("accepted") - .and_then(serde_json::Value::as_bool) - != Some(false); - let message = response - .get("message") - .and_then(serde_json::Value::as_str) - .unwrap_or_default(); - if !accepted || message.starts_with("duplicate:") { - let detail = if message.is_empty() { - "relay did not accept the profile update" - } else { - message - }; - return Err(relay::RelayError::Http(format!( - "profile channel_ids update did not land: {detail}" - ))); - } + profile_write_landed(&response).map_err(|detail| { + relay::RelayError::Http(format!("profile channel_ids update did not land: {detail}")) + })?; Ok(Some(created_at)) } diff --git a/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs b/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs index 4900d38c1c6..7ca9b841802 100644 --- a/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs +++ b/desktop/src/features/agents/lib/agentAutocompleteEligibility.test.mjs @@ -539,9 +539,9 @@ test("coalesceAgentAutocompleteCandidates: leaves non-agents alone", () => { assert.deepEqual(coalesce([first, second]), [first, second]); }); -test("relayAgentCanRespondInChannel: relay membership stands in for a stale channelIds", () => { - // The agent was invited to "fresh" but its kind:10100 still lists only - // "general" — nothing republishes that array when membership changes. +test("relayAgentCanRespondInChannel: the live roster beats a not-yet-polled channelIds", () => { + // The agent was invited to "fresh", but the relay-agents query last ran + // before that and still reports only "general". The roster is live. const agent = { pubkey: PUB_A, respondTo: "anyone", @@ -583,7 +583,7 @@ test("relayAgentCanRespondInChannel: membership does not bypass an allowlist", ( ); }); -test("getMentionableAgentPubkeys: channel scope accepts a member with stale channelIds", () => { +test("getMentionableAgentPubkeys: channel scope accepts a member the agents poll has not caught up to", () => { const agent = { pubkey: PUB_A, respondTo: "anyone", diff --git a/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts b/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts index c227b07a37a..6526794a622 100644 --- a/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts +++ b/desktop/src/features/agents/lib/agentAutocompleteEligibility.ts @@ -10,12 +10,17 @@ export function getSharedChannelIds(channels: readonly Channel[] | undefined) { } /** - * `sharesChannel` lets a caller assert the shared channel from authoritative - * state — relay membership — instead of the agent's self-declared - * `channelIds`. That array is a kind:10100 field nothing republishes when - * membership changes, so an agent invited to a new channel stays absent from - * its picker until someone reissues the profile by hand. Membership is - * maintained by the relay and cannot drift. `respondTo` is still enforced. + * `sharesChannel` lets a caller assert the shared channel from the live channel + * roster instead of `agent.channelIds`. + * + * Both are relay membership — `list_relay_agents` replaces `channelIds` with + * the membership it reads from relay-signed kind:39002 before the frontend sees + * it — but they refresh on different clocks. `useRelayAgentsQuery` polls every + * `AGENTS_FOCUS_STALE_TIME_MS` (5 min) with `refetchOnWindowFocus: false`, + * while the roster updates live. An agent invited to the channel you are + * looking at is therefore missing from the picker until the next poll tick. + * The roster closes that window. `respondTo` is still enforced, so this widens + * discovery, not authority. */ export function relayAgentIsSharedWithUser( agent: Pick< diff --git a/desktop/src/features/profile/ui/UserProfilePanelUtils.test.mjs b/desktop/src/features/profile/ui/UserProfilePanelUtils.test.mjs index 0c983fa3e84..d7d783a5375 100644 --- a/desktop/src/features/profile/ui/UserProfilePanelUtils.test.mjs +++ b/desktop/src/features/profile/ui/UserProfilePanelUtils.test.mjs @@ -2,6 +2,7 @@ import assert from "node:assert/strict"; import test from "node:test"; import { + deriveProfileChannels, parseProfilePanelTab, parseProfilePanelView, personaManagedAgentUpdate, @@ -244,3 +245,71 @@ test("profile target identity stays stable while a requested pubkey is canonical "persona:requested-persona", ); }); + +const CHANNEL_A = "11111111-1111-4111-8111-111111111111"; +const CHANNEL_B = "22222222-2222-4222-8222-222222222222"; + +function relayAgent(overrides = {}) { + return { + pubkey: "aa".repeat(32), + ownerPubkey: null, + name: "Fizz", + agentType: "agent", + channels: [], + channelIds: [], + capabilities: [], + status: "offline", + respondTo: "anyone", + respondToAllowlist: [], + ...overrides, + }; +} + +const COMMUNITY_CHANNELS = [ + { id: CHANNEL_A, name: "general" }, + { id: CHANNEL_B, name: "random" }, +]; + +test("deriveProfileChannels does not pair names with ids by index", () => { + // `channelIds` arrives in channel-id order from relay membership while + // `channels` keeps whatever order the kind:10100 author wrote. Pairing by + // index here labelled a row "general" and opened #random. + const links = deriveProfileChannels( + "ff".repeat(32), + relayAgent({ + channels: ["random", "general"], + channelIds: [CHANNEL_A, CHANNEL_B], + }), + undefined, + COMMUNITY_CHANNELS, + ); + + assert.deepEqual(links, [ + { id: CHANNEL_A, name: "general" }, + { id: CHANNEL_B, name: "random" }, + ]); +}); + +test("deriveProfileChannels surfaces a channel id the profile never named", () => { + // `buzz-acp` appends to `channel_ids` without naming the channel, so an id + // with no matching name still has to resolve. + const links = deriveProfileChannels( + "ff".repeat(32), + relayAgent({ channels: [], channelIds: [CHANNEL_B] }), + undefined, + COMMUNITY_CHANNELS, + ); + + assert.deepEqual(links, [{ id: CHANNEL_B, name: "random" }]); +}); + +test("deriveProfileChannels keeps a named channel the viewer cannot resolve", () => { + const links = deriveProfileChannels( + "ff".repeat(32), + relayAgent({ channels: ["private-room"], channelIds: [] }), + undefined, + COMMUNITY_CHANNELS, + ); + + assert.deepEqual(links, [{ id: "private-room", name: "private-room" }]); +}); diff --git a/desktop/src/features/profile/ui/UserProfilePanelUtils.ts b/desktop/src/features/profile/ui/UserProfilePanelUtils.ts index be1a57c112e..4f13f2623d3 100644 --- a/desktop/src/features/profile/ui/UserProfilePanelUtils.ts +++ b/desktop/src/features/profile/ui/UserProfilePanelUtils.ts @@ -129,11 +129,28 @@ export function deriveProfileChannels( const channelsByName = new Map( channels?.map((channel) => [channel.name, channel]) ?? [], ); + const channelsById = new Map( + channels?.map((channel) => [channel.id, channel]) ?? [], + ); - relayAgent?.channels.forEach((name, index) => { - const channel = channelsByName.get(name); - const id = relayAgent.channelIds[index] ?? channel?.id ?? name; - links.set(id, { id, name }); + // `channels` (names) and `channelIds` are independent arrays, so they cannot + // be paired by index. The backend replaces `channelIds` with relay membership + // in channel-id order before this sees it, and `buzz-acp` appends to the + // profile's copy without touching the names — either one makes position N of + // one array a different channel from position N of the other. Pairing them + // anyway labelled a row with another channel's name and navigated to that + // other channel. Resolve each side against the community's channel list. + for (const id of relayAgent?.channelIds ?? []) { + const channel = channelsById.get(id); + if (channel) { + links.set(channel.id, { id: channel.id, name: channel.name }); + } + } + relayAgent?.channels.forEach((name) => { + const id = channelsByName.get(name)?.id ?? name; + if (!links.has(id)) { + links.set(id, { id, name }); + } }); if (managedAgent && channels) {