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
83 changes: 83 additions & 0 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,89 @@ A `CancellationToken` coordinates shutdown across all three loops.

Slow clients: `ConnectionState::send()` uses `try_send` — if the send buffer is full, a grace counter increments. After `SLOW_CLIENT_GRACE_LIMIT` (3) consecutive full-buffer events, the connection is cancelled. A successful send resets the counter.

### Live Community Authorization Revalidation

Authenticated root and huddle/audio sockets retain their server-resolved
community, principal, admission-captured owner, admission route, whether that
owner was a relay member at admission, and whether closed-relay membership was
granted through that owner. A relay-local worker checks
these records against batched writer-authoritative snapshots every 5 seconds. It
keeps the immediate local disconnect and Redis fan-out paths; this scan is the
durable backstop when a pod misses fan-out.

The maximum authorization-staleness budget is 30 seconds. A database failure
may preserve the last successful decision for 10 seconds, measured from the
start of that session's last successful authoritative check; repeated failures
never extend the grace. The worker handles up to 500 distinct
`(community, principal, captured owner)`
targets per query, runs at most 4 queries at a time, gives each query 1 second,
and starts one shared 5-second lookup deadline before it snapshots and groups
the live-session registry. No database batch starts after that deadline;
unresolved targets are treated as lookup failures. Closing decisions are then applied in one
bounded sweep over the connection snapshot. At the default 10,000-connection
limit, the in-memory regression harness processed a full 10,000-session scan in
70 ms in the latest measured workstation run. A disposable PostgreSQL run
fetched all 10,000 authorization targets in 129 ms through the same 500-row,
four-query batching shape. Database lookups are bounded by the shared 5-second
deadline. The conservative staleness bound is 10 seconds of database grace, one
5-second cadence, two 5-second scan windows, and one 1-second query deadline:
26 seconds, below the 30-second budget. The relay's shared
`BUZZ_MAX_CONNECTIONS` semaphore caps root and audio sockets together at
10,000 by default; that is at most 20 batches and five query waves. The worker
creates no per-socket timers or per-connection revalidation queries. A scan
that exceeds its deadline treats unreturned targets as lookup failures and
closes their sessions when their existing grace has expired. If a result batch
is incomplete or malformed, every target in
that batch is treated as a lookup failure; no partial rows refresh
authorization.

Membership is required only when `require_relay_membership` is enabled. An
open-relay nonmember remains allowed unless directly or through its current
recorded owner it has an active community ban. An owner-derived member remains
allowed only while the relay, owner-delegation setting, and current owner
membership still permit that path. Removal of the admission-captured owner
revokes its session even if the agent's stored owner changed; captured-owner
membership never grants access to a directly admitted principal. The first
stored owner link also invalidates sessions admitted without one, matching the
existing reconnect path when NIP-OA materializes ownership; the durable check
detects that transition if its disconnect fan-out is missed. A directly
admitted member with a stable owner link still relies on its own membership,
and a change between already-recorded owners alone does not create a new
membership requirement for that principal. An active ban on the captured owner
also revokes that session. On open relays, an owner who was already a nonmember
at admission does not become a membership requirement. Moderation timeouts are
not read by this worker because they restrict writes, not an established
socket. All membership, ownership, and ban predicates include the bound
community. Admission reads the captured owner's membership from the writer
when the delegation check did not already establish it; that bounded point
lookup records the baseline needed to distinguish a later membership removal
from pre-existing open-relay nonmembership.

Each batch snapshot is the revalidation ordering point. A removal or ban
committed before that snapshot closes the session; a re-add or unban committed
before it keeps the session authorized. A later concurrent change is observed
on the next scan, and a close based on an earlier snapshot can be recovered by
reconnecting against current state. Admission binds the authenticated identity
before its final writer read and only publishes an admitted live record after
that check, so admission racing a removal is either cancelled by the immediate
path or rejected by the final read / next durable scan.

If the writer remains unavailable past the grace, the relay fails closed and
cancels the socket. NIP-FI root and audio sockets receive their existing
route-specific `authorization_denied` or `authorization_unavailable` terminal
response; other sockets receive the normal access-revoked or unavailable close
reason. This trades connection availability during a prolonged database
outage for a bounded authorization window. Reconnecting after the database
recovers performs the ordinary admission checks. The worker shares graceful
shutdown cancellation and has no migration or rollout flag; all relay pods
must run this code for the bound to hold deployment-wide.

Targets are ordered deterministically within each scan, so repeated deadline
hits can leave the same tail targets unqueried and fail-close them after grace.
This path does not rotate scan order or jitter authorization-driven closes; if
writer saturation lasts beyond grace, many sockets on a pod can close together
and create a reconnect herd.

### Step 5: Cleanup

On disconnect (any cause):
Expand Down
2 changes: 2 additions & 0 deletions Justfile
Original file line number Diff line number Diff line change
Expand Up @@ -591,6 +591,8 @@ test-unit:
+ test(=state::tests::disconnect_community_wins_reason_losing_nip_fi_does_not_enqueue_frame)
+ test(=state::tests::manager_disconnect_sets_reason_enqueues_frame_then_cancels)
+ test(/^api::nip_fi::/)'
# Retain partially consumed mesh receives across housekeeping ticks.
cargo nextest run -p buzz-relay --lib -E 'test(=mesh_boot::tests::demo_echo_retains_pending_receive_across_housekeeping_ticks)'
# boot_lifecycle spawns the real relay binary and asserts its startup
# lifecycle, including that buzz_startup_phase_* reaches /metrics. Its
# non-ignored tests need no Postgres or Redis; the Postgres cases are
Expand Down
9 changes: 5 additions & 4 deletions crates/buzz-db/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,10 +65,10 @@ pub(crate) use runtime::{
};
pub use store::{
admin_moderation, allowlist, api_token, archived_identities, artifact, channel,
channel_members, community, deletion, dm, event, feed, git_repo, moderation, operator_listener,
partition, personal_read, product_feedback, push, reaction, read_state, relay_admin_actions,
relay_invite, relay_members, relay_operators, reminder, replaceable, storage_accounting,
thread, thread_window, usage, user, workflow,
channel_members, community, deletion, dm, event, feed, git_repo, live_authorization,
moderation, operator_listener, partition, personal_read, product_feedback, push, reaction,
read_state, relay_admin_actions, relay_invite, relay_members, relay_operators, reminder,
replaceable, storage_accounting, thread, thread_window, usage, user, workflow,
};

pub use allowlist::AllowlistEntry;
Expand All @@ -81,6 +81,7 @@ pub use community::{
pub use error::{DbError, Result};
pub use event::{ChannelHeadPrecondition, ChannelHeadWriteStatus};
pub use event::{EventQuery, DEFAULT_MAX_PAGE_LIMIT};
pub use live_authorization::{LiveAuthorizationState, LiveAuthorizationTarget};
pub use reaction::ReactionEventInsertOutcome;
pub use reminder::DueReminder;
pub use usage::UsageMetricsLeader;
Expand Down
19 changes: 18 additions & 1 deletion crates/buzz-db/src/store/event_follow_up_postgres_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1018,13 +1018,30 @@ async fn assert_async_isolation(db: &Db, community: CommunityId, keys: &Keys) {
endpoint_grant, max_class, subscriptions, expires_at, updated_at) \
SELECT $1, author, installation_id, source_event_id, source_created_at, generation, \
active, endpoint_enabled, endpoint_hash, endpoint_grant, max_class, \
subscriptions, expires_at, clock_timestamp() FROM push_leases WHERE community_id=$2",
subscriptions, expires_at, \
(SELECT received_at + interval '1 second' FROM events WHERE community_id=$1 AND id=$3) \
FROM push_leases WHERE community_id=$2",
)
.bind(later.as_uuid())
.bind(community.as_uuid())
.bind(before_enrollment.id.as_bytes().as_slice())
.execute(&mut *activation)
.await
.unwrap();
let leases_strictly_newer: bool = sqlx::query_scalar(
"SELECT count(*) > 0 AND bool_and(l.updated_at > e.received_at) \
FROM push_leases l JOIN events e ON e.community_id=l.community_id \
WHERE l.community_id=$1 AND e.id=$2",
)
.bind(later.as_uuid())
.bind(before_enrollment.id.as_bytes().as_slice())
.fetch_one(&mut *activation)
.await
.unwrap();
assert!(
leases_strictly_newer,
"fixture leases must follow recorded receipt"
);
activation.commit().await.unwrap();
assert_eq!(match_count(&isolated, later, &before_enrollment).await, 0);
assert_eq!(
Expand Down
Loading
Loading