mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(serverless): multi-relay support — publish-to-all, read-deduped
Serverless workspaces now take a comma-separated relay list (seeded with the 5 default public relays) for redundancy: - ws_relay query/submit/publish fan out to all relays concurrently; reads merge+dedup by event id, writes succeed if any relay accepts - relay_ws_urls_with_override parses the comma list; live WS uses first relay - setup UIs (Welcome + AddWorkspace) have multi-select relay chips - integration test proves publish-to-2 / read-deduped-from-2
This commit is contained in:
@@ -148,8 +148,9 @@ async fn serverless_create_join_roundtrip() {
|
||||
.expect("publish 39000");
|
||||
eprintln!("39000 publish: accepted={} msg={}", r1.accepted, r1.message);
|
||||
|
||||
let members = events::build_channel_members_serverless(&channel_id, &[creator_pk.clone()])
|
||||
.expect("build members");
|
||||
let members =
|
||||
events::build_channel_members_serverless(&channel_id, std::slice::from_ref(&creator_pk))
|
||||
.expect("build members");
|
||||
let r2 = submit_event(members, &creator_state)
|
||||
.await
|
||||
.expect("publish 39002");
|
||||
@@ -218,9 +219,14 @@ async fn serverless_create_join_roundtrip() {
|
||||
joiner_state.serverless.store(true, Ordering::Relaxed);
|
||||
|
||||
// This is exactly what join_channel does in serverless mode.
|
||||
serverless_set_members(&joiner_state, &channel_id, &[joiner_pk.clone()], &[])
|
||||
.await
|
||||
.expect("join (set members)");
|
||||
serverless_set_members(
|
||||
&joiner_state,
|
||||
&channel_id,
|
||||
std::slice::from_ref(&joiner_pk),
|
||||
&[],
|
||||
)
|
||||
.await
|
||||
.expect("join (set members)");
|
||||
eprintln!("joiner published updated membership");
|
||||
|
||||
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
|
||||
@@ -263,3 +269,71 @@ async fn serverless_create_join_roundtrip() {
|
||||
|
||||
eprintln!("✅ roundtrip OK: create → join → both members visible");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "network: hits multiple public relays"]
|
||||
async fn serverless_multi_relay_fanout() {
|
||||
use crate::app_state::build_app_state;
|
||||
use crate::relay::{query_relay, submit_event};
|
||||
use std::sync::atomic::Ordering;
|
||||
|
||||
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
|
||||
|
||||
let keys = nostr::Keys::generate();
|
||||
let pk = keys.public_key().to_hex();
|
||||
|
||||
// Two relays. Publish fans out to both; reads merge+dedup from both.
|
||||
let relays = "wss://relay.damus.io,wss://nos.lol";
|
||||
|
||||
let state = build_app_state();
|
||||
*state.keys.lock().unwrap() = keys.clone();
|
||||
*state.relay_url_override.lock().unwrap() = Some(relays.to_string());
|
||||
state.serverless.store(true, Ordering::Relaxed);
|
||||
|
||||
let channel_id = uuid::Uuid::new_v4().to_string();
|
||||
let name = format!("multi-{}", &channel_id[..8]);
|
||||
|
||||
let meta =
|
||||
events::build_channel_metadata_serverless(&channel_id, &name, "open", "stream", None, &[])
|
||||
.unwrap();
|
||||
let r = submit_event(meta, &state)
|
||||
.await
|
||||
.expect("publish to all relays");
|
||||
eprintln!(
|
||||
"multi-relay publish: accepted={} msg={}",
|
||||
r.accepted, r.message
|
||||
);
|
||||
assert!(r.accepted);
|
||||
|
||||
let members =
|
||||
events::build_channel_members_serverless(&channel_id, std::slice::from_ref(&pk)).unwrap();
|
||||
submit_event(members, &state)
|
||||
.await
|
||||
.expect("publish members");
|
||||
|
||||
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
|
||||
|
||||
// Read across both relays — should find the channel, deduped to one event.
|
||||
let metas = query_relay(
|
||||
&state,
|
||||
&[serde_json::json!({"kinds":[39000],"#d":[channel_id],"limit":10})],
|
||||
)
|
||||
.await
|
||||
.expect("query metadata");
|
||||
eprintln!("39000 (deduped across 2 relays): {} event(s)", metas.len());
|
||||
assert_eq!(
|
||||
metas.len(),
|
||||
1,
|
||||
"expected exactly one deduped metadata event"
|
||||
);
|
||||
|
||||
let mems = serverless_current_members(&state, &channel_id)
|
||||
.await
|
||||
.expect("read members");
|
||||
assert!(
|
||||
mems.contains(&pk.to_ascii_lowercase()),
|
||||
"membership not found across relays"
|
||||
);
|
||||
|
||||
eprintln!("✅ multi-relay fanout OK: published to 2, read+deduped from 2");
|
||||
}
|
||||
|
||||
@@ -48,7 +48,14 @@ pub fn is_shared_identity() -> bool {
|
||||
|
||||
#[tauri::command]
|
||||
pub fn get_relay_ws_url(state: State<'_, AppState>) -> String {
|
||||
// Serverless workspaces may carry a comma-separated relay list for
|
||||
// redundancy. The live WebSocket connects to the first (primary) relay;
|
||||
// one-shot reads/writes fan out to all of them (see ws_relay).
|
||||
relay_ws_url_with_override(&state)
|
||||
.split(',')
|
||||
.next()
|
||||
.map(|s| s.trim().to_string())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
|
||||
@@ -962,7 +962,7 @@ mod tests {
|
||||
// UI stuck on "join to participate". This guards that fix.
|
||||
let keys = Keys::generate();
|
||||
let me = keys.public_key().to_hex();
|
||||
let event = build_channel_members_serverless("chan-self", &[me.clone()])
|
||||
let event = build_channel_members_serverless("chan-self", std::slice::from_ref(&me))
|
||||
.unwrap()
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
@@ -984,7 +984,7 @@ mod tests {
|
||||
"private",
|
||||
"dm",
|
||||
None,
|
||||
&[me.clone()],
|
||||
std::slice::from_ref(&me),
|
||||
)
|
||||
.unwrap()
|
||||
.sign_with_keys(&keys)
|
||||
|
||||
@@ -40,6 +40,23 @@ pub fn relay_ws_url_with_override(state: &AppState) -> String {
|
||||
workspace_relay_override(state).unwrap_or_else(relay_ws_url)
|
||||
}
|
||||
|
||||
/// Serverless mode: the relay override may be a comma-separated list of relays
|
||||
/// for redundancy (publish-to-all, read-from-all-deduped). Returns the parsed
|
||||
/// list, falling back to the single resolved URL.
|
||||
pub fn relay_ws_urls_with_override(state: &AppState) -> Vec<String> {
|
||||
let raw = relay_ws_url_with_override(state);
|
||||
let list: Vec<String> = raw
|
||||
.split(',')
|
||||
.map(|s| s.trim().to_string())
|
||||
.filter(|s| !s.is_empty())
|
||||
.collect();
|
||||
if list.is_empty() {
|
||||
vec![raw]
|
||||
} else {
|
||||
list
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns the relay HTTP API base URL, checking the workspace override first.
|
||||
/// Precedence: workspace override > env vars > build-time vars > default.
|
||||
pub fn relay_api_base_url_with_override(state: &AppState) -> String {
|
||||
@@ -153,8 +170,8 @@ pub async fn query_relay(
|
||||
) -> Result<Vec<nostr::Event>, String> {
|
||||
// Serverless mode: no HTTP bridge. Query the generic relay over WS.
|
||||
if state.is_serverless() {
|
||||
let relay_url = relay_ws_url_with_override(state);
|
||||
return crate::ws_relay::query_relay_ws(state, &relay_url, filters).await;
|
||||
let relay_urls = relay_ws_urls_with_override(state);
|
||||
return crate::ws_relay::query_relay_ws(state, &relay_urls, filters).await;
|
||||
}
|
||||
|
||||
let url = format!("{}/query", relay_api_base_url_with_override(state));
|
||||
@@ -262,7 +279,12 @@ pub async fn sync_managed_agent_profile(
|
||||
// Serverless mode: publish the agent's profile over plain WS (no HTTP
|
||||
// bridge, no NIP-98). Signed by the agent's keys.
|
||||
if state.is_serverless() {
|
||||
return crate::ws_relay::publish_signed_event_ws(&event, agent_keys, relay_url)
|
||||
let relay_urls: Vec<String> = relay_url
|
||||
.split(',')
|
||||
.map(|s| s.trim().to_string())
|
||||
.filter(|s| !s.is_empty())
|
||||
.collect();
|
||||
return crate::ws_relay::publish_signed_event_ws(&event, agent_keys, &relay_urls)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
format!("Created the agent, but could not sync its profile metadata: {e}")
|
||||
@@ -317,8 +339,8 @@ pub async fn submit_event(
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
// Serverless mode: no HTTP bridge. Publish to the generic relay over WS.
|
||||
if state.is_serverless() {
|
||||
let relay_url = relay_ws_url_with_override(state);
|
||||
return crate::ws_relay::submit_event_ws(builder, state, &relay_url).await;
|
||||
let relay_urls = relay_ws_urls_with_override(state);
|
||||
return crate::ws_relay::submit_event_ws(builder, state, &relay_urls).await;
|
||||
}
|
||||
|
||||
// All synchronous work (signing) must complete before any .await
|
||||
|
||||
@@ -29,7 +29,40 @@ const PUBLISH_TIMEOUT: Duration = Duration::from_secs(10);
|
||||
/// Execute one or more filters as a single `REQ` and collect matching events
|
||||
/// until the relay sends `EOSE`. Mirrors `relay::query_relay` but over a plain
|
||||
/// WebSocket against a generic relay.
|
||||
/// Query a set of relays concurrently and merge results, deduplicating events
|
||||
/// by id. Succeeds if any relay responds; errors only if all fail.
|
||||
pub async fn query_relay_ws(
|
||||
state: &AppState,
|
||||
relay_urls: &[String],
|
||||
filters: &[serde_json::Value],
|
||||
) -> Result<Vec<nostr::Event>, String> {
|
||||
let futures = relay_urls
|
||||
.iter()
|
||||
.map(|url| query_relay_ws_one(state, url, filters));
|
||||
let results = futures_util::future::join_all(futures).await;
|
||||
|
||||
let mut by_id: std::collections::HashMap<String, nostr::Event> =
|
||||
std::collections::HashMap::new();
|
||||
let mut last_err = None;
|
||||
let mut any_ok = false;
|
||||
for r in results {
|
||||
match r {
|
||||
Ok(events) => {
|
||||
any_ok = true;
|
||||
for ev in events {
|
||||
by_id.entry(ev.id.to_hex()).or_insert(ev);
|
||||
}
|
||||
}
|
||||
Err(e) => last_err = Some(e),
|
||||
}
|
||||
}
|
||||
if !any_ok {
|
||||
return Err(last_err.unwrap_or_else(|| "all relays failed".to_string()));
|
||||
}
|
||||
Ok(by_id.into_values().collect())
|
||||
}
|
||||
|
||||
async fn query_relay_ws_one(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
filters: &[serde_json::Value],
|
||||
@@ -137,19 +170,41 @@ pub async fn query_relay_ws(
|
||||
|
||||
/// Publish a signed event over a plain WebSocket and wait for the relay's
|
||||
/// `OK` acknowledgement. Mirrors `relay::submit_event`.
|
||||
/// Sign once, then publish to all relays. Succeeds if any relay accepts.
|
||||
pub async fn submit_event_ws(
|
||||
builder: EventBuilder,
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
relay_urls: &[String],
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
let keys = {
|
||||
let guard = state.keys.lock().map_err(|e| e.to_string())?;
|
||||
guard.clone()
|
||||
};
|
||||
|
||||
let event = builder
|
||||
.sign_with_keys(&keys)
|
||||
.map_err(|e| format!("failed to sign event: {e}"))?;
|
||||
|
||||
let futures = relay_urls
|
||||
.iter()
|
||||
.map(|url| submit_event_ws_one(&event, &keys, url));
|
||||
let results = futures_util::future::join_all(futures).await;
|
||||
|
||||
let mut last_err = None;
|
||||
for r in results {
|
||||
match r {
|
||||
Ok(resp) if resp.accepted => return Ok(resp),
|
||||
Ok(resp) => last_err = Some(format!("relay rejected event: {}", resp.message)),
|
||||
Err(e) => last_err = Some(e),
|
||||
}
|
||||
}
|
||||
Err(last_err.unwrap_or_else(|| "all relays failed".to_string()))
|
||||
}
|
||||
|
||||
async fn submit_event_ws_one(
|
||||
event: &nostr::Event,
|
||||
keys: &nostr::Keys,
|
||||
relay_url: &str,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
let event_id = event.id.to_hex();
|
||||
let event_json = serde_json::json!(["EVENT", event]).to_string();
|
||||
|
||||
@@ -196,7 +251,7 @@ pub async fn submit_event_ws(
|
||||
}
|
||||
"AUTH" => {
|
||||
if let Some(challenge) = arr.get(1).and_then(|v| v.as_str()) {
|
||||
if let Ok(auth_json) = build_auth_message(&keys, relay_url, challenge) {
|
||||
if let Ok(auth_json) = build_auth_message(keys, relay_url, challenge) {
|
||||
let _ = write.send(Message::Text(auth_json.into())).await;
|
||||
// Re-send the event after authenticating.
|
||||
let _ = write.send(Message::Text(event_json.clone().into())).await;
|
||||
@@ -238,6 +293,25 @@ pub async fn submit_event_ws(
|
||||
/// agent-profile sync, where the event is signed by the agent's keys rather
|
||||
/// than the user's identity key.
|
||||
pub async fn publish_signed_event_ws(
|
||||
event: &nostr::Event,
|
||||
keys: &nostr::Keys,
|
||||
relay_urls: &[String],
|
||||
) -> Result<(), String> {
|
||||
let futures = relay_urls
|
||||
.iter()
|
||||
.map(|url| publish_signed_event_ws_one(event, keys, url));
|
||||
let results = futures_util::future::join_all(futures).await;
|
||||
let mut last_err = None;
|
||||
for r in results {
|
||||
match r {
|
||||
Ok(()) => return Ok(()),
|
||||
Err(e) => last_err = Some(e),
|
||||
}
|
||||
}
|
||||
Err(last_err.unwrap_or_else(|| "all relays failed".to_string()))
|
||||
}
|
||||
|
||||
async fn publish_signed_event_ws_one(
|
||||
event: &nostr::Event,
|
||||
keys: &nostr::Keys,
|
||||
relay_url: &str,
|
||||
|
||||
@@ -17,5 +17,36 @@ export const DEFAULT_PUBLIC_RELAYS: readonly string[] = [
|
||||
"wss://nostr.wine",
|
||||
] as const;
|
||||
|
||||
/** The relay pre-filled when a user first enables serverless mode. */
|
||||
export const DEFAULT_SERVERLESS_RELAY = DEFAULT_PUBLIC_RELAYS[0];
|
||||
/** The relays pre-filled (comma-joined) when a user first enables serverless mode. */
|
||||
export const DEFAULT_SERVERLESS_RELAY = DEFAULT_PUBLIC_RELAYS.join(", ");
|
||||
|
||||
/** Parse a comma/whitespace-separated relay string into a clean list. */
|
||||
export function parseRelayList(value: string): string[] {
|
||||
return value
|
||||
.split(",")
|
||||
.map((s) => s.trim())
|
||||
.filter((s) => s.length > 0);
|
||||
}
|
||||
|
||||
/** Whether `relay` is present in the comma-separated `value`. */
|
||||
export function relayListIncludes(value: string, relay: string): boolean {
|
||||
return parseRelayList(value).includes(relay);
|
||||
}
|
||||
|
||||
/** Toggle `relay` in the comma-separated `value`, returning the new value. */
|
||||
export function toggleRelayInList(value: string, relay: string): string {
|
||||
const list = parseRelayList(value);
|
||||
const next = list.includes(relay)
|
||||
? list.filter((r) => r !== relay)
|
||||
: [...list, relay];
|
||||
return next.join(", ");
|
||||
}
|
||||
|
||||
/** Normalize each relay in a comma list to a ws/wss URL; returns comma-joined. */
|
||||
export function normalizeRelayList(value: string): string {
|
||||
return parseRelayList(value)
|
||||
.map((r) =>
|
||||
r.startsWith("ws://") || r.startsWith("wss://") ? r : `wss://${r}`,
|
||||
)
|
||||
.join(",");
|
||||
}
|
||||
|
||||
@@ -3,6 +3,9 @@ import * as React from "react";
|
||||
import {
|
||||
DEFAULT_PUBLIC_RELAYS,
|
||||
DEFAULT_SERVERLESS_RELAY,
|
||||
normalizeRelayList,
|
||||
relayListIncludes,
|
||||
toggleRelayInList,
|
||||
} from "@/features/workspaces/defaultRelays";
|
||||
import type { Workspace } from "@/features/workspaces/types";
|
||||
import {
|
||||
@@ -53,7 +56,9 @@ export function AddWorkspaceDialog({
|
||||
const workspace: Workspace = {
|
||||
id: crypto.randomUUID(),
|
||||
name: name.trim() || deriveWorkspaceName(relayUrl.trim()),
|
||||
relayUrl: normalizeRelayUrl(relayUrl.trim()),
|
||||
relayUrl: serverless
|
||||
? normalizeRelayList(relayUrl.trim())
|
||||
: normalizeRelayUrl(relayUrl.trim()),
|
||||
// Serverless workspaces never use a Sprout API token.
|
||||
token: serverless ? undefined : token.trim() || undefined,
|
||||
mode: serverless ? "serverless" : "sprout",
|
||||
@@ -101,7 +106,7 @@ export function AddWorkspaceDialog({
|
||||
className="text-sm font-medium text-foreground"
|
||||
htmlFor="ws-relay-url"
|
||||
>
|
||||
Relay URL
|
||||
{serverless ? "Relay URLs" : "Relay URL"}
|
||||
</label>
|
||||
<Input
|
||||
autoFocus
|
||||
@@ -116,18 +121,33 @@ export function AddWorkspaceDialog({
|
||||
value={relayUrl}
|
||||
/>
|
||||
{serverless ? (
|
||||
<div className="flex flex-wrap gap-1.5 pt-0.5">
|
||||
{DEFAULT_PUBLIC_RELAYS.map((relay) => (
|
||||
<button
|
||||
className="rounded-full border border-border bg-muted/40 px-2 py-0.5 text-xs text-muted-foreground transition-colors hover:border-primary/60 hover:text-foreground"
|
||||
key={relay}
|
||||
onClick={() => setRelayUrl(relay)}
|
||||
type="button"
|
||||
>
|
||||
{relay.replace("wss://", "")}
|
||||
</button>
|
||||
))}
|
||||
</div>
|
||||
<>
|
||||
<p className="text-xs text-muted-foreground">
|
||||
Connects to all listed relays for redundancy. Tap to
|
||||
add/remove.
|
||||
</p>
|
||||
<div className="flex flex-wrap gap-1.5 pt-0.5">
|
||||
{DEFAULT_PUBLIC_RELAYS.map((relay) => {
|
||||
const selected = relayListIncludes(relayUrl, relay);
|
||||
return (
|
||||
<button
|
||||
className={`rounded-full border px-2 py-0.5 text-xs transition-colors ${
|
||||
selected
|
||||
? "border-primary bg-primary/15 text-foreground"
|
||||
: "border-border bg-muted/40 text-muted-foreground hover:border-primary/60 hover:text-foreground"
|
||||
}`}
|
||||
key={relay}
|
||||
onClick={() =>
|
||||
setRelayUrl(toggleRelayInList(relayUrl, relay))
|
||||
}
|
||||
type="button"
|
||||
>
|
||||
{relay.replace("wss://", "")}
|
||||
</button>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
</>
|
||||
) : null}
|
||||
</div>
|
||||
<div className="flex flex-col gap-1.5">
|
||||
|
||||
@@ -7,6 +7,9 @@ import { Input } from "@/shared/ui/input";
|
||||
import {
|
||||
DEFAULT_PUBLIC_RELAYS,
|
||||
DEFAULT_SERVERLESS_RELAY,
|
||||
normalizeRelayList,
|
||||
relayListIncludes,
|
||||
toggleRelayInList,
|
||||
} from "../defaultRelays";
|
||||
import { initFirstWorkspace, deriveWorkspaceName } from "../workspaceStorage";
|
||||
|
||||
@@ -55,7 +58,7 @@ export function WelcomeSetup({
|
||||
// is the single source of truth — never copied into localStorage.
|
||||
const identity = await getIdentity();
|
||||
initFirstWorkspace(
|
||||
trimmedUrl,
|
||||
serverless ? normalizeRelayList(trimmedUrl) : trimmedUrl,
|
||||
identity.pubkey,
|
||||
serverless ? "serverless" : "sprout",
|
||||
);
|
||||
@@ -130,21 +133,34 @@ export function WelcomeSetup({
|
||||
value={relayUrl}
|
||||
/>
|
||||
{serverless ? (
|
||||
<div className="flex flex-wrap gap-1.5 pt-0.5">
|
||||
{DEFAULT_PUBLIC_RELAYS.map((relay) => (
|
||||
<button
|
||||
className="rounded-full border border-border bg-muted/40 px-2 py-0.5 text-xs text-muted-foreground transition-colors hover:border-primary/60 hover:text-foreground"
|
||||
key={relay}
|
||||
onClick={() => {
|
||||
setRelayUrl(relay);
|
||||
setError(null);
|
||||
}}
|
||||
type="button"
|
||||
>
|
||||
{relay.replace("wss://", "")}
|
||||
</button>
|
||||
))}
|
||||
</div>
|
||||
<>
|
||||
<p className="text-xs text-muted-foreground">
|
||||
Connects to all listed relays for redundancy. Tap to
|
||||
add/remove.
|
||||
</p>
|
||||
<div className="flex flex-wrap gap-1.5 pt-0.5">
|
||||
{DEFAULT_PUBLIC_RELAYS.map((relay) => {
|
||||
const selected = relayListIncludes(relayUrl, relay);
|
||||
return (
|
||||
<button
|
||||
className={`rounded-full border px-2 py-0.5 text-xs transition-colors ${
|
||||
selected
|
||||
? "border-primary bg-primary/15 text-foreground"
|
||||
: "border-border bg-muted/40 text-muted-foreground hover:border-primary/60 hover:text-foreground"
|
||||
}`}
|
||||
key={relay}
|
||||
onClick={() => {
|
||||
setRelayUrl(toggleRelayInList(relayUrl, relay));
|
||||
setError(null);
|
||||
}}
|
||||
type="button"
|
||||
>
|
||||
{relay.replace("wss://", "")}
|
||||
</button>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
</>
|
||||
) : null}
|
||||
</div>
|
||||
) : null}
|
||||
|
||||
Reference in New Issue
Block a user