Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,12 @@ RELAY_URL=ws://localhost:3000
# Stable relay signing key (required). `just bootstrap` generates a random key in
# the gitignored .env file. Preserve that value across restarts and backups.
# BUZZ_RELAY_PRIVATE_KEY=<32-byte hex private key>
# NIP-CL channel labels: default off. Never opt in during a mixed-version rollout.
# Stop/fence old command AND metadata/repair writers, reconcile every community,
# then follow docs/channel-labels-rollout.md before setting these values.
# Disabling commands does not make an older metadata writer safe.
# BUZZ_NIP_CL_ENABLED=false
# BUZZ_NIP_CL_WRITER_CUTOVER=offline-v1
# Optional: path to the web UI dist directory. When set, the relay serves
# the web frontend at / for browser requests. Leave unset for local dev
# (use `just web` for Vite HMR instead).
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

27 changes: 17 additions & 10 deletions PLANS/REPLICA_FULL_READ_ROUTING_DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,13 @@ event during post-write verification and report a false conflict). Every such
read stays on `query_events` / writer.

Display reads — a page the user scrolls, a count on a badge, a history list —
tolerate bounded staleness and take the routed path.
tolerate bounded staleness and take the routed path, except when they also
select current authorization-sensitive state. Queries that can match channel
metadata (kind `39000`, including mixed-kind and ID-only filters) carry
`EventQuery.channel_metadata_read`: current relay head and channel ACL predicates
are evaluated on the writer before filtering and limits. The routed query,
bounded-query and count APIs enforce this pin regardless of client consistency
intent.

Adding, removing, or reclassifying a caller **requires updating the table
below**; the `query_events_routed` doc-comment points here.
Expand All @@ -53,24 +59,25 @@ in `crates/buzz-relay/src/api/bridge.rs` (`extract_consistency`).
## Caller classification table

Rows below are every `query_events_routed` / `query_events_routed_bounded`
call site at head, plus the writer-pinned canvas row this change adds. (The
`count` and feed routed families — `count_events_routed`,
call site at head, plus writer-pinned metadata and canvas reads. (The
other count and feed routed families —
`get_events_by_ids_routed`, `query_feed_*_routed`, and the `get_channel_window`
cursor/head reads — are all display or bounded-count surfaces on the routed
cursor/head reads — are display or bounded-count surfaces on the routed
path; they carry their own soundness notes at their definitions in
`crates/buzz-db/src/lib.rs` and are out of scope for this table.)

| Caller / path label | Pool | Justification |
|---|---|---|
| `bridge_query` (default `/query` filter) | routed | Display reads over the HTTP bridge; bounded staleness acceptable. |
| `bridge_query` (default `/query` filter) | routed, except metadata | Display reads over the HTTP bridge; metadata exception below. |
| `bridge_query` + `consistency: strong` | **writer** | Client-declared write-influencing read (canvas save precondition, restore precondition, post-write ancestry verification). |
| `req_historical` (WS REQ historical page) | routed | Display backfill of a subscription; per-row re-filter absorbs a briefly-stale row. |
| `req_historical` (WS REQ historical page) | routed, except metadata | Display backfill; metadata exception below. |
| `bridge_thread_aux` (`AuxReader::Routed`, thread aux page) | routed | Thread reply hydration; display, post-verified against the fence wall. |
| `bridge_count_fallback` (`query_events_routed_bounded`) | routed (bounded arm) | COUNT fallback that materializes rows; bounded arm only, never covered. |
| `count_req_fallback` (`query_events_routed_bounded`) | routed (bounded arm) | WS COUNT fallback that materializes rows; bounded arm only. |
| `bridge_count_fallback` (`query_events_routed_bounded`) | routed (bounded arm), except metadata | COUNT fallback that materializes rows; bounded arm only, never covered. |
| `count_req_fallback` (`query_events_routed_bounded`) | routed (bounded arm), except metadata | WS COUNT fallback that materializes rows; bounded arm only. |
| Any query / bounded query / count carrying `channel_metadata_read` | **writer** | Current kind-39000 head and current ACL are authority decisions, including mixed-kind and ID-only filters; no replica-staleness budget can relax them. |

The **writer** row is the only write-influencing entry; every other caller is a
display or count surface that tolerates bounded staleness. Client canvas reads
The metadata **writer** row is enforced by `buzz-db` in each routed entry point;
HTTP and WS callers attach the context in `handlers/req.rs`. Client canvas reads
that gate a write set `consistency: strong` (Desktop
`desktop/src-tauri/src/commands/canvas.rs`, CLI
`crates/buzz-cli/src/commands/channels.rs`) so they land on the writer row.
158 changes: 158 additions & 0 deletions crates/buzz-admin/src/channel_metadata.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
//! Operator signing for application-owned canonical metadata repair.

use anyhow::Result;
use buzz_core::{
kind::{KIND_NIP29_GROUP_ADMINS, KIND_NIP29_GROUP_MEMBERS, KIND_NIP29_GROUP_METADATA},
TenantContext,
};
use buzz_db::Db;
use nostr::{EventBuilder, Keys, Kind};
use uuid::Uuid;

pub(crate) async fn repair(
db: &Db,
tenant: &TenantContext,
channel: Uuid,
keys: &Keys,
) -> Result<bool> {
let mut write = db
.begin_channel_metadata_write(tenant.community(), channel, keys.public_key())
.await?;
if !write.needs_snapshot().await? {
write.rollback().await?;
return Ok(false);
}
let tags = write.snapshot_tags().await?;
let timestamp = write.snapshot_timestamp()?;
let keys = keys.clone();
let event = tokio::task::spawn_blocking(move || {
EventBuilder::new(Kind::Custom(KIND_NIP29_GROUP_METADATA as u16), "")
.tags(tags)
.allow_self_tagging()
.custom_created_at(timestamp)
.sign_with_keys(&keys)
})
.await??;
let max_bytes = buzz_core::relay::max_frame_bytes_from_env();
write.store_snapshot(&event, max_bytes).await?;
write.commit().await?;
Ok(true)
}

#[derive(Default)]
pub(crate) struct RepairSummary {
pub repaired: u64,
pub skipped: u64,
pub failed: u64,
}

pub(crate) async fn reconcile(
db: &Db,
tenant: &TenantContext,
target: Option<Uuid>,
keys: &Keys,
) -> Result<RepairSummary> {
let mut cursor = Uuid::nil();
let mut summary = RepairSummary::default();
loop {
let channels = if let Some(channel) = target {
vec![channel]
} else {
db.channel_metadata_repair_page(tenant.community(), cursor)
.await?
};
if channels.is_empty() {
break;
}
for channel in channels {
cursor = channel;
match reconcile_one(db, tenant, channel, keys, target.is_some()).await {
Ok(true) => summary.repaired += 1,
Ok(false) => summary.skipped += 1,
Err(error) => {
summary.failed += 1;
eprintln!("channel {channel}: reconciliation failed: {error:#}");
}
}
}
if target.is_some() {
break;
}
}
Ok(summary)
}

async fn reconcile_one(
db: &Db,
tenant: &TenantContext,
channel: Uuid,
keys: &Keys,
roster_only: bool,
) -> Result<bool> {
// Keep targeted reconciliation's contract: only refresh kind 39002.
db.get_channel(tenant.community(), channel).await?;
let repaired = if roster_only {
false
} else {
repair(db, tenant, channel, keys).await?
};
// Auxiliary snapshots commit independently of metadata. A current 39000 is
// not evidence that a previous bootstrap finished publishing 39001/39002.
let discovery = if roster_only {
Vec::new()
} else {
db.query_events_for_bootstrap(&buzz_db::event::EventQuery {
kinds: Some(vec![
KIND_NIP29_GROUP_ADMINS as i32,
KIND_NIP29_GROUP_MEMBERS as i32,
]),
authors: Some(vec![keys.public_key().to_bytes().to_vec()]),
d_tag: Some(channel.to_string()),
..buzz_db::event::EventQuery::for_community(tenant.community())
})
.await?
};
let missing = |kind| {
!discovery
.iter()
.any(|stored| buzz_core::kind::event_kind_u32(&stored.event) == kind)
};
let admins = !roster_only && missing(KIND_NIP29_GROUP_ADMINS);
let members = roster_only || missing(KIND_NIP29_GROUP_MEMBERS);
if !admins && !members {
return Ok(repaired);
}
for kind in [KIND_NIP29_GROUP_ADMINS, KIND_NIP29_GROUP_MEMBERS] {
if (kind == KIND_NIP29_GROUP_ADMINS && !admins)
|| (kind == KIND_NIP29_GROUP_MEMBERS && !members)
{
continue;
}
let relay = keys.public_key().to_bytes();
let mut snapshot = if kind == KIND_NIP29_GROUP_ADMINS {
db.lock_admin_snapshot(tenant.community(), channel, &relay)
.await?
} else {
db.lock_member_snapshot(tenant.community(), channel, &relay)
.await?
};
let tags = snapshot.snapshot_tags()?;
let timestamp = snapshot.snapshot_timestamp().await?;
let keys = keys.clone();
let event = tokio::task::spawn_blocking(move || {
EventBuilder::new(Kind::Custom(kind as u16), "")
.tags(tags)
.allow_self_tagging()
.custom_created_at(timestamp)
.sign_with_keys(&keys)
})
.await??;
let (_, inserted) = snapshot.replace_discovery_event(&event).await?;
anyhow::ensure!(
inserted,
"discovery publication superseded; rerun reconciliation"
);
snapshot.release().await?;
}
Ok(true)
}
Loading
Loading