diff --git a/desktop/src-tauri/src/archive/sync.rs b/desktop/src-tauri/src/archive/sync.rs index 6ce4df92c..d759fb58c 100644 --- a/desktop/src-tauri/src/archive/sync.rs +++ b/desktop/src-tauri/src/archive/sync.rs @@ -31,7 +31,7 @@ use super::{ store::SaveSubscription, ArchiveBatchResult, ArchiveCandidate, MatchedScope, ScopeType, }; use crate::app_state::AppState; -use crate::native_relay_client::{self, MatchedEvent, RelaySession, Subscription}; +use crate::native_relay_client::{MatchedEvent, NativeRelayClient, RelaySession, Subscription}; /// Flush once this many events are buffered. Parity with the renderer manager. const FLUSH_BATCH_SIZE: usize = 25; @@ -496,6 +496,7 @@ pub async fn start_archive_sync( app: AppHandle, state: State<'_, AppState>, sync_state: State<'_, ArchiveSyncState>, + relay_client: State<'_, NativeRelayClient>, epoch: u64, lease: u64, ) -> Result<(), String> { @@ -516,7 +517,7 @@ pub async fn start_archive_sync( // No NIP-OA auth tag: this is the owner's own session, authenticated as // the identity itself, exactly like the renderer's relay client. - let (session, events) = native_relay_client::start(relay_url, keys, None); + let (session, events) = relay_client.archive_session(relay_url, keys).await; let io = AppIo { app: app.clone(), @@ -524,7 +525,7 @@ pub async fn start_archive_sync( }; tauri::async_runtime::spawn(async move { run_sync(&io, reload, events, cancel).await; - session.shutdown(); + session.set_subscriptions(Vec::new()).await; }); Ok(()) } diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index ecc08f29a..bb0512c4f 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -31,6 +31,7 @@ mod native_websocket; mod native_websocket_batch; mod nostr_bind; pub mod nostr_convert; +mod persona_catalog; mod prevent_sleep; mod ptt_shortcut; mod relay; @@ -229,6 +230,7 @@ pub fn run() { .manage(commands::pairing::PairingHandle::new()) .manage(terminal_runtime::TerminalSessions::default()) .manage(archive::sync::ArchiveSyncState::default()) + .manage(native_relay_client::NativeRelayClient::default()) .setup(move |app| { let app_handle = app.handle().clone(); #[cfg(target_os = "macos")] @@ -713,6 +715,7 @@ pub fn run() { update_managed_agent, discover_backend_providers, probe_backend_provider, + persona_catalog::fetch_persona_catalog, list_personas, create_persona, update_persona, diff --git a/desktop/src-tauri/src/native_relay_client.rs b/desktop/src-tauri/src/native_relay_client.rs index 4f36801c5..3f80dcb70 100644 --- a/desktop/src-tauri/src/native_relay_client.rs +++ b/desktop/src-tauri/src/native_relay_client.rs @@ -27,7 +27,7 @@ use std::{ use buzz_ws_client_pkg::{NostrWsConnection, RelayMessage}; use nostr::{Event, Keys}; use tokio::{ - sync::{mpsc, Mutex}, + sync::{broadcast, mpsc, oneshot, Mutex}, time::Instant, }; use tokio_util::sync::CancellationToken; @@ -69,6 +69,7 @@ pub(crate) struct Subscription { /// An event delivered to the session owner, tagged with the subscription that /// matched it. Callers demultiplex on `subscription_id`. +#[derive(Clone)] pub(crate) struct MatchedEvent { pub(crate) subscription_id: String, pub(crate) event: Box, @@ -76,12 +77,84 @@ pub(crate) struct MatchedEvent { /// Handle to a running session. Dropping it does not stop the session; call /// [`RelaySession::shutdown`] so the socket closes deterministically. +/// App-wide owner of the one native socket for the active `(relay, pubkey)` +/// scope. Features subscribe independently, while scope replacement cancels +/// the old socket before exposing the new one. +#[derive(Default)] +pub(crate) struct NativeRelayClient { + current: Mutex>, +} + +struct ManagedSession { + scope: (String, String), + session: Arc, +} + +impl NativeRelayClient { + async fn ensure_session(&self, relay_url: String, keys: Keys) -> Arc { + let scope = (relay_url.clone(), keys.public_key().to_hex()); + let mut current = self.current.lock().await; + if let Some(managed) = current.as_ref().filter(|managed| managed.scope == scope) { + return Arc::clone(&managed.session); + } + if let Some(previous) = current.take() { + previous.session.shutdown(); + } + let session = start_managed(relay_url, keys, None); + *current = Some(ManagedSession { + scope, + session: Arc::clone(&session), + }); + session + } + + pub(crate) async fn session(&self, relay_url: String, keys: Keys) -> Arc { + self.ensure_session(relay_url, keys).await + } + + pub(crate) async fn archive_session( + &self, + relay_url: String, + keys: Keys, + ) -> (Arc, mpsc::Receiver) { + let session = self.ensure_session(relay_url, keys).await; + let mut events = session.subscribe(); + let (event_tx, event_rx) = mpsc::channel(256); + tauri::async_runtime::spawn(async move { + loop { + match events.recv().await { + Ok(event) => { + if event_tx.send(event).await.is_err() { + return; + } + } + Err(broadcast::error::RecvError::Lagged(skipped)) => { + eprintln!( + "buzz-desktop: archive relay receiver lagged by {skipped} events; stopping sync rather than silently losing archive data" + ); + return; + } + Err(broadcast::error::RecvError::Closed) => return, + } + } + }); + (session, event_rx) + } +} + pub(crate) struct RelaySession { state: Arc>, + requests: Arc>>, + events: broadcast::Sender, wake: mpsc::Sender<()>, cancel: CancellationToken, } +struct PendingRequest { + events: Vec, + complete: oneshot::Sender, String>>, +} + /// Desired set plus the write-time record of what has left it. /// /// One lock covers both because reconcile must read them together: snapshotting @@ -91,6 +164,7 @@ pub(crate) struct RelaySession { #[derive(Default)] struct SessionState { desired: Vec, + transient: Vec, /// Ids whose exact subscription has left `desired` since the last /// reconcile drained this. Written here rather than derived at reconcile /// time because reconcile cannot derive it: wakes coalesce, so a remove @@ -128,6 +202,57 @@ impl SessionState { } impl RelaySession { + pub(crate) fn subscribe(&self) -> broadcast::Receiver { + self.events.subscribe() + } + + /// Fetches one finite page over this session without disturbing persistent + /// feature subscriptions. Request ids are fresh, so CLOSED/backoff history + /// can never leak between pages or into a long-lived subscription. + pub(crate) async fn fetch_events( + &self, + filter: serde_json::Value, + timeout: Duration, + ) -> Result, String> { + let id = format!("native-fetch-{}", uuid::Uuid::new_v4()); + let (complete, result) = oneshot::channel(); + self.requests.lock().await.insert( + id.clone(), + PendingRequest { + events: Vec::new(), + complete, + }, + ); + { + let mut state = self.state.lock().await; + state.transient.push(Subscription { + id: id.clone(), + filter, + }); + } + let _ = self.wake.try_send(()); + + let outcome = tokio::select! { + _ = self.cancel.cancelled() => Err("relay session cancelled".to_string()), + value = tokio::time::timeout(timeout, result) => match value { + Ok(Ok(value)) => value, + Ok(Err(_)) => Err("relay request ended before EOSE".to_string()), + Err(_) => Err("relay request timed out".to_string()), + } + }; + self.finish_request(&id).await; + outcome + } + + async fn finish_request(&self, id: &str) { + self.requests.lock().await.remove(id); + let mut state = self.state.lock().await; + state.transient.retain(|subscription| subscription.id != id); + state.removed.insert(id.to_string()); + drop(state); + let _ = self.wake.try_send(()); + } + /// Replaces the desired subscription set and wakes the loop to reconcile. /// /// Reconciliation is declarative rather than incremental: callers state @@ -169,15 +294,42 @@ impl RelaySession { /// reconnects on drop with exponential backoff and resubscribes the current /// desired set โ€” never a snapshot captured at connect time, so a subscription /// change during an outage is honored by the reconnect that follows. +#[cfg(test)] pub(crate) fn start( relay_url: String, keys: Keys, auth_tag: Option, ) -> (Arc, mpsc::Receiver) { + let session = start_managed(relay_url, keys, auth_tag); + let mut events = session.subscribe(); let (event_tx, event_rx) = mpsc::channel(256); + tauri::async_runtime::spawn(async move { + loop { + match events.recv().await { + Ok(event) => { + if event_tx.send(event).await.is_err() { + return; + } + } + Err(broadcast::error::RecvError::Lagged(skipped)) => { + eprintln!( + "buzz-desktop: native_relay_client: legacy receiver lagged by {skipped} events" + ); + } + Err(broadcast::error::RecvError::Closed) => return, + } + } + }); + (session, event_rx) +} + +fn start_managed(relay_url: String, keys: Keys, auth_tag: Option) -> Arc { + let (events, _) = broadcast::channel(256); let (wake, wake_rx) = mpsc::channel(1); let session = Arc::new(RelaySession { state: Arc::new(Mutex::new(SessionState::default())), + requests: Arc::new(Mutex::new(HashMap::new())), + events, wake, cancel: CancellationToken::new(), }); @@ -188,10 +340,9 @@ pub(crate) fn start( auth_tag, Arc::clone(&session), wake_rx, - event_tx, )); - (session, event_rx) + session } async fn run_session( @@ -200,7 +351,6 @@ async fn run_session( auth_tag: Option, session: Arc, mut wake_rx: mpsc::Receiver<()>, - event_tx: mpsc::Sender, ) { let mut delay = RECONNECT_BASE_DELAY; loop { @@ -215,7 +365,7 @@ async fn run_session( // clean exit โ€” a socket that drops after one event must not // inherit the previous failure's delay. delay = RECONNECT_BASE_DELAY; - run_connection(conn, &session, &mut wake_rx, &event_tx).await; + run_connection(conn, &session, &mut wake_rx).await; } Err(error) => { eprintln!("buzz-desktop: native_relay_client: connect failed: {error}"); @@ -238,7 +388,6 @@ async fn run_connection( mut conn: NostrWsConnection, session: &RelaySession, wake_rx: &mut mpsc::Receiver<()>, - event_tx: &mpsc::Sender, ) { // Subscription ids currently open ON THIS SOCKET. Deliberately local: a new // socket has none, so reconnect resubscribes the full desired set without @@ -321,19 +470,40 @@ async fn run_connection( if !open.contains_key(&subscription_id) { continue; } + let pending = session + .requests + .lock() + .await + .contains_key(&subscription_id); + if pending { + // Reject forged finite-request events before + // retaining them, bounding memory at the transport + // seam. The catalog re-verifies defensively before + // head selection. + if event.verify().is_err() { + continue; + } + if let Some(request) = session + .requests + .lock() + .await + .get_mut(&subscription_id) + { + request.events.push(*event); + } + continue; + } // Delivery proves the subscription is healthy, so any // accumulated backoff for it is stale. Mirrors the JS // port's per-event `closedRetryAttempt = 0`. retries.remove(&subscription_id); - if event_tx - .send(MatchedEvent { subscription_id, event }) - .await - .is_err() - { - // Receiver gone: nobody is consuming this session. - let _ = conn.disconnect().await; - return; - } + // A broadcast session remains healthy when no feature + // currently observes persistent events (catalog fetches + // are fulfilled above through `requests`). + let _ = session.events.send(MatchedEvent { + subscription_id, + event, + }); } Ok(RelayMessage::Closed { subscription_id, message }) => { // The relay dropped it; forget it so a reopen re-sends @@ -348,6 +518,15 @@ async fn run_connection( if open.remove(&subscription_id).is_none() { continue; } + if let Some(request) = session.requests.lock().await.remove(&subscription_id) { + let _ = request.complete.send(Err(format!("relay closed request: {message}"))); + let mut state = session.state.lock().await; + state.transient.retain(|subscription| subscription.id != subscription_id); + state.removed.insert(subscription_id.clone()); + drop(state); + let _ = session.wake.try_send(()); + continue; + } let retry = retries.entry(subscription_id.clone()).or_default(); retry.schedule(&message); eprintln!( @@ -361,6 +540,15 @@ async fn run_connection( // is what keeps an intermittent relay from ratcheting // its way to the 30s ceiling and staying there. let was_open = open.contains_key(&subscription_id); + if let Some(request) = session.requests.lock().await.remove(&subscription_id) { + let _ = request.complete.send(Ok(request.events)); + let mut state = session.state.lock().await; + state.transient.retain(|subscription| subscription.id != subscription_id); + state.removed.insert(subscription_id.clone()); + drop(state); + let _ = session.wake.try_send(()); + continue; + } retries.remove(&subscription_id); // The relay is running a subscription this socket does // not think is open, so the two disagree. EOSE is the @@ -411,7 +599,15 @@ async fn reconcile( let (desired, removed) = { let mut state = session.state.lock().await; let removed = std::mem::take(&mut state.removed); - (state.desired.clone(), removed) + ( + state + .desired + .iter() + .chain(&state.transient) + .cloned() + .collect::>(), + removed, + ) }; // Retry state is only valid while its id has been continuously desired diff --git a/desktop/src-tauri/src/native_relay_client_tests.rs b/desktop/src-tauri/src/native_relay_client_tests.rs index b1b8a632a..25cc69f52 100644 --- a/desktop/src-tauri/src/native_relay_client_tests.rs +++ b/desktop/src-tauri/src/native_relay_client_tests.rs @@ -152,6 +152,83 @@ async fn settle() { tokio::time::sleep(Duration::from_millis(500)).await; } +/// C's acceptance edge: a finite request shares the authenticated real socket +/// with a persistent subscription, completes on wire EOSE, and does not steal +/// later persistent delivery. A fake connection cannot establish any of those +/// transport/lifetime properties. +#[tokio::test] +async fn finite_fetch_multiplexes_with_persistent_delivery_on_a_real_websocket() { + let (relay_url, mut frames, commands) = stub_relay().await; + let (session, mut events) = start(relay_url, Keys::generate(), None); + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the persistent REQ").await, PROBE_ID); + + let fetch = { + let session = Arc::clone(&session); + tokio::spawn(async move { + session + .fetch_events( + serde_json::json!({ "kinds": [buzz_core_pkg::kind::KIND_PERSONA], "limit": 500 }), + Duration::from_secs(10), + ) + .await + }) + }; + let request_id = next_req(&mut frames, "the finite fetch REQ").await; + assert_ne!(request_id, PROBE_ID); + + let relay_keys = Keys::generate(); + let mut forged = EventBuilder::text_note("forged catalog page event") + .sign_with_keys(&relay_keys) + .unwrap(); + forged.content = "tampered after signing".into(); + commands + .send(StubCommand::Event( + request_id.clone(), + serde_json::to_value(forged).unwrap(), + )) + .await + .unwrap(); + let fetched = EventBuilder::text_note("catalog page event") + .sign_with_keys(&relay_keys) + .unwrap(); + commands + .send(StubCommand::Event( + request_id.clone(), + serde_json::to_value(&fetched).unwrap(), + )) + .await + .unwrap(); + commands + .send(StubCommand::Eose(request_id.clone())) + .await + .unwrap(); + + assert_eq!(fetch.await.unwrap().unwrap(), vec![fetched]); + assert_eq!( + next_frame(&mut frames, "finite fetch CLOSE").await, + Frame::Close(request_id) + ); + + let persistent = EventBuilder::text_note("persistent event after fetch") + .sign_with_keys(&relay_keys) + .unwrap(); + commands + .send(StubCommand::Event( + PROBE_ID.into(), + serde_json::to_value(&persistent).unwrap(), + )) + .await + .unwrap(); + let delivered = tokio::time::timeout(Duration::from_secs(10), events.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!(delivered.subscription_id, PROBE_ID); + assert_eq!(*delivered.event, persistent); + session.shutdown(); +} + /// The blocker: a CLOSED with the desired set never changing again must /// still reopen the subscription. /// diff --git a/desktop/src-tauri/src/persona_catalog.rs b/desktop/src-tauri/src/persona_catalog.rs new file mode 100644 index 000000000..5d1717d67 --- /dev/null +++ b/desktop/src-tauri/src/persona_catalog.rs @@ -0,0 +1,296 @@ +//! Native persona-catalog fetch and trust-boundary projection. +//! +//! The renderer owns presentation/linkage to local personas. Relay paging, +//! signature verification, NIP-33 head selection, and untrusted-content parsing +//! stay here so a catalog refresh crosses IPC once instead of once per page and +//! never performs Schnorr verification on the webview thread. + +use std::{collections::HashMap, time::Duration}; + +use buzz_core_pkg::kind::KIND_PERSONA; +use nostr::Event; +use regex::Regex; +use serde::Serialize; +use serde_json::Value; +use std::sync::LazyLock; +use tauri::State; + +use crate::{ + app_state::AppState, managed_agents::validate_agent_definition_text, + native_relay_client::NativeRelayClient, +}; + +const CATALOG_PAGE_SIZE: usize = 500; +const MAX_CATALOG_PAGES: usize = 40; +const PAGE_TIMEOUT: Duration = Duration::from_secs(10); +const MAX_HTTP_AVATAR_LENGTH: usize = 2_048; +const INLINE_SVG_AVATAR_PREFIX: &str = "data:image/svg+xml,"; +const MAX_INLINE_SVG_AVATAR_LENGTH: usize = 8_192; +const MAX_INLINE_RASTER_AVATAR_LENGTH: usize = 256 * 1_024; + +static INLINE_RASTER_AVATAR: LazyLock> = LazyLock::new(|| { + Regex::new(r"^data:image/(?:png|jpeg|gif|webp);base64,([A-Za-z0-9+/]+={0,2})$").ok() +}); + +#[derive(Debug, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct PersonaCatalogPublication { + event_id: String, + owner_pubkey: String, + source_persona_id: String, + created_at: u64, + agent: CatalogAgentProjection, +} + +#[derive(Debug, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +struct CatalogAgentProjection { + display_name: String, + avatar_url: Option, + system_prompt: String, + runtime: Option, + model: Option, + provider: Option, + name_pool: Vec, + respond_to: Option, + parallelism: Option, +} + +/// Fetches the active community's relay-confirmed persona catalog. +/// +/// The command accepts no relay or identity input: both are snapshotted from +/// `AppState`, then checked again before return so an in-flight old-community +/// response cannot populate the new community's query cache. +#[tauri::command] +pub(crate) async fn fetch_persona_catalog( + state: State<'_, AppState>, + relay_client: State<'_, NativeRelayClient>, +) -> Result, String> { + let keys = state.signing_keys()?; + let owner = keys.public_key().to_hex(); + let relay_url = crate::relay::relay_ws_url_with_override(&state); + let session = relay_client.session(relay_url.clone(), keys).await; + let mut by_id = HashMap::new(); + let mut until = None; + + for _ in 0..MAX_CATALOG_PAGES { + let mut filter = serde_json::json!({ + "kinds": [KIND_PERSONA], + "limit": CATALOG_PAGE_SIZE, + }); + if let Some(until) = until { + filter["until"] = serde_json::json!(until); + } + let page = session.fetch_events(filter, PAGE_TIMEOUT).await?; + let page_len = page.len(); + // Schnorr verification is CPU-bound. Keep the complete page off the + // async executor (and therefore off Tauri command scheduling). + let verified = tauri::async_runtime::spawn_blocking(move || { + page.into_iter() + .filter(|event| event.verify().is_ok()) + .collect::>() + }) + .await + .map_err(|error| format!("catalog signature verification failed: {error}"))?; + + let progress = merge_verified_page(&mut by_id, page_len, verified); + match progress { + PageProgress::Done => break, + PageProgress::Next(next_until) => until = Some(next_until), + } + } + + let current_keys = state.signing_keys()?; + if current_keys.public_key().to_hex() != owner + || crate::relay::relay_ws_url_with_override(&state) != relay_url + { + return Err("persona catalog scope changed while fetching".to_string()); + } + + Ok(publications_from_verified_events( + by_id.into_values().collect(), + )) +} + +#[derive(Debug, PartialEq)] +enum PageProgress { + Done, + Next(u64), +} + +fn merge_verified_page( + by_id: &mut HashMap, + wire_page_len: usize, + verified: Vec, +) -> PageProgress { + let size_before = by_id.len(); + let oldest = verified + .iter() + .map(|event| event.created_at.as_secs()) + .min(); + for event in verified { + by_id.insert(event.id.to_hex(), event); + } + + // A short page is the end of the catalog; a page of only repeats means the + // inclusive `until` cursor cannot advance past tied timestamps. + if wire_page_len < CATALOG_PAGE_SIZE || by_id.len() == size_before { + return PageProgress::Done; + } + // A full page of invalid signatures cannot supply a trusted cursor. + oldest.map_or(PageProgress::Done, PageProgress::Next) +} + +fn publications_from_verified_events(mut events: Vec) -> Vec { + events.sort_by(|left, right| { + right + .created_at + .cmp(&left.created_at) + .then_with(|| left.id.cmp(&right.id)) + }); + let mut claimed = std::collections::HashSet::new(); + let mut publications = Vec::new(); + + for event in events { + if event.kind.as_u16() as u32 != KIND_PERSONA { + continue; + } + let Some(source_persona_id) = coordinate_tag(&event, "d") else { + continue; + }; + if source_persona_id.is_empty() { + continue; + } + let owner_pubkey = event.pubkey.to_hex().to_ascii_lowercase(); + let coordinate = (owner_pubkey.clone(), source_persona_id.clone()); + if !claimed.insert(coordinate) { + continue; + } + + // Claim happens before visibility or parsing. A valid newest unshared + // or malformed head is still the NIP-33 head and must not resurrect an + // older shared definition. + if exact_tag(&event, "shared").as_deref() != Some("true") { + continue; + } + let Some(agent) = parse_agent(&event.content) else { + continue; + }; + publications.push(PersonaCatalogPublication { + event_id: event.id.to_hex(), + owner_pubkey, + source_persona_id, + created_at: event.created_at.as_secs(), + agent, + }); + } + publications +} + +fn coordinate_tag(event: &Event, name: &str) -> Option { + let matches = event + .tags + .iter() + .filter_map(|tag| { + let values = tag.as_slice(); + (values.len() >= 2 && values.first().is_some_and(|value| value == name)) + .then(|| values[1].clone()) + }) + .collect::>(); + (matches.len() == 1).then(|| matches[0].clone()) +} + +fn exact_tag(event: &Event, name: &str) -> Option { + let matches = event + .tags + .iter() + .filter_map(|tag| { + let values = tag.as_slice(); + (values.len() == 2 && values.first().is_some_and(|value| value == name)) + .then(|| values[1].clone()) + }) + .collect::>(); + (matches.len() == 1).then(|| matches[0].clone()) +} + +fn parse_agent(content: &str) -> Option { + let value: Value = serde_json::from_str(content).ok()?; + let object = value.as_object()?; + let display_name = object.get("display_name")?.as_str()?.to_string(); + let system_prompt = object + .get("system_prompt") + .and_then(Value::as_str) + .unwrap_or_default() + .to_string(); + validate_agent_definition_text(&display_name, &system_prompt).ok()?; + + let respond_to = match object.get("respond_to").and_then(Value::as_str) { + Some("allowlist") => Some("owner-only".to_string()), + Some(value @ ("owner-only" | "anyone")) => Some(value.to_string()), + _ => None, + }; + let parallelism = object + .get("parallelism") + .and_then(Value::as_u64) + .filter(|value| (1..=32).contains(value)); + let name_pool = object + .get("name_pool") + .and_then(Value::as_array) + .map(|values| { + values + .iter() + .filter_map(Value::as_str) + .map(ToOwned::to_owned) + .collect() + }) + .unwrap_or_default(); + + Some(CatalogAgentProjection { + display_name, + avatar_url: object + .get("avatar_url") + .and_then(Value::as_str) + .filter(|value| safe_avatar(value)) + .map(ToOwned::to_owned), + system_prompt, + runtime: optional_string(object.get("runtime")), + model: optional_string(object.get("model")), + provider: optional_string(object.get("provider")), + name_pool, + respond_to, + parallelism, + }) +} + +fn optional_string(value: Option<&Value>) -> Option { + value + .and_then(Value::as_str) + .filter(|value| !value.trim().is_empty()) + .map(ToOwned::to_owned) +} + +fn safe_avatar(value: &str) -> bool { + if value.starts_with(INLINE_SVG_AVATAR_PREFIX) { + return value.len() <= MAX_INLINE_SVG_AVATAR_LENGTH; + } + if value.len() <= MAX_INLINE_RASTER_AVATAR_LENGTH { + if let Some(captures) = INLINE_RASTER_AVATAR + .as_ref() + .and_then(|pattern| pattern.captures(value)) + { + return captures + .get(1) + .is_some_and(|payload| payload.as_str().len() % 4 == 0); + } + } + value.len() <= MAX_HTTP_AVATAR_LENGTH + && !value.chars().any(char::is_whitespace) + && !value.contains(['(', ')']) + && url::Url::parse(value) + .ok() + .is_some_and(|url| matches!(url.scheme(), "http" | "https")) +} + +#[cfg(test)] +#[path = "persona_catalog_tests.rs"] +mod tests; diff --git a/desktop/src-tauri/src/persona_catalog_tests.rs b/desktop/src-tauri/src/persona_catalog_tests.rs new file mode 100644 index 000000000..0532660b4 --- /dev/null +++ b/desktop/src-tauri/src/persona_catalog_tests.rs @@ -0,0 +1,192 @@ +use super::*; +use nostr::{EventBuilder, Keys, Kind, Tag, Timestamp}; +use serde_json::json; + +fn event(keys: &Keys, created_at: u64, source: &str, shared: bool, content: Value) -> Event { + let mut tags = vec![Tag::parse(["d", source]).unwrap()]; + if shared { + tags.push(Tag::parse(["shared", "true"]).unwrap()); + } + EventBuilder::new(Kind::Custom(KIND_PERSONA as u16), content.to_string()) + .tags(tags) + .custom_created_at(Timestamp::from(created_at)) + .sign_with_keys(keys) + .unwrap() +} + +fn valid_content(name: &str) -> Value { + json!({ + "display_name": name, + "system_prompt": "Review changes.", + "avatar_url": "https://relay.example/avatar.png", + "runtime": " goose ", + "model": "claude", + "provider": null, + "name_pool": ["Reviewer", 7], + "respond_to": "allowlist", + "parallelism": 4 + }) +} + +#[test] +fn paging_uses_oldest_verified_cursor_and_stops_on_ties_or_short_pages() { + let keys = Keys::generate(); + let newest = event(&keys, 9, "newest", true, valid_content("Newest")); + let oldest = event(&keys, 4, "oldest", true, valid_content("Oldest")); + let mut by_id = HashMap::new(); + + assert_eq!( + merge_verified_page( + &mut by_id, + CATALOG_PAGE_SIZE, + vec![newest.clone(), oldest.clone()] + ), + PageProgress::Next(4) + ); + assert_eq!( + merge_verified_page(&mut by_id, CATALOG_PAGE_SIZE, vec![newest, oldest]), + PageProgress::Done + ); + + let short = event(&keys, 1, "short", true, valid_content("Short")); + assert_eq!( + merge_verified_page(&mut by_id, CATALOG_PAGE_SIZE - 1, vec![short]), + PageProgress::Done + ); + assert_eq!( + merge_verified_page(&mut HashMap::new(), CATALOG_PAGE_SIZE, Vec::new()), + PageProgress::Done + ); +} + +#[test] +fn forged_newest_head_is_dropped_before_it_can_claim_the_coordinate() { + let keys = Keys::generate(); + let older = event(&keys, 1, "reviewer", true, valid_content("Older")); + let mut forged = event(&keys, 2, "reviewer", true, valid_content("Forged")); + forged.content = valid_content("Tampered").to_string(); + + let verified = [older.clone(), forged] + .into_iter() + .filter(|candidate| candidate.verify().is_ok()) + .collect(); + let publications = publications_from_verified_events(verified); + assert_eq!(publications.len(), 1); + assert_eq!(publications[0].event_id, older.id.to_hex()); +} + +#[test] +fn valid_newest_head_claims_before_visibility_and_content_parsing() { + let keys = Keys::generate(); + for newest in [ + event(&keys, 2, "reviewer", false, valid_content("Unshared")), + event(&keys, 2, "reviewer", true, json!({})), + ] { + let older = event(&keys, 1, "reviewer", true, valid_content("Older")); + assert!(publications_from_verified_events(vec![older, newest]).is_empty()); + } +} + +#[test] +fn equal_second_heads_use_lowest_event_id_and_authors_are_independent() { + let alice = Keys::generate(); + let bob = Keys::generate(); + let shared = event(&alice, 1, "reviewer", true, valid_content("Shared")); + let unshared = event(&alice, 1, "reviewer", false, valid_content("Hidden")); + let bob_head = event(&bob, 1, "reviewer", true, valid_content("Bob")); + let expected_alice = if shared.id < unshared.id { 1 } else { 0 }; + + let publications = publications_from_verified_events(vec![shared, unshared, bob_head]); + assert_eq!(publications.len(), expected_alice + 1); +} + +#[test] +fn parser_projects_types_and_foreign_allowlists_exactly() { + let projection = parse_agent(&valid_content("Reviewer").to_string()).unwrap(); + assert_eq!(projection.display_name, "Reviewer"); + assert_eq!(projection.runtime.as_deref(), Some(" goose ")); + assert_eq!(projection.provider, None); + assert_eq!(projection.name_pool, vec!["Reviewer"]); + assert_eq!(projection.respond_to.as_deref(), Some("owner-only")); + assert_eq!(projection.parallelism, Some(4)); + + for bad in [0, 33] { + let mut content = valid_content("Reviewer"); + content["parallelism"] = json!(bad); + assert_eq!(parse_agent(&content.to_string()).unwrap().parallelism, None); + } +} + +#[test] +fn parser_rejects_malformed_and_invisible_definition_text() { + for content in [ + "not-json".to_string(), + "[]".to_string(), + json!({"display_name": 7}).to_string(), + valid_content("Review\u{202e}er").to_string(), + ] { + assert!(parse_agent(&content).is_none()); + } + let visible = parse_agent( + &json!({ + "display_name": "Reviewer ๐Ÿ", + "system_prompt": "Review.\n\t||literal markdown||" + }) + .to_string(), + ) + .unwrap(); + assert_eq!(visible.display_name, "Reviewer ๐Ÿ"); +} + +#[test] +fn avatar_allowlist_and_bounds_match_the_renderer_contract() { + assert!(safe_avatar("https://relay.example/avatar.png")); + assert!(!safe_avatar("javascript:alert(1)")); + assert!(safe_avatar("data:image/svg+xml,")); + assert!(!safe_avatar(&format!( + "data:image/svg+xml,{}", + "a".repeat(MAX_INLINE_SVG_AVATAR_LENGTH) + ))); + for mime in ["png", "jpeg", "gif", "webp"] { + assert!(safe_avatar(&format!( + "data:image/{mime};base64,iVBORw0KGgo=" + ))); + } + assert!(!safe_avatar("data:image/bmp;base64,aA==")); + assert!(!safe_avatar("data:image/png;base64,not base64")); +} + +#[test] +fn exact_tags_reject_duplicates_and_extra_fields() { + let keys = Keys::generate(); + let base = event(&keys, 1, "reviewer", true, valid_content("Reviewer")); + assert_eq!(exact_tag(&base, "shared").as_deref(), Some("true")); + + let duplicate = EventBuilder::new( + Kind::Custom(KIND_PERSONA as u16), + valid_content("x").to_string(), + ) + .tags([ + Tag::parse(["d", "reviewer"]).unwrap(), + Tag::parse(["shared", "true"]).unwrap(), + Tag::parse(["shared", "true"]).unwrap(), + ]) + .sign_with_keys(&keys) + .unwrap(); + assert_eq!(exact_tag(&duplicate, "shared"), None); + + // The old renderer accepts an extended d tag (it reads tag[1]) but shared + // is opt-in only for the exact two-field shape. + let extended = EventBuilder::new( + Kind::Custom(KIND_PERSONA as u16), + valid_content("x").to_string(), + ) + .tags([ + Tag::parse(["d", "reviewer", "relay hint"]).unwrap(), + Tag::parse(["shared", "true", "extra"]).unwrap(), + ]) + .sign_with_keys(&keys) + .unwrap(); + assert_eq!(coordinate_tag(&extended, "d").as_deref(), Some("reviewer")); + assert_eq!(exact_tag(&extended, "shared"), None); +} diff --git a/desktop/src/features/agents/lib/personaCatalogRelay.invoke.test.mjs b/desktop/src/features/agents/lib/personaCatalogRelay.invoke.test.mjs new file mode 100644 index 000000000..5ce7bda4c --- /dev/null +++ b/desktop/src/features/agents/lib/personaCatalogRelay.invoke.test.mjs @@ -0,0 +1,26 @@ +import assert from "node:assert/strict"; +import test, { mock } from "node:test"; + +import { fetchPersonaCatalogPublications } from "./personaCatalogRelay.ts"; + +function installTauriInvoke(handler) { + globalThis.window ??= {}; + window.__TAURI_INTERNALS__ = { invoke: handler }; +} + +test("persona catalog wrapper invokes the one native active-scope command", async (t) => { + const prior = globalThis.window; + t.after(() => { + mock.restoreAll(); + globalThis.window = prior; + }); + const expected = [{ eventId: "event-1", ownerPubkey: "alice" }]; + const calls = []; + installTauriInvoke((command, args) => { + calls.push([command, args]); + return Promise.resolve(expected); + }); + + assert.deepEqual(await fetchPersonaCatalogPublications(), expected); + assert.deepEqual(calls, [["fetch_persona_catalog", {}]]); +}); diff --git a/desktop/src/features/agents/lib/personaCatalogRelay.test.mjs b/desktop/src/features/agents/lib/personaCatalogRelay.test.mjs index 5eda8a195..49e7b45eb 100644 --- a/desktop/src/features/agents/lib/personaCatalogRelay.test.mjs +++ b/desktop/src/features/agents/lib/personaCatalogRelay.test.mjs @@ -1,450 +1,32 @@ import assert from "node:assert/strict"; -import test, { mock } from "node:test"; -import { finalizeEvent, getPublicKey } from "nostr-tools/pure"; +import test from "node:test"; -import { relayClient } from "@/shared/api/relayClient"; -import { emojiAvatarDataUrl } from "@/features/profile/ui/ProfileAvatarEditor.utils.ts"; -import { - catalogPersonasFromPublications, - catalogPublicationsFromEvents, - fetchPersonaCatalogPublications, - personaEventIsShared, -} from "./personaCatalogRelay.ts"; +import { catalogPersonasFromPublications } from "./personaCatalogRelay.ts"; -const ALICE_SECRET = new Uint8Array(32); -ALICE_SECRET[31] = 1; -const BOB_SECRET = new Uint8Array(32); -BOB_SECRET[31] = 2; -const ALICE = getPublicKey(ALICE_SECRET); -const BOB = getPublicKey(BOB_SECRET); +const ALICE = "a".repeat(64); +const BOB = "b".repeat(64); -function secretForOwner(owner) { - if (owner === ALICE) return ALICE_SECRET; - if (owner === BOB) return BOB_SECRET; - throw new Error(`No test secret for catalog owner ${owner}`); -} - -function personaEvent({ - createdAt, - id, - owner = ALICE, - sourcePersonaId = "reviewer", - shared = true, - avatarUrl = null, - displayName = "Relay Reviewer", - respondTo = null, - systemPrompt = "Review changes.", - sharedTag, - contentOverride, -}) { - return finalizeEvent( - { - created_at: createdAt, - kind: 30175, - tags: [ - ["d", sourcePersonaId], - ["test-id", id], - ...(shared - ? [sharedTag ?? ["shared", "true"]] - : sharedTag - ? [sharedTag] - : []), - ], - content: - contentOverride ?? - JSON.stringify({ - display_name: displayName, - system_prompt: systemPrompt, - avatar_url: avatarUrl, - runtime: "goose", - model: "claude", - provider: null, - name_pool: ["Reviewer"], - respond_to: respondTo, - respond_to_allowlist: respondTo === "allowlist" ? [BOB] : undefined, - parallelism: 4, - }), +function publication(overrides = {}) { + return { + eventId: "event-1", + ownerPubkey: ALICE, + sourcePersonaId: "reviewer", + createdAt: 1, + agent: { + displayName: "Relay Reviewer", + avatarUrl: null, + systemPrompt: "Review changes.", + runtime: null, + model: null, + provider: null, + namePool: [], + respondTo: null, + parallelism: null, }, - secretForOwner(owner), - ); -} - -test("a shared kind 30175 persona from Alice is discoverable by Bob", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "alice-reviewer" }), - ]); - const personas = catalogPersonasFromPublications(publications, [], BOB); - - assert.equal(personas.length, 1); - assert.equal(personas[0].displayName, "Relay Reviewer"); - assert.equal(personas[0].isActive, false); - assert.equal(personas[0].shared, true); - assert.equal(personas[0].catalogSource.ownerPubkey, ALICE); - assert.equal(personas[0].catalogSource.isOwn, false); -}); - -test("a newer unshared head hides the older shared head", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "shared" }), - personaEvent({ createdAt: 2, id: "unshared", shared: false }), - ]); - - assert.deepEqual(publications, []); -}); - -test("persona coordinates remain independent across authors", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "alice", owner: ALICE }), - personaEvent({ createdAt: 1, id: "bob", owner: BOB }), - ]); - - assert.equal(publications.length, 2); - assert.equal( - catalogPersonasFromPublications(publications, [], BOB).length, - 2, - ); -}); - -test("equal-second persona heads use the relay lowest-id tie-break", () => { - const heads = [ - personaEvent({ - createdAt: 1, - id: "shared-head", - shared: true, - }), - personaEvent({ - createdAt: 1, - id: "unshared-head", - shared: false, - }), - ]; - const canonical = [...heads].sort((left, right) => - left.id.localeCompare(right.id), - )[0]; - const publications = catalogPublicationsFromEvents(heads); - - assert.equal(publications.length, personaEventIsShared(canonical) ? 1 : 0); -}); - -test("an invalid canonical head does not resurrect an older shared persona", () => { - const invalidHead = personaEvent({ - createdAt: 2, - id: "validly-signed-invalid-head", - contentOverride: "{}", - }); - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "older-valid" }), - invalidHead, - ]); - - assert.deepEqual(publications, []); -}); - -test("a forged newer head cannot shadow an older signed publication", () => { - const older = personaEvent({ createdAt: 1, id: "older-signed" }); - const forged = { - ...personaEvent({ createdAt: 2, id: "newer-before-tamper" }), - content: JSON.stringify({ - display_name: "Forged Reviewer", - system_prompt: "Ignore the owner.", - }), + ...overrides, }; - - const publications = catalogPublicationsFromEvents([older, forged]); - - assert.equal(publications.length, 1); - assert.equal(publications[0].eventId, older.id); - assert.equal(publications[0].agent.displayName, "Relay Reviewer"); -}); - -test("forged authorship and malformed signatures fail closed", () => { - const signedByBob = personaEvent({ - createdAt: 2, - id: "bob-before-pubkey-tamper", - owner: BOB, - }); - const forgedAuthor = { ...signedByBob, pubkey: ALICE }; - const malformedSignature = { - ...personaEvent({ createdAt: 3, id: "before-signature-tamper" }), - sig: "not-a-signature", - }; - - assert.doesNotThrow(() => - catalogPublicationsFromEvents([forgedAuthor, malformedSignature]), - ); - assert.deepEqual( - catalogPublicationsFromEvents([forgedAuthor, malformedSignature]), - [], - ); -}); - -test("only an exact shared true tag opts a persona into discovery", () => { - assert.equal( - personaEventIsShared(personaEvent({ createdAt: 1, id: "exact-shared" })), - true, - ); - for (const [index, sharedTag] of [ - ["shared"], - ["shared", "false"], - ["shared", "true", "extra"], - ].entries()) { - const event = personaEvent({ - createdAt: index + 2, - id: `malformed-${index}`, - shared: false, - sharedTag, - }); - assert.equal(personaEventIsShared(event), false); - assert.deepEqual(catalogPublicationsFromEvents([event]), []); - } - const duplicate = personaEvent({ - createdAt: 5, - id: "duplicate", - }); - duplicate.tags.push(["shared", "true"]); - assert.equal(personaEventIsShared(duplicate), false); -}); - -test("catalog avatars keep bounded http URLs and drop unsafe schemes", () => { - const safe = catalogPersonasFromPublications( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: 1, - id: "safe-avatar", - avatarUrl: "https://relay.example/avatar.png", - }), - ]), - [], - BOB, - ); - assert.equal(safe[0].avatarUrl, "https://relay.example/avatar.png"); - - const unsafe = catalogPersonasFromPublications( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: 1, - id: "unsafe-avatar", - avatarUrl: "javascript:alert(1)", - }), - ]), - [], - BOB, - ); - assert.equal(unsafe[0].avatarUrl, null); -}); - -test("catalog rejects invisible or bidirectional formatting characters", () => { - for (const [index, character] of [ - "\u00ad", - "\u034f", - "\u200b", - "\u202e", - "\u2060", - "\u2066", - "\u3164", - "\u{e007f}", - ].entries()) { - assert.deepEqual( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - displayName: `Review${character}er`, - id: `unsafe-name-${index}`, - }), - ]), - [], - ); - assert.deepEqual( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - id: `unsafe-prompt-${index}`, - systemPrompt: `Review code.${character}`, - }), - ]), - [], - ); - } -}); - -test("catalog keeps rendered emoji sequences in names and instructions", () => { - for (const [index, emoji] of [ - "โค๏ธ", - "โ˜•๏ธ", - "๐Ÿ‘ฉโ€๐Ÿ’ป", - "๐Ÿง‘๐Ÿฝโ€๐Ÿ’ป", - "๐Ÿ‘จโ€๐Ÿ‘ฉโ€๐Ÿ‘งโ€๐Ÿ‘ฆ", - "1๏ธโƒฃ", - ].entries()) { - const publications = catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - displayName: `Reviewer ${emoji}`, - id: `rendered-emoji-${index}`, - systemPrompt: `Review changes ${emoji}`, - }), - ]); - - assert.equal(publications.length, 1); - assert.equal(publications[0].agent.displayName, `Reviewer ${emoji}`); - assert.equal(publications[0].agent.systemPrompt, `Review changes ${emoji}`); - } -}); - -test("catalog rejects detached emoji formatting and tag sequences", () => { - const taggedFlag = "๐Ÿด\u{e0067}\u{e0062}\u{e0073}\u{e0063}\u{e0074}\u{e007f}"; - for (const [index, value] of [ - "Review\ufe0fer", - "Review\u200der", - "Review code.\u200d", - taggedFlag, - ].entries()) { - assert.deepEqual( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - displayName: value, - id: `detached-emoji-name-${index}`, - }), - ]), - [], - ); - assert.deepEqual( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - id: `detached-emoji-prompt-${index}`, - systemPrompt: value, - }), - ]), - [], - ); - } -}); - -test("catalog rejects layout controls in display names", () => { - for (const [index, character] of ["\n", "\t"].entries()) { - assert.deepEqual( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: index + 1, - displayName: `Relay${character}Reviewer`, - id: `unsafe-layout-name-${index}`, - }), - ]), - [], - ); - } -}); - -test("catalog keeps visible unicode and literal markdown instructions", () => { - const systemPrompt = - "Review changes.\n\t||This syntax must be shown literally.||"; - const publications = catalogPublicationsFromEvents([ - personaEvent({ - createdAt: 1, - displayName: "Relay Reviewer ๐Ÿ", - id: "visible-unicode", - systemPrompt, - }), - ]); - - assert.equal(publications[0].agent.displayName, "Relay Reviewer ๐Ÿ"); - assert.equal(publications[0].agent.systemPrompt, systemPrompt); -}); - -/** The avatar a catalog entry projects for `avatarUrl`, or null if dropped. */ -function catalogAvatarUrl(avatarUrl) { - const personas = catalogPersonasFromPublications( - catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "avatar-vector", avatarUrl }), - ]), - [], - BOB, - ); - return personas[0].avatarUrl; } -// An emoji avatar is self-contained, so it is the one `data:` avatar that can -// render on another member's machine. Dropping it left shared agents looking -// avatar-less in the catalog. -test("test_percent_encoded_emoji_svg_avatar_survives_the_catalog", () => { - const emojiAvatar = emojiAvatarDataUrl("๐Ÿ", "#FFCC00"); - - assert.equal(catalogAvatarUrl(emojiAvatar), emojiAvatar); -}); - -test("test_base64_svg_avatar_is_rejected", () => { - assert.equal( - catalogAvatarUrl(`data:image/svg+xml;base64,${btoa("")}`), - null, - ); -}); - -test("test_non_svg_data_avatar_is_rejected", () => { - assert.equal(catalogAvatarUrl("data:image/png,%89PNG"), null); -}); - -test("test_legacy_inline_raster_avatar_survives_the_catalog", () => { - for (const mime of ["png", "jpeg", "gif", "webp"]) { - const avatar = `data:image/${mime};base64,iVBORw0KGgo=`; - assert.equal(catalogAvatarUrl(avatar), avatar); - } -}); - -test("test_inline_raster_avatar_rejects_unbounded_or_malformed_payloads", () => { - const prefix = "data:image/png;base64,"; - const payloadLength = 256 * 1_024 - prefix.length; - const validPayloadLength = payloadLength - (payloadLength % 4); - const withinCap = `${prefix}${"a".repeat(validPayloadLength - 2)}==`; - assert.ok(withinCap.length <= 256 * 1_024); - assert.equal(catalogAvatarUrl(withinCap), withinCap); - assert.equal( - catalogAvatarUrl( - `${withinCap}${"a".repeat(256 * 1_024 - withinCap.length + 1)}`, - ), - null, - ); - assert.equal(catalogAvatarUrl("data:image/png;base64,not base64"), null); - assert.equal(catalogAvatarUrl("data:image/bmp;base64,aA=="), null); -}); - -test("test_oversized_inline_svg_avatar_is_rejected", () => { - const withinCap = `data:image/svg+xml,${"a".repeat(8_192 - "data:image/svg+xml,".length)}`; - assert.equal(withinCap.length, 8_192); - assert.equal(catalogAvatarUrl(withinCap), withinCap); - assert.equal(catalogAvatarUrl(`${withinCap}a`), null); -}); - -// Catalog avatars render through `` (ProfileAvatar โ†’ AvatarImage), -// where an SVG document is never scripted, so a script-bearing avatar is -// accepted and inert rather than filtered โ€” the projection must not silently -// start sanitizing markup it does not render. -test("test_script_bearing_inline_svg_avatar_is_accepted_and_rendered_inert", () => { - const scripted = `data:image/svg+xml,${encodeURIComponent( - '', - )}`; - - assert.equal(catalogAvatarUrl(scripted), scripted); -}); - -test("foreign allowlist behavior imports as owner-only", () => { - const personas = catalogPersonasFromPublications( - catalogPublicationsFromEvents([ - personaEvent({ - createdAt: 1, - id: "allowlist", - respondTo: "allowlist", - }), - ]), - [], - BOB, - ); - - assert.equal(personas[0].respondTo, "owner-only"); - assert.deepEqual(personas[0].respondToAllowlist, []); -}); - test("a pending local share does not appear before relay confirmation", () => { const localPersona = { id: "local-reviewer", @@ -501,13 +83,11 @@ function localPersona(overrides = {}) { // stored catalogSource coordinate links the copy back to the publication. test("test_added_foreign_catalog_entry_keeps_publisher_identity_and_local_selection", () => { const publisherAvatar = "https://relay.example/publisher.png"; - const publications = catalogPublicationsFromEvents([ - personaEvent({ - createdAt: 1, - id: "alice-reviewer", - avatarUrl: publisherAvatar, + const publications = [ + publication({ + agent: { ...publication().agent, avatarUrl: publisherAvatar }, }), - ]); + ]; const copy = localPersona({ id: "a-fresh-uuid", displayName: "Locally Renamed Reviewer", @@ -535,9 +115,7 @@ test("test_added_foreign_catalog_entry_keeps_publisher_identity_and_local_select }); test("test_foreign_entry_with_no_local_copy_stays_unselected", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "alice-reviewer" }), - ]); + const publications = [publication()]; // A same-named local persona with no provenance is a different agent. const unrelated = localPersona({ id: "unrelated" }); @@ -554,9 +132,7 @@ test("test_foreign_entry_with_no_local_copy_stays_unselected", () => { // Provenance is per-owner: the same d-tag under a different publisher is a // different agent, so a copy of Alice's must not mask Bob's entry. test("test_catalog_source_match_is_scoped_to_the_publishing_owner", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "bob-reviewer", owner: BOB }), - ]); + const publications = [publication({ ownerPubkey: BOB })]; const copyOfAlices = localPersona({ id: "copy-of-alices", catalogSource: { ownerPubkey: ALICE, personaId: "reviewer" }, @@ -573,9 +149,7 @@ test("test_catalog_source_match_is_scoped_to_the_publishing_owner", () => { }); test("test_own_publication_still_resolves_by_local_id", () => { - const publications = catalogPublicationsFromEvents([ - personaEvent({ createdAt: 1, id: "alice-reviewer" }), - ]); + const publications = [publication()]; const own = localPersona({ id: "reviewer", shared: true }); const personas = catalogPersonasFromPublications(publications, [own], ALICE); @@ -583,126 +157,3 @@ test("test_own_publication_still_resolves_by_local_id", () => { assert.equal(personas[0].id, "reviewer"); assert.equal(personas[0].catalogSource.isOwn, true); }); - -function pageOfEvents(count, startId, createdAt) { - return Array.from({ length: count }, (_, index) => - personaEvent({ - createdAt: typeof createdAt === "function" ? createdAt(index) : createdAt, - id: `event-${startId + index}`, - sourcePersonaId: `persona-${startId + index}`, - }), - ); -} - -function stubPagedRelay(pages) { - const filters = []; - mock.method(relayClient, "fetchEvents", (filter) => { - filters.push(filter); - return Promise.resolve(pages[filters.length - 1] ?? []); - }); - return filters; -} - -// A single limit-capped fetch drops every entry past the relay's clamp, making -// those agents undiscoverable. The walk must keep going while pages come back -// full, and must carry an `until` cursor derived from the oldest event seen. -test("test_full_page_is_followed_by_a_cursored_request_for_older_events", async (t) => { - t.after(() => mock.restoreAll()); - const filters = stubPagedRelay([ - pageOfEvents(500, 0, (index) => 10_000 - index), - pageOfEvents(3, 500, 9_000), - ]); - - const publications = await fetchPersonaCatalogPublications(); - - assert.equal(filters.length, 2, "a full page must be followed by another"); - assert.equal(filters[0].until, undefined, "the first page has no cursor"); - assert.equal( - filters[1].until, - 10_000 - 499, - "the cursor must be the oldest created_at from the previous page", - ); - assert.equal( - publications.length, - 503, - "entries past the first page must still be discoverable", - ); -}); - -test("test_invalid_events_cannot_control_the_catalog_cursor", async (t) => { - t.after(() => mock.restoreAll()); - const validEvents = pageOfEvents(499, 0, (index) => 10_000 - index); - const invalidOldest = { - ...personaEvent({ - createdAt: 1, - id: "invalid-oldest-cursor", - sourcePersonaId: "invalid-oldest-cursor", - }), - sig: "not-a-signature", - }; - const filters = stubPagedRelay([ - [...validEvents, invalidOldest], - pageOfEvents(1, 500, 9_000), - ]); - - const publications = await fetchPersonaCatalogPublications(); - - assert.equal(filters.length, 2); - assert.equal( - filters[1].until, - 10_000 - 498, - "the cursor must be derived only from verified events", - ); - assert.equal(publications.length, 500); - assert.equal( - publications.some( - (publication) => publication.sourcePersonaId === "invalid-oldest-cursor", - ), - false, - ); -}); - -test("test_short_first_page_does_not_issue_a_second_request", async (t) => { - t.after(() => mock.restoreAll()); - const filters = stubPagedRelay([pageOfEvents(2, 0, 10_000)]); - - const publications = await fetchPersonaCatalogPublications(); - - assert.equal(filters.length, 1); - assert.equal(publications.length, 2); -}); - -// `until` is inclusive on the relay, so consecutive pages overlap on the -// boundary timestamp. Without id dedupe the repeats would be counted twice. -test("test_overlapping_pages_are_deduped_by_event_id", async (t) => { - t.after(() => mock.restoreAll()); - const firstPage = pageOfEvents(500, 0, (index) => 10_000 - index); - const secondPage = [ - // The boundary event repeats because `until` includes its timestamp. - firstPage[firstPage.length - 1], - ...pageOfEvents(2, 500, 9_000), - ]; - stubPagedRelay([firstPage, secondPage]); - - const publications = await fetchPersonaCatalogPublications(); - - assert.equal(publications.length, 502, "the repeated event must count once"); -}); - -// The stop-on-no-progress guard: a full page whose events all share one -// created_at cannot advance the cursor, so paging must terminate instead of -// re-requesting the same page forever. -test("test_full_page_of_tied_timestamps_terminates_the_walk", async (t) => { - t.after(() => mock.restoreAll()); - const tiedPage = pageOfEvents(500, 0, 10_000); - const filters = stubPagedRelay([tiedPage, tiedPage, tiedPage, tiedPage]); - - const publications = await fetchPersonaCatalogPublications(); - - assert.equal( - filters.length, - 2, - "the walk must stop once a page contributes nothing new", - ); - assert.equal(publications.length, 500); -}); diff --git a/desktop/src/features/agents/lib/personaCatalogRelay.ts b/desktop/src/features/agents/lib/personaCatalogRelay.ts index 3f7cd9fdd..63a357e44 100644 --- a/desktop/src/features/agents/lib/personaCatalogRelay.ts +++ b/desktop/src/features/agents/lib/personaCatalogRelay.ts @@ -1,12 +1,9 @@ -import { relayClient } from "@/shared/api/relayClient"; import type { AgentPersona, CatalogSourceCoordinate, - RelayEvent, RespondToMode, } from "@/shared/api/types"; -import { KIND_PERSONA } from "@/shared/constants/kinds"; -import { verifyEvent } from "nostr-tools/pure"; +import { invokeTauri } from "@/shared/api/tauri"; export type CatalogPersonaShareLevel = "not-shared" | "none"; @@ -41,385 +38,19 @@ export type CatalogPersona = AgentPersona & { type JsonObject = Record; -const MAX_AGENT_DISPLAY_NAME_CHARACTERS = 128; -const MAX_AGENT_SYSTEM_PROMPT_BYTES = 64 * 1_024; -const EMOJI_VARIATION_SELECTOR = 0xfe0f; -const ZERO_WIDTH_JOINER = 0x200d; -const EXTENDED_PICTOGRAPHIC_RE = /^\p{Extended_Pictographic}$/u; - -function isProhibitedAgentTextCharacter( - characters: readonly string[], - index: number, - allowLayoutControls: boolean, -): boolean { - const character = characters[index]; - if (character === undefined) return false; - const codePoint = character.codePointAt(0); - if (codePoint === undefined) return false; - - const isControl = - codePoint <= 0x1f || (codePoint >= 0x7f && codePoint <= 0x9f); - const isAllowedLayoutControl = - allowLayoutControls && (codePoint === 0x09 || codePoint === 0x0a); - if (isControl && !isAllowedLayoutControl) return true; - if (isAllowedEmojiFormatCharacter(characters, index)) return false; - - return ( - codePoint === 0x00ad || - codePoint === 0x034f || - codePoint === 0x061c || - (codePoint >= 0x115f && codePoint <= 0x1160) || - (codePoint >= 0x17b4 && codePoint <= 0x17b5) || - (codePoint >= 0x180b && codePoint <= 0x180f) || - (codePoint >= 0x200b && codePoint <= 0x200f) || - (codePoint >= 0x202a && codePoint <= 0x202e) || - (codePoint >= 0x2060 && codePoint <= 0x206f) || - codePoint === 0x3164 || - (codePoint >= 0xfe00 && codePoint <= 0xfe0f) || - codePoint === 0xfeff || - codePoint === 0xffa0 || - (codePoint >= 0xfff0 && codePoint <= 0xfff8) || - (codePoint >= 0x1bca0 && codePoint <= 0x1bca3) || - (codePoint >= 0x1d173 && codePoint <= 0x1d17a) || - (codePoint >= 0xe0000 && codePoint <= 0xe0fff) - ); -} - -function isAllowedEmojiFormatCharacter( - characters: readonly string[], - index: number, -): boolean { - const codePoint = characters[index]?.codePointAt(0); - if (codePoint === EMOJI_VARIATION_SELECTOR) { - const previous = characters[index - 1]; - return previous !== undefined && isEmojiVariationBase(previous); - } - if (codePoint !== ZERO_WIDTH_JOINER) return false; - - const next = characters[index + 1]; - return ( - hasPrecedingEmojiBase(characters, index) && - next !== undefined && - EXTENDED_PICTOGRAPHIC_RE.test(next) - ); -} - -function hasPrecedingEmojiBase( - characters: readonly string[], - index: number, -): boolean { - for (let previous = index - 1; previous >= 0; previous -= 1) { - const character = characters[previous]; - const codePoint = character?.codePointAt(0); - if ( - codePoint === EMOJI_VARIATION_SELECTOR || - (codePoint !== undefined && codePoint >= 0x1f3fb && codePoint <= 0x1f3ff) - ) { - continue; - } - return character !== undefined && EXTENDED_PICTOGRAPHIC_RE.test(character); - } - return false; -} - -function isEmojiVariationBase(character: string): boolean { - return ( - /^[#*0-9]$/u.test(character) || EXTENDED_PICTOGRAPHIC_RE.test(character) - ); -} - -function isSafeAgentDefinitionText( - displayName: string, - systemPrompt: string, -): boolean { - const displayNameCharacters = [...displayName]; - const systemPromptCharacters = [...systemPrompt]; - return ( - displayName.trim().length > 0 && - displayNameCharacters.length <= MAX_AGENT_DISPLAY_NAME_CHARACTERS && - new TextEncoder().encode(systemPrompt).length <= - MAX_AGENT_SYSTEM_PROMPT_BYTES && - !displayNameCharacters.some((_character, index) => - isProhibitedAgentTextCharacter(displayNameCharacters, index, false), - ) && - !systemPromptCharacters.some((_character, index) => - isProhibitedAgentTextCharacter(systemPromptCharacters, index, true), - ) - ); -} - -function eventHasValidSignature(event: RelayEvent): boolean { - try { - // Verify a fresh wire-shaped value. nostr-tools memoizes successful checks - // on event objects; relay input must never inherit a stale verification - // marker from an object that was subsequently mutated. - return verifyEvent({ - id: event.id, - pubkey: event.pubkey, - created_at: event.created_at, - kind: event.kind, - tags: event.tags, - content: event.content, - sig: event.sig, - }); - } catch { - return false; - } -} - function isObject(value: unknown): value is JsonObject { return typeof value === "object" && value !== null && !Array.isArray(value); } -function extractTag(event: RelayEvent, name: string): string | null { - const matches = event.tags.filter( - (tag) => tag.length >= 2 && tag[0] === name && typeof tag[1] === "string", - ); - return matches.length === 1 ? (matches[0]?.[1] ?? null) : null; -} - -export function personaEventIsShared(event: RelayEvent): boolean { - const sharedTags = event.tags.filter((tag) => tag[0] === "shared"); - return ( - sharedTags.length === 1 && - sharedTags[0]?.length === 2 && - sharedTags[0]?.[1] === "true" - ); -} - -function isSafeHttpUrl(value: unknown): value is string { - if ( - typeof value !== "string" || - value.length === 0 || - value.length > 2_048 || - /[\s()]/u.test(value) - ) { - return false; - } - try { - const parsed = new URL(value); - return parsed.protocol === "https:" || parsed.protocol === "http:"; - } catch { - return false; - } -} - /** - * Emoji avatars are the one `data:` avatar a catalog entry keeps. - * - * They persist as inline, percent-encoded SVG (`emojiAvatarDataUrl` in - * `ProfileAvatarEditor.utils.ts`), so they are self-contained and render on - * any member's machine โ€” unlike a bundled runtime-default avatar, whose local - * asset path means nothing to another install. The accepted shape is exactly - * that prefix: the trailing comma is what rejects `;base64` payloads, and - * every other `data:` MIME stays rejected. Catalog avatars render through - * `` (`ProfileAvatar` โ†’ `AvatarImage`), where SVG script never - * executes, so bounding the length is the remaining concern โ€” 8 KiB is an - * order of magnitude above the ~700 characters an emoji avatar encodes to. + * Fetch the active community catalog through the shared native relay session. + * Relay scoping, paging, signature verification, and head selection are native; + * this boundary intentionally accepts no caller-supplied relay or identity. */ -const INLINE_SVG_AVATAR_PREFIX = "data:image/svg+xml,"; -const MAX_INLINE_SVG_AVATAR_LENGTH = 8_192; - -/** - * Shared persona heads can carry an uploaded avatar as an inline raster. Keep - * those self-contained images renderable without accepting arbitrary `data:` - * URLs: only the raster MIME types browsers decode in ``, strict base64 - * shape, and a bound no larger than the relay's event-content ceiling. - */ -const MAX_INLINE_RASTER_AVATAR_LENGTH = 256 * 1_024; -const INLINE_RASTER_AVATAR_RE = - /^data:image\/(?:png|jpeg|gif|webp);base64,([A-Za-z0-9+/]+={0,2})$/u; - -function isInlineSvgAvatar(value: unknown): value is string { - return ( - typeof value === "string" && - value.startsWith(INLINE_SVG_AVATAR_PREFIX) && - value.length <= MAX_INLINE_SVG_AVATAR_LENGTH - ); -} - -function isInlineRasterAvatar(value: unknown): value is string { - if ( - typeof value !== "string" || - value.length > MAX_INLINE_RASTER_AVATAR_LENGTH - ) { - return false; - } - const match = INLINE_RASTER_AVATAR_RE.exec(value); - return match !== null && (match[1]?.length ?? 0) % 4 === 0; -} - -function optionalString(value: unknown): string | null { - return typeof value === "string" && value.trim().length > 0 ? value : null; -} - -function parsePersonaContent(event: RelayEvent): CatalogAgentProjection | null { - let parsed: unknown; - try { - parsed = JSON.parse(event.content); - } catch { - return null; - } - if (!isObject(parsed)) return null; - - const displayName = parsed.display_name; - const systemPrompt = - typeof parsed.system_prompt === "string" ? parsed.system_prompt : ""; - if ( - typeof displayName !== "string" || - !isSafeAgentDefinitionText(displayName, systemPrompt) - ) { - return null; - } - - const avatarUrl = - isSafeHttpUrl(parsed.avatar_url) || - isInlineSvgAvatar(parsed.avatar_url) || - isInlineRasterAvatar(parsed.avatar_url) - ? parsed.avatar_url - : null; - const namePool = Array.isArray(parsed.name_pool) - ? parsed.name_pool.filter( - (candidate): candidate is string => typeof candidate === "string", - ) - : []; - const respondTo = - parsed.respond_to === "allowlist" - ? "owner-only" - : parsed.respond_to === "owner-only" || parsed.respond_to === "anyone" - ? parsed.respond_to - : null; - const parallelism = - typeof parsed.parallelism === "number" && - Number.isInteger(parsed.parallelism) && - parsed.parallelism >= 1 && - parsed.parallelism <= 32 - ? parsed.parallelism - : null; - - return { - displayName, - avatarUrl, - systemPrompt, - runtime: optionalString(parsed.runtime), - model: optionalString(parsed.model), - provider: optionalString(parsed.provider), - namePool, - respondTo, - parallelism, - }; -} - -/** - * Collapse relay results to the canonical NIP-33 head for each persona - * coordinate, then keep only exact `["shared", "true"]` heads. - * - * The relay normally returns one replaceable head. The client-side collapse is - * defense in depth for older relays and fixtures, and deliberately claims the - * coordinate before parsing so an invalid or unshared newest head cannot - * resurrect an older shared definition. - */ -export function catalogPublicationsFromEvents( - events: readonly RelayEvent[], -): PersonaCatalogPublication[] { - return catalogPublicationsFromVerifiedEvents( - events.filter(eventHasValidSignature), - ); -} - -function catalogPublicationsFromVerifiedEvents( - events: readonly RelayEvent[], -): PersonaCatalogPublication[] { - const sorted = [...events].sort( - (left, right) => - right.created_at - left.created_at || left.id.localeCompare(right.id), - ); - const seenCoordinates = new Set(); - const publications: PersonaCatalogPublication[] = []; - - for (const event of sorted) { - if (event.kind !== KIND_PERSONA) continue; - const sourcePersonaId = extractTag(event, "d"); - if (!sourcePersonaId) continue; - const ownerPubkey = event.pubkey.toLowerCase(); - const coordinate = `${ownerPubkey}:${sourcePersonaId}`; - if (seenCoordinates.has(coordinate)) continue; - seenCoordinates.add(coordinate); - - if (!personaEventIsShared(event)) continue; - const agent = parsePersonaContent(event); - if (!agent) continue; - publications.push({ - eventId: event.id, - ownerPubkey, - sourcePersonaId, - createdAt: event.created_at, - agent, - }); - } - - return publications; -} - -/** - * Events per catalog page. - * - * Kept well under the relay's 1,000-row `query_events` clamp so a page that - * comes back full is a reliable "there may be more" signal rather than a - * silently truncated result. - */ -const CATALOG_PAGE_SIZE = 500; - -/** - * Hard bound on pages walked, so a relay that keeps returning full pages can - * never spin this forever. - */ -const MAX_CATALOG_PAGES = 40; - -/** - * Read every shared persona event, page by page. - * - * A single `limit`-capped fetch silently truncates once a community publishes - * more agents than the relay's clamp, and the entries that fall off are simply - * undiscoverable. Paging walks backwards through `created_at` using the only - * cursor a WS `REQ` filter carries โ€” `until` โ€” which the relay treats as - * *inclusive*, so consecutive pages overlap on tied timestamps. Two things - * follow, and both are load-bearing: - * - * - dedupe by event id, because the boundary events repeat; and - * - stop when a page contributes nothing new, because a page whose events all - * share one `created_at` would otherwise be requested forever. - */ -export async function fetchPersonaCatalogPublications(): Promise< +export function fetchPersonaCatalogPublications(): Promise< PersonaCatalogPublication[] > { - const byId = new Map(); - let until: number | undefined; - - for (let page = 0; page < MAX_CATALOG_PAGES; page += 1) { - const events = await relayClient.fetchEvents({ - kinds: [KIND_PERSONA], - limit: CATALOG_PAGE_SIZE, - ...(until === undefined ? {} : { until }), - }); - - const sizeBefore = byId.size; - let oldestCreatedAt = Number.POSITIVE_INFINITY; - for (const event of events) { - if (!eventHasValidSignature(event)) continue; - byId.set(event.id, event); - oldestCreatedAt = Math.min(oldestCreatedAt, event.created_at); - } - - // A short page is the end of the catalog; a page of only-repeats means the - // cursor cannot advance past a run of tied timestamps. - if (events.length < CATALOG_PAGE_SIZE || byId.size === sizeBefore) { - break; - } - until = oldestCreatedAt; - } - - return catalogPublicationsFromVerifiedEvents([...byId.values()]); + return invokeTauri("fetch_persona_catalog"); } function publicationToPersona(