mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): capture one relay target per nest regeneration
The archive read re-resolved the workspace relay override twice — once in fetch_relay_self for the NIP-11 signer and again in query_relay for the snapshot — while the rendered footer captured it a third time. A workspace switch between those reads could pair one relay's advertised signer with another relay's snapshot, fail open, and render archived agents as active. Capture the effective relay target (ws + api base) once, before any network work, via capture_relay_target, and thread it through fetch_archived_pubkeys_at (fetch_relay_self_at + query_relay_at) and the rendered footer, so signer, snapshot, and footer all belong to one relay. A deterministic seam test mutates the override between capture and fetch and proves the fetch never crosses relays. Also rename NestRegenCoalescer to NestRegenGate: it is an ordered, latest-write-wins gate, not a work coalescer — superseded generations still perform their relay reads and drop the result at commit. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
@@ -18,11 +18,43 @@ use crate::{
|
||||
app_state::AppState,
|
||||
events,
|
||||
relay::{
|
||||
classify_request_error, query_relay, relay_http_base_url, relay_ws_url_with_override,
|
||||
submit_event, SubmitEventResponse,
|
||||
classify_request_error, query_relay, query_relay_at, relay_api_base_url,
|
||||
relay_http_base_url, relay_ws_url, relay_ws_url_with_override, submit_event,
|
||||
workspace_relay_override, SubmitEventResponse,
|
||||
},
|
||||
};
|
||||
|
||||
/// A relay target resolved from a single workspace-override read, so a caller
|
||||
/// that performs several relay requests cannot mix two relays if the workspace
|
||||
/// override changes mid-flight.
|
||||
///
|
||||
/// `relay_ws_url_with_override` and `relay_api_base_url_with_override` each read
|
||||
/// the override independently; a workspace switch between two such reads can
|
||||
/// pair one relay's NIP-11 signer with another relay's snapshot query.
|
||||
/// Capturing both fields from one read — matching those two functions' exact
|
||||
/// precedence, including the standalone `BUZZ_RELAY_HTTP` path when no override
|
||||
/// is set — guarantees the pair is internally consistent.
|
||||
pub(crate) struct RelayTarget {
|
||||
/// Relay WebSocket URL (drives the NIP-11 fetch and the rendered footer).
|
||||
pub ws_url: String,
|
||||
/// Relay HTTP API base URL (drives `/query`).
|
||||
pub api_base_url: String,
|
||||
}
|
||||
|
||||
/// Capture the effective relay target once, before any network work.
|
||||
pub(crate) fn capture_relay_target(state: &AppState) -> RelayTarget {
|
||||
match workspace_relay_override(state) {
|
||||
Some(url) => RelayTarget {
|
||||
api_base_url: relay_http_base_url(&url),
|
||||
ws_url: url,
|
||||
},
|
||||
None => RelayTarget {
|
||||
ws_url: relay_ws_url(),
|
||||
api_base_url: relay_api_base_url(),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// ── Helpers ─────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Read `target`'s live `kind:0` event and extract the first valid NIP-OA
|
||||
@@ -228,8 +260,18 @@ struct RelayInformationDocument {
|
||||
}
|
||||
|
||||
pub(crate) async fn fetch_relay_self(state: &AppState) -> Result<Option<String>, String> {
|
||||
let relay_url = relay_ws_url_with_override(state);
|
||||
let http_url = relay_http_base_url(&relay_url);
|
||||
fetch_relay_self_at(state, &relay_ws_url_with_override(state)).await
|
||||
}
|
||||
|
||||
/// Like [`fetch_relay_self`] but reads NIP-11 from an explicit relay WS URL
|
||||
/// instead of re-resolving the workspace override. Used by
|
||||
/// [`fetch_archived_pubkeys_at`] so the advertised signer and the snapshot
|
||||
/// query belong to the same captured relay target.
|
||||
pub(crate) async fn fetch_relay_self_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
) -> Result<Option<String>, String> {
|
||||
let http_url = relay_http_base_url(relay_url);
|
||||
let response = state
|
||||
.http_client
|
||||
.get(&http_url)
|
||||
@@ -286,12 +328,25 @@ fn archived_pubkeys_from_snapshot(snapshot: &nostr::Event) -> Vec<String> {
|
||||
/// signature or wrong author, or a query error — **fails open** with an empty
|
||||
/// set rather than trusting unauthenticated relay-authoritative state.
|
||||
pub(crate) async fn fetch_archived_pubkeys(state: &AppState) -> Vec<String> {
|
||||
let Ok(Some(relay_self)) = fetch_relay_self(state).await else {
|
||||
fetch_archived_pubkeys_at(state, &capture_relay_target(state)).await
|
||||
}
|
||||
|
||||
/// Like [`fetch_archived_pubkeys`] but resolves both the NIP-11 signer and the
|
||||
/// snapshot query against one captured [`RelayTarget`] instead of re-reading
|
||||
/// the workspace override for each. This keeps a regeneration's advertised
|
||||
/// signer and its snapshot query on the same relay even if the workspace
|
||||
/// override changes between the two awaits.
|
||||
pub(crate) async fn fetch_archived_pubkeys_at(
|
||||
state: &AppState,
|
||||
target: &RelayTarget,
|
||||
) -> Vec<String> {
|
||||
let Ok(Some(relay_self)) = fetch_relay_self_at(state, &target.ws_url).await else {
|
||||
return vec![];
|
||||
};
|
||||
|
||||
let query = query_relay(
|
||||
let query = query_relay_at(
|
||||
state,
|
||||
&target.api_base_url,
|
||||
&[serde_json::json!({
|
||||
"authors": [relay_self.clone()],
|
||||
"kinds": [13535],
|
||||
@@ -490,4 +545,89 @@ mod tests {
|
||||
assert_eq!(minimal.content, "");
|
||||
assert!(minimal.reason.is_none());
|
||||
}
|
||||
|
||||
/// Regression for the cross-relay capture defect: `fetch_archived_pubkeys_at`
|
||||
/// must resolve BOTH the NIP-11 signer and the `/query` snapshot against the
|
||||
/// single captured [`RelayTarget`], never re-reading the live workspace
|
||||
/// override. Two loopback relays advertise distinct signers and archive
|
||||
/// distinct pubkeys; we capture relay A, then mutate the override to relay B
|
||||
/// before the fetch. Because capture happens once up front, the override's
|
||||
/// value at any later instant — including between the two archive awaits —
|
||||
/// is irrelevant by construction, so setting it to B is the strongest form
|
||||
/// of that perturbation. A must supply both the signer and the snapshot.
|
||||
///
|
||||
/// RED-on-revert: restore `fetch_archived_pubkeys` to read the override for
|
||||
/// each leg (`fetch_relay_self` + `query_relay`) and this returns B's pubkey.
|
||||
#[tokio::test]
|
||||
async fn archived_fetch_never_crosses_relays_mid_flight() {
|
||||
use crate::app_state::build_app_state;
|
||||
use crate::relay_admission::{reset_rate_limit_gate, TEST_SERIAL};
|
||||
use axum::{routing::get, routing::post, Json, Router};
|
||||
|
||||
let _serial = TEST_SERIAL.lock().await;
|
||||
reset_rate_limit_gate();
|
||||
|
||||
// Build a loopback relay that advertises `relay_keys` as its NIP-11
|
||||
// `self` and serves a relay-signed 13535 snapshot archiving `archived`.
|
||||
async fn spawn_relay(relay_keys: Keys, archived: String) -> String {
|
||||
let self_hex = relay_keys.public_key().to_hex();
|
||||
let snapshot = EventBuilder::new(Kind::Custom(13535), "")
|
||||
.tags([
|
||||
Tag::parse(["-"]).unwrap(),
|
||||
Tag::parse(["p", &archived]).unwrap(),
|
||||
])
|
||||
.sign_with_keys(&relay_keys)
|
||||
.unwrap();
|
||||
let snapshot_json = serde_json::to_value(&snapshot).unwrap();
|
||||
|
||||
let router = Router::new()
|
||||
.route(
|
||||
"/",
|
||||
get(move || {
|
||||
let self_hex = self_hex.clone();
|
||||
async move { Json(serde_json::json!({ "self": self_hex })) }
|
||||
}),
|
||||
)
|
||||
.route(
|
||||
"/query",
|
||||
post(move || {
|
||||
let snapshot_json = snapshot_json.clone();
|
||||
async move { Json(serde_json::json!([snapshot_json])) }
|
||||
}),
|
||||
);
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let addr = listener.local_addr().unwrap();
|
||||
tokio::spawn(async move {
|
||||
axum::serve(listener, router).await.ok();
|
||||
});
|
||||
format!("ws://{addr}")
|
||||
}
|
||||
|
||||
let relay_a_keys = Keys::generate();
|
||||
let relay_b_keys = Keys::generate();
|
||||
// Distinct archived pubkeys, unrelated to either relay's signing key —
|
||||
// nostr 0.37's EventBuilder silently drops a `p` tag that references the
|
||||
// event's own signer, so the archived key must not equal the relay key.
|
||||
let archived_on_a = Keys::generate().public_key().to_hex();
|
||||
let archived_on_b = Keys::generate().public_key().to_hex();
|
||||
let relay_a = spawn_relay(relay_a_keys, archived_on_a.clone()).await;
|
||||
let relay_b = spawn_relay(relay_b_keys, archived_on_b.clone()).await;
|
||||
|
||||
let state = build_app_state();
|
||||
|
||||
// Capture relay A, then swap the override to relay B before the fetch.
|
||||
*state.relay_url_override.lock().unwrap() = Some(relay_a.clone());
|
||||
let target = capture_relay_target(&state);
|
||||
*state.relay_url_override.lock().unwrap() = Some(relay_b.clone());
|
||||
|
||||
let archived = fetch_archived_pubkeys_at(&state, &target).await;
|
||||
|
||||
assert_eq!(
|
||||
archived,
|
||||
vec![archived_on_a],
|
||||
"signer and snapshot must both come from the captured relay A, \
|
||||
never the mutated override (relay B)"
|
||||
);
|
||||
reset_rate_limit_gate();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,8 +11,7 @@ use super::{load_managed_agents, load_personas, AgentDefinition, ManagedAgentRec
|
||||
#[cfg(test)]
|
||||
use super::{BackendKind, RespondTo};
|
||||
use crate::app_state::AppState;
|
||||
use crate::commands::fetch_archived_pubkeys;
|
||||
use crate::relay::relay_ws_url_with_override;
|
||||
use crate::commands::{capture_relay_target, fetch_archived_pubkeys_at};
|
||||
use std::collections::HashSet;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
@@ -666,19 +665,22 @@ pub fn upsert_managed_section(file_path: &Path, new_section_content: &str) -> io
|
||||
}
|
||||
|
||||
/// Serializes nest-context writes so a slow, stale regeneration cannot roll the
|
||||
/// file back over a newer one.
|
||||
/// file back over a newer one. This is an ordered, latest-write-wins gate — not
|
||||
/// a work coalescer: every superseded generation still performs its relay reads,
|
||||
/// then drops its result at commit time. Adding a true dirty-loop owner would be
|
||||
/// a larger change and is unwarranted at this user-driven trigger rate.
|
||||
///
|
||||
/// Each regeneration request claims a monotonic generation *synchronously* at
|
||||
/// request time (see [`NestRegenCoalescer::claim`]), so the generation encodes
|
||||
/// request time (see [`NestRegenGate::claim`]), so the generation encodes
|
||||
/// program order: boot's regen is claimed before `apply_workspace`'s, an edit's
|
||||
/// regen before the next edit's. The claimed generation travels with the
|
||||
/// spawned task and gates its write in [`NestRegenCoalescer::commit`]: a task
|
||||
/// spawned task and gates its write in [`NestRegenGate::commit`]: a task
|
||||
/// whose generation is below the high-water mark drops its result instead of
|
||||
/// overwriting the newer file. The commit lock is held across the compare and
|
||||
/// the synchronous file write, so two tasks cannot both read a stale mark and
|
||||
/// both write — a bare mutex acquired around only the write would still permit
|
||||
/// that stale-last rollback.
|
||||
struct NestRegenCoalescer {
|
||||
struct NestRegenGate {
|
||||
/// Next generation to hand out. `fetch_add` under a single atomic preserves
|
||||
/// the program order of `claim` calls regardless of memory ordering.
|
||||
next_gen: AtomicU64,
|
||||
@@ -686,7 +688,7 @@ struct NestRegenCoalescer {
|
||||
last_written: Mutex<u64>,
|
||||
}
|
||||
|
||||
impl NestRegenCoalescer {
|
||||
impl NestRegenGate {
|
||||
const fn new() -> Self {
|
||||
Self {
|
||||
next_gen: AtomicU64::new(1),
|
||||
@@ -709,7 +711,7 @@ impl NestRegenCoalescer {
|
||||
let mut last = self
|
||||
.last_written
|
||||
.lock()
|
||||
.expect("nest regen coalescer lock poisoned");
|
||||
.expect("nest regen gate lock poisoned");
|
||||
if generation < *last {
|
||||
return Ok(false);
|
||||
}
|
||||
@@ -719,8 +721,8 @@ impl NestRegenCoalescer {
|
||||
}
|
||||
}
|
||||
|
||||
/// Process-wide coalescer for nest-context regeneration.
|
||||
static NEST_REGEN: NestRegenCoalescer = NestRegenCoalescer::new();
|
||||
/// Process-wide ordered write gate for nest-context regeneration.
|
||||
static NEST_REGEN: NestRegenGate = NestRegenGate::new();
|
||||
|
||||
pub async fn regenerate_nest_context(app: &AppHandle, generation: u64) -> Result<(), String> {
|
||||
let nest = nest_dir().ok_or("cannot resolve home directory for nest")?;
|
||||
@@ -733,15 +735,22 @@ pub async fn regenerate_nest_context(app: &AppHandle, generation: u64) -> Result
|
||||
let personas = load_personas(app)?;
|
||||
let agents = load_managed_agents(app)?;
|
||||
let state = app.state::<AppState>();
|
||||
let relay_url = relay_ws_url_with_override(&state);
|
||||
// Capture the relay target once, before any network work, so this
|
||||
// generation's rendered footer, NIP-11 signer, and snapshot query all
|
||||
// belong to one relay even if a workspace switch changes the override
|
||||
// between the two archive awaits below.
|
||||
let target = capture_relay_target(&state);
|
||||
// Identity-archived agents live only in the relay's `kind:13535` snapshot;
|
||||
// local records all read `is_active: true`. Fails open (empty set → render
|
||||
// everyone) so an unreachable relay can't blank the roster. The archive read
|
||||
// and the relay URL above are read for this generation; a later generation's
|
||||
// uses the same captured target as the rendered relay; a later generation's
|
||||
// task always wins the commit, so a fallback-relay boot render cannot bury a
|
||||
// later apply_workspace render.
|
||||
let archived: HashSet<String> = fetch_archived_pubkeys(&state).await.into_iter().collect();
|
||||
let content = render_dynamic_section(&personas, &agents, &archived, &relay_url);
|
||||
let archived: HashSet<String> = fetch_archived_pubkeys_at(&state, &target)
|
||||
.await
|
||||
.into_iter()
|
||||
.collect();
|
||||
let content = render_dynamic_section(&personas, &agents, &archived, &target.ws_url);
|
||||
NEST_REGEN
|
||||
.commit(&agents_md, &content, generation)
|
||||
.map_err(|e| format!("regenerate nest context: {e}"))?;
|
||||
@@ -755,7 +764,7 @@ pub async fn regenerate_nest_context(app: &AppHandle, generation: u64) -> Result
|
||||
/// All call sites treat regeneration as fire-and-forget — agents run fine with
|
||||
/// a stale AGENTS.md, so we warn and continue rather than propagating the error.
|
||||
/// The generation is claimed *here*, synchronously, so it encodes call order;
|
||||
/// the spawned task carries it into [`NestRegenCoalescer::commit`], which drops
|
||||
/// the spawned task carries it into [`NestRegenGate::commit`], which drops
|
||||
/// a stale render rather than letting a slow task overwrite a newer file. A
|
||||
/// just-archived agent may linger for one regen cycle until the next regen (any
|
||||
/// agent/team edit or the next launch).
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
//! Tests for the dynamic AGENTS.md section renderer, the managed-section
|
||||
//! upsert, and the regeneration coalescer. Split from `tests.rs` to keep
|
||||
//! upsert, and the regeneration gate. Split from `tests.rs` to keep
|
||||
//! each test file under the repository's per-file line ratchet.
|
||||
|
||||
use super::*;
|
||||
@@ -517,19 +517,19 @@ fn commit_newer_generation_wins_over_a_stale_finisher() {
|
||||
// When A finally finishes and commits LAST, its lower generation is dropped
|
||||
// so the file still reflects B. Ordering of *finishing* is the only variable —
|
||||
// the generation, claimed at request time, decides the winner.
|
||||
let coalescer = NestRegenCoalescer::new();
|
||||
let gate = NestRegenGate::new();
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let file = agents_md_with_markers(tmp.path());
|
||||
|
||||
let gen_a = coalescer.claim(); // pre-edit request
|
||||
let gen_b = coalescer.claim(); // post-edit request
|
||||
let gen_a = gate.claim(); // pre-edit request
|
||||
let gen_b = gate.claim(); // post-edit request
|
||||
assert!(gen_a < gen_b);
|
||||
|
||||
// B (newer) commits first.
|
||||
assert!(coalescer.commit(&file, "post-edit roster", gen_b).unwrap());
|
||||
assert!(gate.commit(&file, "post-edit roster", gen_b).unwrap());
|
||||
// A (older) finishes last and must be dropped.
|
||||
assert!(
|
||||
!coalescer.commit(&file, "pre-edit roster", gen_a).unwrap(),
|
||||
!gate.commit(&file, "pre-edit roster", gen_a).unwrap(),
|
||||
"a stale (lower-generation) render must not overwrite a newer one"
|
||||
);
|
||||
|
||||
@@ -547,15 +547,15 @@ fn commit_boot_fallback_relay_cannot_bury_apply_workspace_relay() {
|
||||
// fallback relay) is claimed first but finishes last; the apply_workspace
|
||||
// regen (generation 2, workspace relay) commits first. The workspace relay
|
||||
// render must survive even though the fallback-relay task writes afterward.
|
||||
let coalescer = NestRegenCoalescer::new();
|
||||
let gate = NestRegenGate::new();
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let file = agents_md_with_markers(tmp.path());
|
||||
|
||||
let boot_gen = coalescer.claim(); // boot, fallback relay
|
||||
let apply_gen = coalescer.claim(); // apply_workspace, workspace relay
|
||||
let boot_gen = gate.claim(); // boot, fallback relay
|
||||
let apply_gen = gate.claim(); // apply_workspace, workspace relay
|
||||
|
||||
// apply_workspace's render lands first.
|
||||
assert!(coalescer
|
||||
assert!(gate
|
||||
.commit(
|
||||
&file,
|
||||
"## Workspace\n- Relay: wss://workspace.example",
|
||||
@@ -563,7 +563,7 @@ fn commit_boot_fallback_relay_cannot_bury_apply_workspace_relay() {
|
||||
)
|
||||
.unwrap());
|
||||
// Boot's slower fallback-relay render finishes last and is dropped.
|
||||
assert!(!coalescer
|
||||
assert!(!gate
|
||||
.commit(
|
||||
&file,
|
||||
"## Workspace\n- Relay: wss://fallback.example",
|
||||
@@ -583,14 +583,14 @@ fn commit_boot_fallback_relay_cannot_bury_apply_workspace_relay() {
|
||||
fn commit_equal_generation_is_allowed() {
|
||||
// The gate rejects only strictly-lower generations. Re-committing the same
|
||||
// generation (e.g. a retried request) is permitted and refreshes the file.
|
||||
let coalescer = NestRegenCoalescer::new();
|
||||
let gate = NestRegenGate::new();
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let file = agents_md_with_markers(tmp.path());
|
||||
|
||||
let gen = coalescer.claim();
|
||||
assert!(coalescer.commit(&file, "first", gen).unwrap());
|
||||
let gen = gate.claim();
|
||||
assert!(gate.commit(&file, "first", gen).unwrap());
|
||||
assert!(
|
||||
coalescer.commit(&file, "second", gen).unwrap(),
|
||||
gate.commit(&file, "second", gen).unwrap(),
|
||||
"an equal generation must still be allowed to write"
|
||||
);
|
||||
|
||||
|
||||
@@ -31,7 +31,7 @@ pub fn relay_ws_url() -> String {
|
||||
|
||||
/// Read the workspace relay URL override, if set. Returns `None` when no
|
||||
/// override is active or when the mutex is poisoned (best-effort).
|
||||
fn workspace_relay_override(state: &AppState) -> Option<String> {
|
||||
pub(crate) fn workspace_relay_override(state: &AppState) -> Option<String> {
|
||||
state
|
||||
.relay_url_override
|
||||
.lock()
|
||||
|
||||
Reference in New Issue
Block a user