perf(desktop): move persona catalog fetching into Rust

Fetch persona catalog pages through the shared native relay session, verify relay events outside the renderer, preserve NIP-33 head and parser trust semantics, and return one projected DTO across IPC. Keep only catalog-to-local-persona linkage in TypeScript.

Finite requests use fresh subscription ids alongside archive subscriptions on the same authenticated socket. Real-WebSocket coverage proves request EVENT/EOSE/CLOSE flow, forged-event rejection, and continued persistent delivery.

Co-authored-by: Tyler Longwell <tlongwell@squareup.com>
Signed-off-by: Tyler Longwell <tlongwell@squareup.com>
This commit is contained in:
Wren
2026-08-15 23:46:57 -04:00
co-authored by Tyler Longwell
parent 8b047a35e0
commit 7396e442b0
9 changed files with 844 additions and 971 deletions
+4 -3
View File
@@ -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(())
}
+3
View File
@@ -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,
+212 -16
View File
@@ -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<Event>,
@@ -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<Option<ManagedSession>>,
}
struct ManagedSession {
scope: (String, String),
session: Arc<RelaySession>,
}
impl NativeRelayClient {
async fn ensure_session(&self, relay_url: String, keys: Keys) -> Arc<RelaySession> {
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<RelaySession> {
self.ensure_session(relay_url, keys).await
}
pub(crate) async fn archive_session(
&self,
relay_url: String,
keys: Keys,
) -> (Arc<RelaySession>, mpsc::Receiver<MatchedEvent>) {
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<Mutex<SessionState>>,
requests: Arc<Mutex<HashMap<String, PendingRequest>>>,
events: broadcast::Sender<MatchedEvent>,
wake: mpsc::Sender<()>,
cancel: CancellationToken,
}
struct PendingRequest {
events: Vec<Event>,
complete: oneshot::Sender<Result<Vec<Event>, 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<Subscription>,
transient: Vec<Subscription>,
/// 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<MatchedEvent> {
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<Vec<Event>, 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<nostr::Tag>,
) -> (Arc<RelaySession>, mpsc::Receiver<MatchedEvent>) {
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<nostr::Tag>) -> Arc<RelaySession> {
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<nostr::Tag>,
session: Arc<RelaySession>,
mut wake_rx: mpsc::Receiver<()>,
event_tx: mpsc::Sender<MatchedEvent>,
) {
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<MatchedEvent>,
) {
// 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::<Vec<_>>(),
removed,
)
};
// Retry state is only valid while its id has been continuously desired
@@ -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.
///
+296
View File
@@ -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<Option<Regex>> = 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<String>,
system_prompt: String,
runtime: Option<String>,
model: Option<String>,
provider: Option<String>,
name_pool: Vec<String>,
respond_to: Option<String>,
parallelism: Option<u64>,
}
/// 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<Vec<PersonaCatalogPublication>, 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::<Vec<_>>()
})
.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<String, Event>,
wire_page_len: usize,
verified: Vec<Event>,
) -> 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<Event>) -> Vec<PersonaCatalogPublication> {
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<String> {
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::<Vec<_>>();
(matches.len() == 1).then(|| matches[0].clone())
}
fn exact_tag(event: &Event, name: &str) -> Option<String> {
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::<Vec<_>>();
(matches.len() == 1).then(|| matches[0].clone())
}
fn parse_agent(content: &str) -> Option<CatalogAgentProjection> {
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<String> {
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;
@@ -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,<svg/>"));
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);
}
@@ -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", {}]]);
});
@@ -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("<svg/>")}`),
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 `<img src>` (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(
'<svg xmlns="http://www.w3.org/2000/svg"><script>alert(1)</script></svg>',
)}`;
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);
});
@@ -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<string, unknown>;
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
* `<img src>` (`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 `<img>`, 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<string>();
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<string, RelayEvent>();
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<PersonaCatalogPublication[]>("fetch_persona_catalog");
}
function publicationToPersona(