feat(relay): expose production authorization reachability

Signed-off-by: Cea Stapleton Cordasco <261786559+cea-block@users.noreply.github.com>
This commit is contained in:
Cea Stapleton Cordasco
2026-08-04 12:19:13 -05:00
parent 592f267902
commit 3a682c6fc1
14 changed files with 3070 additions and 90 deletions
+191 -1
View File
@@ -15,8 +15,16 @@ pub mod admin_moderation;
pub mod api_token;
/// Relay-scoped archived identity persistence (NIP-IA).
pub mod archived_identities;
/// Transaction-owned admission records for protected audio sessions.
pub mod audio_admission;
/// Durable provider-neutral authorization invalidation authority.
pub mod authorization_invalidation;
/// Restore-independent high-water snapshots for protected authority.
pub mod authorization_version;
/// Channel and membership persistence.
pub mod channel;
/// Transaction-owned current-only client verification-status revisions.
pub mod client_status;
/// Direct message channel persistence.
pub mod dm;
/// Database error types.
@@ -39,6 +47,12 @@ pub mod moderation;
pub mod partition;
/// Buzz product-feedback sidecar persistence.
pub mod product_feedback;
/// PostgreSQL-authoritative visibility for protected object-store content.
pub mod protected_publication;
/// Monotonic migration and cutover authority for protected object visibility.
pub mod protected_visibility;
/// Durable reconciliation for optional relay-authored identity projections.
pub mod public_projection;
/// Community-scoped push lease and durable wake-outbox persistence.
pub mod push;
/// Reaction persistence.
@@ -69,7 +83,7 @@ use uuid::Uuid;
use buzz_core::{CommunityId, StoredEvent};
fn event_replacement_lock_key(
pub(crate) fn event_replacement_lock_key(
community_id: CommunityId,
kind: i32,
pubkey: &[u8],
@@ -172,6 +186,43 @@ pub async fn insert_mentions(
Ok(())
}
/// Transaction-aware mention-index projection for a protected event commit.
pub async fn insert_mentions_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
community_id: CommunityId,
event: &nostr::Event,
channel_id: Option<Uuid>,
) -> Result<()> {
let created_at_secs = event.created_at.as_secs() as i64;
let created_at = DateTime::from_timestamp(created_at_secs, 0)
.ok_or(DbError::InvalidTimestamp(created_at_secs))?;
for pubkey in event.tags.iter().filter_map(|tag| {
let parts = tag.as_slice();
(parts.len() >= 2
&& parts[0] == "p"
&& parts[1].len() == 64
&& parts[1]
.chars()
.all(|character| character.is_ascii_hexdigit()))
.then(|| parts[1].to_ascii_lowercase())
}) {
sqlx::query(
"INSERT INTO event_mentions \
(community_id, pubkey_hex, event_id, event_created_at, channel_id, event_kind) \
VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT DO NOTHING",
)
.bind(community_id.as_uuid())
.bind(pubkey)
.bind(event.id.as_bytes().as_slice())
.bind(created_at)
.bind(channel_id)
.bind(event.kind.as_u16() as i32)
.execute(&mut **tx)
.await?;
}
Ok(())
}
/// Database handle. Clone is cheap (Arc-backed pool).
#[derive(Clone, Debug)]
pub struct Db {
@@ -1933,6 +1984,17 @@ impl Db {
push::claim_due_match_batch(&self.pool, limit, lease_until).await
}
/// Claim a matcher batch outside exact protected Enforce domains.
pub async fn claim_due_push_match_batch_excluding(
&self,
limit: i64,
lease_until: DateTime<Utc>,
excluded_communities: &[Uuid],
) -> Result<Option<push::ClaimedMatchBatch>> {
push::claim_due_match_batch_excluding(&self.pool, limit, lease_until, excluded_communities)
.await
}
/// Load active endpoint-enabled leases eligible for push matching.
pub async fn active_push_match_leases(
&self,
@@ -1967,6 +2029,14 @@ impl Db {
push::reap_exhausted_matches(&self.pool).await
}
/// Reap matcher jobs outside exact protected Enforce domains.
pub async fn reap_exhausted_push_matches_excluding(
&self,
excluded_communities: &[Uuid],
) -> Result<u64> {
push::reap_exhausted_matches_excluding(&self.pool, excluded_communities).await
}
/// Idempotently enqueue a wake for a matched lease and event.
pub async fn enqueue_push_wake(
&self,
@@ -2320,6 +2390,27 @@ impl Db {
channel::get_accessible_channel_ids(&self.pool, community_id, pubkey).await
}
/// Revalidate uncached read access to one channel at an outbound release
/// boundary.
pub async fn channel_read_authorized(
&self,
community_id: CommunityId,
channel_id: Uuid,
pubkey: &[u8],
) -> Result<bool> {
channel::channel_read_authorized(&self.pool, community_id, channel_id, pubkey).await
}
/// Revalidate uncached read access to a complete channel set in one query.
pub async fn channel_set_read_authorized(
&self,
community_id: CommunityId,
channel_ids: &[Uuid],
pubkey: &[u8],
) -> Result<bool> {
channel::channel_set_read_authorized(&self.pool, community_id, channel_ids, pubkey).await
}
/// Lists channels, optionally filtered by visibility.
pub async fn list_channels(
&self,
@@ -2454,6 +2545,14 @@ impl Db {
channel::reap_expired_ephemeral_channels(&self.pool).await
}
/// Archive expired ephemeral channels outside protected Enforce domains.
pub async fn reap_expired_ephemeral_channels_excluding(
&self,
excluded_communities: &[Uuid],
) -> Result<Vec<channel::ReapedEphemeralChannel>> {
channel::reap_expired_ephemeral_channels_excluding(&self.pool, excluded_communities).await
}
/// Query due reminders ready for delivery.
pub async fn query_due_reminders(
&self,
@@ -2463,6 +2562,22 @@ impl Db {
event::query_due_reminders(&self.pool, now_secs, batch_limit).await
}
/// Query reminders outside protected Enforce domains.
pub async fn query_due_reminders_excluding(
&self,
now_secs: i64,
batch_limit: i64,
excluded_communities: &[Uuid],
) -> Result<Vec<event::DueReminder>> {
event::query_due_reminders_excluding(
&self.pool,
now_secs,
batch_limit,
excluded_communities,
)
.await
}
/// Atomically claim a due reminder for delivery (cross-pod dedup).
pub async fn claim_due_reminder(
&self,
@@ -4333,6 +4448,16 @@ impl Db {
relay_invite::reap_expired_relay_invites(&self.pool, cutoff).await
}
/// Delete expired invites outside protected Enforce domains.
pub async fn reap_expired_relay_invites_excluding(
&self,
cutoff: chrono::DateTime<chrono::Utc>,
excluded_communities: &[Uuid],
) -> Result<u64> {
relay_invite::reap_expired_relay_invites_excluding(&self.pool, cutoff, excluded_communities)
.await
}
/// Atomically claims a v2 relay invite. The full redemption (membership
/// insert, policy evidence, use_count increment) runs in one PostgreSQL
/// transaction with `FOR UPDATE` on the invite row.
@@ -4570,6 +4695,16 @@ impl Db {
git_repo::count_repos_for_owner(&self.pool, community, owner_pubkey).await
}
/// Return an existing Git reservation's immutable publication origin.
pub async fn repo_publication_origin(
&self,
community_id: CommunityId,
repo_id: &str,
owner_pubkey: &str,
) -> Result<Option<String>> {
git_repo::repo_publication_origin(&self.pool, community_id, repo_id, owner_pubkey).await
}
/// Release a git repo name reservation held by `owner_pubkey` (rollback).
///
/// Returns the number of rows removed (0 or 1). See [`git_repo::release_repo_name`].
@@ -6281,6 +6416,61 @@ mod tests {
assert_eq!(retry, restored);
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn protected_community_lifecycle_is_fail_closed_while_off_remains_legacy() {
let db = setup_db().await;
let owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple());
let protected_host = format!("protected-lifecycle-{}.example", Uuid::new_v4().simple());
let created = db
.create_community_with_owner(&protected_host, &owner)
.await
.expect("create protected fixture");
let CreateCommunityWithOwnerResult::Created(protected) = created else {
panic!("expected new protected fixture");
};
sqlx::query("INSERT INTO authorization_invalidation_domains (community_id) VALUES ($1)")
.bind(protected.id.as_uuid())
.execute(&db.pool)
.await
.expect("activate protected marker");
assert!(db
.archive_community_owned_by(&protected_host, &owner, "reserved.example")
.await
.is_err());
assert!(sqlx::query("DELETE FROM communities WHERE id=$1")
.bind(protected.id.as_uuid())
.execute(&db.pool)
.await
.is_err());
assert!(db
.lookup_community_by_host(&protected_host)
.await
.expect("protected lookup")
.is_some());
let off_host = format!("off-lifecycle-{}.example", Uuid::new_v4().simple());
let created = db
.create_community_with_owner(&off_host, &owner)
.await
.expect("create Off fixture");
let CreateCommunityWithOwnerResult::Created(off) = created else {
panic!("expected new Off fixture");
};
assert!(db
.archive_community_owned_by(&off_host, &owner, "reserved.example")
.await
.expect("Off archive keeps legacy behavior")
.is_some());
assert!(db
.unarchive_community_owned_by(&off_host, &owner)
.await
.expect("Off restore keeps legacy behavior")
.is_some());
assert!(!off.id.as_uuid().is_nil());
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn create_community_with_owner_enforces_per_owner_limit() {
+203 -2
View File
@@ -566,7 +566,7 @@ mod tests {
let mut migrations: Vec<_> = MIGRATOR.iter().collect();
migrations.sort_by_key(|migration| migration.version);
assert_eq!(migrations.len(), 29);
assert_eq!(migrations.len(), 44);
assert_eq!(migrations[0].version, 1);
assert_eq!(&*migrations[0].description, "initial schema");
assert!(migrations[0]
@@ -617,6 +617,39 @@ mod tests {
.as_str()
.contains("CREATE INDEX idx_events_tags_gin"));
assert!(!migrations[0].sql.as_str().contains("idx_events_tags_gin"));
assert_eq!(migrations[34].version, 35);
assert!(migrations[34]
.sql
.as_str()
.contains("ADD COLUMN publication_origin"));
assert!(!migrations[0].sql.as_str().contains("publication_origin"));
assert_eq!(migrations[39].version, 40);
assert!(migrations[39]
.sql
.as_str()
.contains("protected_domain_marker_delete_guard"));
assert_eq!(migrations[40].version, 41);
assert!(migrations[40].sql.as_str().contains("cleanup_requested_at"));
assert_eq!(migrations[41].version, 42);
assert!(migrations[41]
.sql
.as_str()
.contains("git_policy_update_authority_epoch"));
assert_eq!(migrations[42].version, 43);
let audio_visibility = migrations[42].sql.as_str();
assert!(audio_visibility.contains("visibility_observed_at"));
assert!(audio_visibility.contains("'reserved', 'active', 'visible', 'aborted', 'finished'"));
assert!(audio_visibility.contains("audio_admission_visibility_transition_guard"));
assert!(audio_visibility
.contains("OLD.state = 'active' AND NEW.state IN ('visible', 'aborted')"));
assert_eq!(migrations[43].version, 44);
let projection_retirement = migrations[43].sql.as_str();
assert!(projection_retirement.contains("identity_public_projection_heads"));
assert!(projection_retirement.contains("identity_public_projection_retirements"));
assert!(projection_retirement.contains("source_binding_version"));
assert!(!projection_retirement.contains("issuer"));
assert!(!projection_retirement.contains("subject TEXT"));
assert!(!projection_retirement.contains("display_name"));
// NIP-AM (kind 44200) FTS exclusion: additive migration, never folded
// into 0001 — folding would change 0001's checksum and break brownfield
@@ -983,6 +1016,74 @@ mod tests {
"migration 0029 is missing {required}"
);
}
assert_eq!(migrations[29].version, 30);
let invalidation = migrations[29].sql.as_str();
assert!(invalidation.contains("CREATE TABLE authorization_invalidation_domains"));
assert!(invalidation.contains("CREATE TABLE authorization_invalidation_receipts"));
assert!(invalidation.contains("CREATE TABLE authorization_invalidation_floors"));
assert_eq!(migrations[30].version, 31);
let operation_receipts = migrations[30].sql.as_str();
assert!(operation_receipts.contains("CREATE TABLE authorization_operation_receipts"));
assert!(operation_receipts.contains("request_fingerprint"));
assert!(operation_receipts.contains("result_payload"));
assert!(operation_receipts.contains("authorization_operation_expiry_guard"));
assert_eq!(migrations[31].version, 32);
let protected_publications = migrations[31].sql.as_str();
assert!(protected_publications.contains("CREATE TABLE git_repo_publications"));
assert!(protected_publications.contains("CREATE TABLE media_publications"));
assert_eq!(migrations[32].version, 33);
let audio_admissions = migrations[32].sql.as_str();
assert!(audio_admissions.contains("CREATE TABLE audio_session_admissions"));
assert!(audio_admissions.contains("lease_expires_at"));
assert!(audio_admissions.contains("audio_session_admissions_channel_fk"));
assert_eq!(migrations[33].version, 34);
let protected_object_authority = migrations[33].sql.as_str();
assert!(protected_object_authority.contains("CREATE TABLE protected_object_authority"));
assert!(protected_object_authority.contains("inventory_sha256"));
assert!(
protected_object_authority.contains("state IN ('legacy', 'importing', 'postgresql')")
);
assert_eq!(migrations[34].version, 35);
assert!(migrations[34]
.sql
.as_str()
.contains("ADD COLUMN publication_origin"));
assert_eq!(migrations[35].version, 36);
let audio_lifecycle = migrations[35].sql.as_str();
assert!(audio_lifecycle.contains("ADD COLUMN state TEXT"));
assert!(audio_lifecycle.contains("'reserved', 'active', 'aborted', 'finished'"));
assert!(audio_lifecycle.contains("idx_audio_session_admissions_reconcile"));
assert_eq!(migrations[36].version, 37);
assert_eq!(migrations[37].version, 38);
assert_eq!(migrations[38].version, 39);
assert!(migrations[38]
.sql
.as_str()
.contains("protected_community_lifecycle_guard"));
let authority_epochs = migrations[36].sql.as_str();
assert!(authority_epochs.contains("CREATE TABLE authorization_authority_epochs"));
assert!(authority_epochs.contains("CREATE TABLE client_status_revisions"));
assert!(authority_epochs.contains("advance_authorization_authority_epoch"));
assert!(authority_epochs.contains("pg_trigger_depth() > 1"));
assert!(authority_epochs.contains("ON DELETE CASCADE"));
for protected_table in [
"identity_bindings",
"identity_principals",
"identity_revoked_keys",
"identity_retired_pairs",
"relay_members",
"channel_members",
"community_bans",
"channels",
"users",
"authorization_invalidation_domains",
"git_repo_publications",
"media_publications",
"protected_object_authority",
"audio_session_admissions",
] {
assert!(authority_epochs.contains(protected_table));
}
}
fn additive_identity_executable_sql(sql: &str) -> String {
@@ -1419,7 +1520,15 @@ mod tests {
run_migrations(&pool)
.await
.expect("retry succeeds after operator repair");
assert_eq!(applied_versions(&pool).await.last().copied(), Some(28));
let latest_version = MIGRATOR
.iter()
.map(|migration| migration.version)
.max()
.expect("embedded migration set is non-empty");
assert_eq!(
applied_versions(&pool).await.last().copied(),
Some(latest_version)
);
}
#[tokio::test]
@@ -2596,4 +2705,96 @@ mod tests {
"fresh installs must default non-allowlisted kinds to NULL: {search_expression}"
);
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn authority_triggers_preserve_off_and_deny_unwitnessed_protected_teardown() {
let pool = connect_test_pool().await;
reset_public_schema(&pool).await;
run_migrations(&pool).await.expect("apply all migrations");
let community_id = uuid::Uuid::new_v4();
sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)")
.bind(community_id)
.bind(format!(
"authority-trigger-{}.example",
community_id.simple()
))
.execute(&pool)
.await
.expect("insert legacy community");
sqlx::query(
"INSERT INTO relay_members (community_id, pubkey, role) VALUES ($1, $2, 'member')",
)
.bind(community_id)
.bind("11".repeat(32))
.execute(&pool)
.await
.expect("legacy membership remains writable");
let legacy_domains: i64 = sqlx::query_scalar(
"SELECT count(*) FROM authorization_invalidation_domains WHERE community_id=$1",
)
.bind(community_id)
.fetch_one(&pool)
.await
.expect("read legacy authorization rows");
assert_eq!(legacy_domains, 0, "Off must not acquire protected state");
sqlx::query("INSERT INTO authorization_invalidation_domains (community_id) VALUES ($1)")
.bind(community_id)
.execute(&pool)
.await
.expect("initialize protected domain");
sqlx::query(
"INSERT INTO relay_members (community_id, pubkey, role) VALUES ($1, $2, 'member')",
)
.bind(community_id)
.bind("22".repeat(32))
.execute(&pool)
.await
.expect("protected membership mutation");
let generation: i64 = sqlx::query_scalar(
"SELECT generation FROM authorization_invalidation_domains WHERE community_id=$1",
)
.bind(community_id)
.fetch_one(&pool)
.await
.expect("read protected generation");
assert_eq!(generation, 1, "one mutation advances generation once");
sqlx::query("DELETE FROM relay_members WHERE community_id=$1")
.bind(community_id)
.execute(&pool)
.await
.expect("delete ordinary community-owned rows first");
assert!(sqlx::query("DELETE FROM communities WHERE id=$1")
.bind(community_id)
.execute(&pool)
.await
.is_err());
let retained: i64 = sqlx::query_scalar(
"SELECT count(*) FROM authorization_authority_epochs WHERE community_id=$1",
)
.bind(community_id)
.fetch_one(&pool)
.await
.expect("read retained authority state");
assert_eq!(retained, 1, "denied teardown retains the monotonic floor");
assert!(sqlx::query(
"DELETE FROM authorization_invalidation_domains WHERE community_id=$1"
)
.bind(community_id)
.execute(&pool)
.await
.is_err());
let marker: i64 = sqlx::query_scalar(
"SELECT count(*) FROM authorization_invalidation_domains WHERE community_id=$1",
)
.bind(community_id)
.fetch_one(&pool)
.await
.expect("read retained activation marker");
assert_eq!(marker, 1, "protected activation is a one-way cutover");
}
}
+7 -6
View File
@@ -223,12 +223,13 @@ async fn feedback_attachment(
return Err(ApiError::not_found());
}
let response = crate::api::media::serve_blob_for_tenant(&state, &tenant, &sha256, &headers)
.await
.map_err(|error| match error {
buzz_media::MediaError::NotFound => ApiError::not_found(),
_ => ApiError::internal(),
})?;
let response =
crate::api::media::serve_blob_for_tenant(&state, &tenant, &sha256, &headers, None)
.await
.map_err(|error| match error {
buzz_media::MediaError::NotFound => ApiError::not_found(),
_ => ApiError::internal(),
})?;
tracing::info!(
feedback_id = %feedback.id,
community_id = %feedback.community_id,
+4 -7
View File
@@ -6,6 +6,7 @@ pub mod events;
pub mod git;
pub mod invites;
pub mod media;
pub mod media_migration;
pub mod mesh_demo;
pub mod nip05;
pub mod operator;
@@ -92,16 +93,12 @@ pub mod relay_members {
.await
.map_err(|e| format!("relay membership check (owner) failed: {e}"))?;
if owner_is_member {
debug!(
agent = %pubkey_hex,
owner = %owner_hex,
"NIP-OA membership granted via owner"
);
debug!("NIP-OA membership granted via owner");
return Ok(MembershipDecision::ViaOwner(owner_pubkey));
}
}
Err(e) => {
info!(agent = %pubkey_hex, "NIP-OA auth tag invalid: {e}");
info!("NIP-OA auth tag invalid: {e}");
}
}
}
@@ -186,7 +183,7 @@ pub mod relay_members {
Ok(true) => {
metrics::counter!(
"buzz_users_created_total",
"community" => tenant.host().to_owned()
"community" => crate::metrics::community_label(tenant.community())
)
.increment(1);
}
@@ -0,0 +1,22 @@
//! Provider-neutral runtime authorization seams.
//!
//! This commit registers the complete O4 module shape while implementing only
//! exact-domain provider selection, provider-evidence finalization, and bounded
//! leases. Transport adoption, invalidation, and client status remain separate
//! extension lanes.
pub(crate) mod ephemeral;
/// Transaction-owned protected mutation execution and idempotency.
pub mod executor;
/// Exact-domain provider selection and authorization finalization.
pub mod finalization;
/// Durable provider-neutral invalidation, reconciliation, and use fences.
pub mod invalidation;
/// Disabled-by-default production runtime construction.
pub mod production;
/// Independent high-water protection against stale PostgreSQL restoration.
pub mod restore;
/// Reserved provider-neutral client-status extension seam.
pub mod status;
/// Reserved provider-neutral transport-adoption extension seam.
pub mod transport;
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,155 @@
//! Durable PostgreSQL implementation behind the storage-agnostic status seam.
use async_trait::async_trait;
use uuid::Uuid;
use super::{
ClientStatusCurrentRequirement, ClientStatusIssuanceReceipt, ClientStatusRevisionScope,
DurableClientStatusRevision, DurableClientStatusRevisionSource,
};
/// PostgreSQL-backed revision source coupled to the independent restore witness.
///
/// Construction does not enable presentation; the unconstructible presentation
/// permit remains the separate runtime gate.
pub struct PostgresClientStatusRevisionSource {
db: buzz_db::Db,
restore: std::sync::Arc<super::super::restore::RestoreProtectionRuntime>,
}
impl PostgresClientStatusRevisionSource {
/// Bind the writer database and the exact initialized restore runtime.
pub fn new(
db: buzz_db::Db,
restore: std::sync::Arc<super::super::restore::RestoreProtectionRuntime>,
) -> Self {
Self { db, restore }
}
async fn reconcile_allocation(
&self,
scope: ClientStatusRevisionScope,
operation_id: Uuid,
request_fingerprint: [u8; 32],
result: Result<
buzz_db::client_status::AllocatedStatusRevision,
buzz_db::client_status::ClientStatusAllocationError,
>,
witness: super::super::restore::RestoreMutationGuard,
) -> Option<DurableClientStatusRevision> {
match result {
Ok(revision) => {
witness.commit().await.ok()?;
DurableClientStatusRevision::from_durable_state(revision.revision, revision.floor)
.ok()
}
Err(buzz_db::client_status::ClientStatusAllocationError::CommitUnknown(_)) => {
witness.commit().await.ok()?;
let revision = self
.db
.committed_status_revision(
scope.authorization_domain(),
operation_id,
request_fingerprint,
)
.await
.ok()??;
DurableClientStatusRevision::from_durable_state(revision.revision, revision.floor)
.ok()
}
Err(_) => {
let _ = witness.abort().await;
None
}
}
}
}
#[async_trait]
impl DurableClientStatusRevisionSource for PostgresClientStatusRevisionSource {
async fn current_revision_for(
&self,
requirement: &ClientStatusCurrentRequirement<'_>,
issuance_fingerprint: [u8; 32],
) -> Option<DurableClientStatusRevision> {
let scope = requirement.scope();
let operation_id = super::super::executor::ProtectedOperationId::derive(
scope.authorization_domain(),
"client.status.current.v1",
&issuance_fingerprint,
)
.ok()?
.as_uuid();
let witness = self
.restore
.begin(
scope.authorization_domain(),
operation_id,
issuance_fingerprint,
)
.await
.ok()?;
let event_author_pubkey = scope.event_author_pubkey().to_bytes();
let result = self
.db
.allocate_current_status_revision(buzz_db::client_status::CurrentStatusAllocation {
community_id: scope.authorization_domain(),
event_author_pubkey: &event_author_pubkey,
binding_id: requirement.binding_id(),
binding_version: requirement.binding_version().get(),
policy_version: requirement.policy_version().as_str(),
evaluation_generation: requirement.evaluation_generation(),
fresh_until: requirement.fresh_until(),
operation_id,
request_fingerprint: issuance_fingerprint,
})
.await;
self.reconcile_allocation(scope, operation_id, issuance_fingerprint, result, witness)
.await
}
async fn withdrawal_revision_for(
&self,
receipt: &ClientStatusIssuanceReceipt,
withdrawal_fingerprint: [u8; 32],
) -> Option<DurableClientStatusRevision> {
let operation_id = super::super::executor::ProtectedOperationId::derive(
receipt.scope.authorization_domain(),
"client.status.withdraw.v1",
&withdrawal_fingerprint,
)
.ok()?
.as_uuid();
let witness = self
.restore
.begin(
receipt.scope.authorization_domain(),
operation_id,
withdrawal_fingerprint,
)
.await
.ok()?;
let event_author_pubkey = receipt.scope.event_author_pubkey().to_bytes();
let result = self
.db
.allocate_withdrawn_status_revision(
buzz_db::client_status::WithdrawalStatusAllocation {
community_id: receipt.scope.authorization_domain(),
event_author_pubkey: &event_author_pubkey,
supersedes_revision: receipt.revision,
issuance_fingerprint: receipt.issuance_fingerprint,
operation_id,
request_fingerprint: withdrawal_fingerprint,
},
)
.await;
self.reconcile_allocation(
receipt.scope,
operation_id,
withdrawal_fingerprint,
result,
witness,
)
.await
}
}
+204 -55
View File
@@ -12,6 +12,9 @@
use std::sync::Arc;
use axum::extract::ws::Message as WsMessage;
use buzz_auth::{
AuthTransport, VerifiedDelegationOutput, VerifiedEvidenceAdapter, VerifiedNostrProof,
};
use tracing::{debug, info, warn};
use crate::connection::{AuthState, ConnectionState};
@@ -42,6 +45,7 @@ pub fn extract_auth_tag_json(event: &nostr::Event) -> Option<String> {
#[tracing::instrument(skip_all, fields(event_id, conn_id))]
pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state: Arc<AppState>) {
let event_id_hex = event.id.to_hex();
let verified_event = event.clone();
let (challenge, conn_id) = {
let auth = conn.auth_state.read().await;
match &*auth {
@@ -146,7 +150,7 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
Ok(state) if state.banned => BanOutcome::Banned,
Ok(_) => BanOutcome::Clear,
Err(e) => {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), error = %e,
warn!(conn_id = %conn_id, error = %e,
"ban-state DB lookup failed, denying (fail-closed)");
BanOutcome::DbError
}
@@ -168,7 +172,7 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
Ok(state) if state.banned => BanOutcome::Banned,
Ok(_) => BanOutcome::Clear,
Err(e) => {
warn!(conn_id = %conn_id, owner = %owner.to_hex(), error = %e,
warn!(conn_id = %conn_id, error = %e,
"owner ban-state DB lookup failed, denying (fail-closed)");
BanOutcome::DbError
}
@@ -188,7 +192,7 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
};
if let Some((metric_reason, deny_reason)) = denial {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), reason = deny_reason, "principal denied at ban seam");
warn!(conn_id = %conn_id, reason = deny_reason, "principal denied at ban seam");
metrics::counter!("buzz_auth_failures_total", "reason" => metric_reason)
.increment(1);
*conn.auth_state.write().await = AuthState::Failed;
@@ -205,28 +209,63 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
}
}
let identity_lane = crate::authorization_runtime::transport::legacy_identity_lane(
&state,
conn.tenant.community(),
);
let identity_proof = match crate::corporate_identity::verify_corporate_identity(
&state,
conn.tenant.community(),
pubkey,
conn.corporate_identity_jwt.as_deref(),
conn.corporate_identity_assertion.as_ref(),
auth_tag_json.as_deref(),
)
.await
{
Ok(proof) => proof,
Ok(proof) => Some(proof),
Err(e) => {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), error = %e, "corporate identity denied");
*conn.auth_state.write().await = AuthState::Failed;
conn.send(RelayMessage::ok(
&event_id_hex,
false,
&format!("restricted: {}", e.public_message()),
));
return;
warn!(conn_id = %conn_id, error = ?e, "corporate identity denied");
if identity_lane
== crate::authorization_runtime::transport::LegacyIdentityLane::ObserveOnly
{
None
} else {
*conn.auth_state.write().await = AuthState::Failed;
conn.send(RelayMessage::ok(
&event_id_hex,
false,
&format!("restricted: {}", e.public_message()),
));
return;
}
}
};
let verified_assertion = match identity_proof.as_ref() {
Some(proof) => {
match crate::corporate_identity::current_verified_assertion_for_proof(
&state,
proof,
conn.tenant.community(),
AuthTransport::RelayWebSocket,
) {
Ok(assertion) => assertion.map(Arc::new),
Err(error) => {
warn!(conn_id = %conn_id, error = %error, "federated evidence sealing failed");
if identity_lane
== crate::authorization_runtime::transport::LegacyIdentityLane::ObserveOnly
{
None
} else {
*conn.auth_state.write().await = AuthState::Failed;
return;
}
}
}
}
None => None,
};
// Pubkey allowlist gate — only for pubkey-only auth.
if state.config.pubkey_allowlist_enabled
&& auth_ctx.auth_method == buzz_auth::AuthMethod::Nip42
@@ -238,13 +277,13 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
{
Ok(v) => v,
Err(e) => {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), error = %e,
warn!(conn_id = %conn_id, error = %e,
"allowlist DB lookup failed, denying (fail-closed)");
false
}
};
if !allowed {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), "pubkey not in allowlist");
warn!(conn_id = %conn_id, "pubkey not in allowlist");
metrics::counter!("buzz_auth_failures_total", "reason" => "allowlist_denied")
.increment(1);
*conn.auth_state.write().await = AuthState::Failed;
@@ -268,7 +307,7 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
{
Ok(owner) => owner,
Err(e) => {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), error = ?e, "not a relay member");
warn!(conn_id = %conn_id, error = ?e, "not a relay member");
metrics::counter!("buzz_auth_failures_total", "reason" => "not_relay_member")
.increment(1);
*conn.auth_state.write().await = AuthState::Failed;
@@ -281,30 +320,40 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
}
};
let identity_decision = match crate::corporate_identity::finalize_corporate_identity(
&state,
conn.tenant.community(),
pubkey,
identity_proof,
)
.await
let identity_decision = if identity_lane
== crate::authorization_runtime::transport::LegacyIdentityLane::Legacy
{
Ok(decision) => decision,
Err(e) => {
warn!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), error = %e, "corporate identity finalization denied");
*conn.auth_state.write().await = AuthState::Failed;
conn.send(RelayMessage::ok(
&event_id_hex,
false,
&format!("restricted: {}", e.public_message()),
));
return;
if let Some(identity_proof) = identity_proof.clone() {
match crate::corporate_identity::finalize_corporate_identity(
&state,
conn.tenant.community(),
pubkey,
identity_proof,
)
.await
{
Ok(decision) => Some(decision),
Err(e) => {
warn!(conn_id = %conn_id, error = ?e, "corporate identity finalization denied");
*conn.auth_state.write().await = AuthState::Failed;
conn.send(RelayMessage::ok(
&event_id_hex,
false,
&format!("restricted: {}", e.public_message()),
));
return;
}
}
} else {
None
}
} else {
None
};
if let crate::corporate_identity::CorporateIdentityDecision::Delegated {
if let Some(crate::corporate_identity::CorporateIdentityDecision::Delegated {
owner_pubkey,
..
} = &identity_decision
}) = &identity_decision
{
auth_ctx.agent_owner_pubkey = Some(*owner_pubkey);
}
@@ -327,37 +376,95 @@ pub async fn handle_auth(event: nostr::Event, conn: Arc<ConnectionState>, state:
// Stash NIP-OA owner on the auth context only after the shared
// backfill confirms the first-write-wins relationship.
if let Some(owner) = nip_oa_owner {
if crate::api::relay_members::materialize_nip_oa_owner(
&state,
&conn.tenant,
&pubkey,
&owner,
)
.await
{
let owner_is_current = identity_lane
!= crate::authorization_runtime::transport::LegacyIdentityLane::Legacy
|| crate::api::relay_members::materialize_nip_oa_owner(
&state,
&conn.tenant,
&pubkey,
&owner,
)
.await;
if owner_is_current {
auth_ctx.agent_owner_pubkey = Some(owner);
} else {
warn!(
conn_id = %conn_id,
agent = %pubkey.to_hex(),
nip_oa_owner = %owner.to_hex(),
"NIP-OA owner could not be materialized"
);
}
}
info!(conn_id = %conn_id, pubkey = %pubkey.to_hex(), "NIP-42 auth successful");
info!(conn_id = %conn_id, "NIP-42 auth successful");
let transport_delegation =
crate::corporate_identity::verify_unconditional_nip_oa_owner(
pubkey,
auth_tag_json.as_deref(),
)
.map(|owner| {
VerifiedDelegationOutput::from_workspace_verifier(owner, pubkey, None, true)
});
let verified_proof: Arc<VerifiedNostrProof> = match VerifiedEvidenceAdapter::new()
.verify_nip42(
conn.tenant.community(),
AuthTransport::RelayWebSocket,
&verified_event,
&challenge,
&relay_url,
transport_delegation,
) {
Ok(proof) => Arc::new(proof),
Err(error) => {
warn!(conn_id = %conn_id, error = %error, "sealed NIP-42 evidence creation failed");
*conn.auth_state.write().await = AuthState::Failed;
conn.send(RelayMessage::ok(
&event_id_hex,
false,
"auth-required: verification failed",
));
return;
}
};
*conn.auth_state.write().await = AuthState::Authenticated(auth_ctx);
state
.conn_manager
.set_authenticated_pubkey(conn_id, pubkey.to_bytes().to_vec());
crate::corporate_identity::spawn_session_revalidation(
Arc::clone(&state),
conn.tenant.community(),
pubkey,
identity_decision,
conn.cancel.clone(),
state.conn_manager.set_authenticated_authority(
conn_id,
Arc::clone(&verified_proof),
verified_assertion.clone(),
);
if let (Some(runtime), Some(assertion)) =
(state.client_status_runtime().cloned(), verified_assertion)
{
if let Err(error) = runtime
.present_after_auth(
Arc::clone(&state),
verified_proof,
assertion,
conn_id,
conn.cancel.clone(),
)
.await
{
// Presentation failure never widens or narrows access. The
// client receives no current indicator and clears any old
// status on its existing freshness/disconnect boundary.
metrics::counter!("buzz_client_status_degradation_total").increment(1);
warn!(
conn_id = %conn_id,
reason = "client_status_unavailable",
"client binding status withheld"
);
tracing::debug!(error = %error, "client binding status detail");
}
}
if let Some(identity_decision) = identity_decision {
crate::corporate_identity::spawn_session_revalidation(
Arc::clone(&state),
conn.tenant.community(),
pubkey,
identity_decision,
conn.cancel.clone(),
);
}
conn.send(RelayMessage::ok(&event_id_hex, true, ""));
}
Err(e) => {
@@ -378,6 +485,48 @@ mod tests {
use super::extract_auth_tag_json;
use nostr::{EventBuilder, Keys, Kind, Tag};
#[test]
fn observational_auth_cannot_enter_mutating_identity_lane() {
use crate::authorization_runtime::{
finalization::AuthorizationMode,
transport::{legacy_identity_lane_for_mode, LegacyIdentityLane},
};
let mut binding_writes = 0;
let mut membership_writes = 0;
let mut public_projection_writes = 0;
for mode in [
AuthorizationMode::Shadow,
AuthorizationMode::VerifyOnly,
AuthorizationMode::Enforce,
] {
if legacy_identity_lane_for_mode(Some(mode)) == LegacyIdentityLane::Legacy {
binding_writes += 1;
membership_writes += 1;
public_projection_writes += 1;
}
}
assert_eq!(binding_writes, 0);
assert_eq!(membership_writes, 0);
assert_eq!(public_projection_writes, 0);
assert_eq!(
legacy_identity_lane_for_mode(Some(AuthorizationMode::Off)),
LegacyIdentityLane::Legacy
);
assert_eq!(
legacy_identity_lane_for_mode(Some(AuthorizationMode::Shadow)),
LegacyIdentityLane::ObserveOnly
);
assert_eq!(
legacy_identity_lane_for_mode(Some(AuthorizationMode::VerifyOnly)),
LegacyIdentityLane::ObserveOnly
);
assert_eq!(
legacy_identity_lane_for_mode(Some(AuthorizationMode::Enforce)),
LegacyIdentityLane::ProtectedEnforce
);
}
/// Build a signed NIP-98 (kind 27235) event carrying the given tags. The
/// `auth` tag lives inside the signed event exactly as the git and
/// WebSocket auth paths receive it.
+5
View File
@@ -4,6 +4,9 @@
mod admission;
/// Provider-neutral runtime authorization and bounded finalization.
pub mod authorization_runtime;
/// REST API route handlers.
pub mod api;
/// WebSocket audio relay for huddle voice channels.
@@ -31,6 +34,8 @@ pub mod mesh_boot;
pub mod metrics;
/// NIP-11 relay information document.
pub mod nip11;
/// Provider-neutral inventory of every protected relay surface.
pub mod protected_surface;
/// NIP-01 client/relay message parsing.
pub mod protocol;
/// Durable NIP-PL matcher and delivery worker.
+63 -7
View File
@@ -319,7 +319,7 @@ async fn main() -> anyhow::Result<()> {
(deployment_community, config.relay_owner_pubkey.as_ref())
{
match db.bootstrap_owner(community, owner_pubkey).await {
Ok(()) => info!(pubkey = %owner_pubkey, "Relay owner bootstrapped"),
Ok(()) => info!("Relay owner bootstrapped"),
Err(e) => {
if config.require_relay_membership {
// Membership enforcement is on — a missing owner means no one
@@ -427,7 +427,6 @@ async fn main() -> anyhow::Result<()> {
"0000000000000000000000000000000000000000000000000000000000000001";
let keys = nostr::Keys::parse(DEV_RELAY_PRIVKEY).expect("hardcoded dev key is valid");
tracing::warn!(
pubkey = %keys.public_key().to_hex(),
"Using hardcoded dev relay keypair (BUZZ_REQUIRE_AUTH_TOKEN=false). \
Set BUZZ_RELAY_PRIVATE_KEY for production."
);
@@ -461,6 +460,15 @@ async fn main() -> anyhow::Result<()> {
);
let state = Arc::new(app_state);
// Protected authorization is absent unless exact domains are named in
// server configuration. When present, durable invalidation snapshots are
// initialized before the runtime becomes reachable by any transport.
if buzz_relay::authorization_runtime::production::install_from_environment(&state).await?
== buzz_relay::authorization_runtime::production::ProtectedRuntimeInstallation::Installed
{
info!("Protected authorization runtime installed");
}
// Inter-relay mesh (BUZZ_MESH seam). `boot_mesh` returns None when the
// kill switch is off — nothing is bound, published, or spawned, so the
// relay behaves byte-identically to a build without the mesh. When
@@ -480,6 +488,8 @@ async fn main() -> anyhow::Result<()> {
// BUZZ_MESH_DEMO_ECHO) before peers can route traffic here.
handle.wire_consumers(
Arc::clone(&state.audio_rooms),
state.db.clone(),
state.relay_keypair.secret_key().as_secret_bytes(),
state.config.mesh_demo_echo,
Arc::clone(&state.shutting_down),
);
@@ -527,6 +537,41 @@ async fn main() -> anyhow::Result<()> {
);
}
// Enforce startup is verification-only for protected-object cutover.
// The resumable one-way preparation must complete before the independent
// restore anchor is provisioned; mutating PostgreSQL after anchor
// verification would create an unwitnessed authority advance.
if let Some(runtime) = state.protected_transport() {
let enforcing = runtime.enforcing_domains();
if !enforcing.is_empty() {
let hosts = state
.db
.usage_community_hosts()
.await?
.into_iter()
.map(|record| (buzz_core::CommunityId::from_uuid(record.id), record.host))
.collect::<std::collections::HashMap<_, _>>();
for community_id in enforcing {
let host = hosts.get(&community_id).ok_or_else(|| {
anyhow::anyhow!(
"protected object verification domain has no active community mapping"
)
})?;
let tenant = buzz_core::TenantContext::resolved(community_id, host);
let verification = async {
buzz_relay::api::git::migration::require_reconciled_authority(&state, &tenant)
.await?;
buzz_relay::api::media_migration::require_reconciled_authority(&state, &tenant)
.await?;
anyhow::Ok(())
};
tokio::time::timeout(std::time::Duration::from_secs(600), verification)
.await
.map_err(|_| anyhow::anyhow!("protected object verification timed out"))??;
}
}
}
// NIP-43: reconcile the event-backed roster for every provisioned
// community before opening the listener. `relay_members` is canonical;
// this repairs pre-snapshot communities and any publication that failed
@@ -616,8 +661,12 @@ async fn main() -> anyhow::Result<()> {
});
}
// Wire the action sink — must happen after AppState (which creates
// sub_registry, conn_manager) and before the cron loop starts.
// Wire the provider-neutral mutation gate and action sink after AppState
// construction and before any scheduled workflow can start.
let mutation_gate = Arc::new(buzz_relay::workflow_sink::RelayWorkflowMutationGate::new(
&state,
));
workflow_engine.set_mutation_gate(mutation_gate);
let action_sink = Arc::new(buzz_relay::workflow_sink::RelayActionSink::new(&state));
workflow_engine.set_action_sink(action_sink);
@@ -645,7 +694,12 @@ async fn main() -> anyhow::Result<()> {
loop {
tokio::time::sleep(std::time::Duration::from_secs(reaper_interval_secs)).await;
let expired = match reaper_state.db.reap_expired_ephemeral_channels().await {
let excluded = reaper_state.enforcing_protected_domain_ids();
let expired = match reaper_state
.db
.reap_expired_ephemeral_channels_excluding(&excluded)
.await
{
Ok(ids) => ids,
Err(e) => {
error!("Ephemeral reaper tick failed: {e}");
@@ -748,9 +802,10 @@ async fn main() -> anyhow::Result<()> {
tokio::time::sleep(std::time::Duration::from_secs(scheduler_interval_secs)).await;
let now_secs = chrono::Utc::now().timestamp();
let excluded = scheduler_state.enforcing_protected_domain_ids();
let due = match scheduler_state
.db
.query_due_reminders(now_secs, scheduler_batch_limit)
.query_due_reminders_excluding(now_secs, scheduler_batch_limit, &excluded)
.await
{
Ok(reminders) => reminders,
@@ -1506,9 +1561,10 @@ async fn run_usage_metrics_tick(
return Err(error);
}
let invite_retention_cutoff = chrono::Utc::now() - chrono::Duration::days(30);
let excluded = state.enforcing_protected_domain_ids();
match state
.db
.reap_expired_relay_invites(invite_retention_cutoff)
.reap_expired_relay_invites_excluding(invite_retention_cutoff, &excluded)
.await
{
Ok(deleted) if deleted > 0 => {
+40 -10
View File
@@ -157,6 +157,9 @@ pub struct MeshHandle {
///
/// [`MeshAudioRouter`]: crate::audio::mesh::MeshAudioRouter
pub audio_fence: Arc<crate::audio::mesh::GenerationFloor>,
/// Live reliable-control attachments accepted by the realtime media lane.
/// Datagrams cannot create entries in this registry.
pub audio_attachments: Arc<crate::audio::mesh::MediaAttachmentRegistry>,
/// The running mesh (status snapshots, shutdown).
runtime: MeshRuntime,
/// Per-room huddle owner-lease coordination. Shared with the WS-join owner
@@ -180,6 +183,8 @@ impl MeshHandle {
pub fn wire_consumers(
&self,
rooms: Arc<crate::audio::AudioRoomManager>,
db: buzz_db::Db,
relay_secret: &[u8],
demo_echo: bool,
shutting_down: Arc<AtomicBool>,
) {
@@ -189,8 +194,15 @@ impl MeshHandle {
Arc::clone(&self.transport),
self.local_runtime_id,
Arc::clone(&self.audio_fence),
Arc::clone(&self.audio_attachments),
rooms,
Arc::clone(&self.owners),
Some(
crate::authorization_runtime::ephemeral::AuthorityTokenVerifier::new(
db,
relay_secret,
),
),
demo_echo,
shutting_down,
)
@@ -221,14 +233,16 @@ impl MeshHandle {
/// [`HuddleControlAcceptor::accept_inbound`]: crate::audio::join::HuddleControlAcceptor::accept_inbound
/// [`ReliableJoin::Owned`]: crate::tunnel::reliable::ReliableJoin::Owned
#[allow(clippy::too_many_arguments)] // boot-only parts bundle, one caller + tests
pub fn wire_mesh_consumers(
pub(crate) fn wire_mesh_consumers(
dispatcher: &MeshInboundDispatcher,
directory: SessionDirectory,
transport: Arc<dyn RelayPeerTransport>,
local_runtime_id: RuntimeId,
audio_fence: Arc<crate::audio::mesh::GenerationFloor>,
audio_attachments: Arc<crate::audio::mesh::MediaAttachmentRegistry>,
rooms: Arc<crate::audio::AudioRoomManager>,
owners: Arc<crate::audio::join::HuddleOwnerRegistry>,
authority_verifier: Option<crate::authorization_runtime::ephemeral::AuthorityTokenVerifier>,
demo_echo: bool,
shutting_down: Arc<AtomicBool>,
) {
@@ -239,21 +253,29 @@ pub fn wire_mesh_consumers(
Arc::clone(&rooms),
local_runtime_id,
audio_fence,
Arc::clone(&audio_attachments),
);
dispatcher.register_datagrams(Box::new(move |_from, dgram| {
audio_router.on_media_datagram(&dgram);
dispatcher.register_datagrams(Box::new(move |from, dgram| {
audio_router.on_media_datagram(from, &dgram);
}));
// HuddleControl streams: owner-side peer registration for cross-pod
// huddles. The acceptor validates structurally, then Redis-fences every
// stateful frame in its control loop.
let acceptor = Arc::new(crate::audio::join::HuddleControlAcceptor::new(
let acceptor = crate::audio::join::HuddleControlAcceptor::new(
rooms,
Arc::clone(&transport),
Arc::new(directory.clone()),
local_runtime_id,
Arc::clone(&owners),
));
audio_attachments,
);
let acceptor = if let Some(verifier) = authority_verifier {
acceptor.with_authority_verifier(verifier)
} else {
acceptor
};
let acceptor = Arc::new(acceptor);
dispatcher.register_huddle_control(Box::new(move |from, hello, stream| {
let acceptor = Arc::clone(&acceptor);
tokio::spawn(async move {
@@ -486,6 +508,7 @@ pub async fn boot_mesh(
let runtime = MeshRuntime::start(endpoint, membership, Some(registry));
let owners = Arc::new(crate::audio::join::HuddleOwnerRegistry::new());
let audio_attachments = Arc::new(crate::audio::mesh::MediaAttachmentRegistry::default());
// Dial seed peers now rather than waiting for the first reconcile tick.
runtime.reconcile_now().await;
@@ -528,6 +551,7 @@ pub async fn boot_mesh(
local_runtime_id: runtime_id,
dispatcher,
audio_fence: Arc::new(crate::audio::mesh::GenerationFloor::new()),
audio_attachments,
runtime,
owners,
}))
@@ -719,6 +743,7 @@ mod tests {
let dispatcher = MeshInboundDispatcher::default();
let fence = Arc::new(crate::audio::mesh::GenerationFloor::new());
let attachments = Arc::new(crate::audio::mesh::MediaAttachmentRegistry::default());
let pool = deadpool_redis::Config::from_url("redis://127.0.0.1:1") // never dialed
.create_pool(Some(deadpool_redis::Runtime::Tokio1))
.unwrap();
@@ -728,25 +753,30 @@ mod tests {
Arc::new(NoopTransport),
rid(9),
Arc::clone(&fence),
Arc::clone(&attachments),
Arc::new(crate::audio::AudioRoomManager::new()),
Arc::new(crate::audio::join::HuddleOwnerRegistry::new()),
None,
false,
Arc::new(AtomicBool::new(false)),
);
let session = uuid::Uuid::new_v4();
let fenced = FencedHeader {
session_id: session,
generation: 7,
owner_runtime_id: rid(1),
};
let _attachment = attachments.register_owner_fanout(fenced, uuid::Uuid::new_v4(), u64::MAX);
dispatcher.on_datagram(
rid(1),
MeshDatagram {
fenced: FencedHeader {
session_id: session,
generation: 7,
owner_runtime_id: rid(9),
},
fenced,
seq: 0,
payload: vec![0, 1, 2],
},
);
tokio::task::yield_now().await;
// The shared fence observed the datagram's generation: a stale check
// through the HANDLE's Arc is rejected, proving one floor, not two.
+332
View File
@@ -1,6 +1,9 @@
//! NIP-11 relay information document.
use std::num::NonZeroU64;
use serde::{Deserialize, Serialize};
use thiserror::Error;
#[cfg(test)]
use crate::config::DEFAULT_MAX_FRAME_BYTES;
@@ -20,6 +23,143 @@ pub(crate) const SUPPORTED_NIPS: &[u32] = &[1, 2, 10, 11, 16, 17, 23, 25, 29, 33
/// to be verifiable by clients.
pub(crate) const NIP_RELAY_MEMBERSHIP: u32 = 43;
/// Provider-neutral NIP-FI assertion transport profile.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum NipFiTransportProfile {
/// Assertions are injected only by an origin-isolated trusted proxy that
/// strips untrusted inbound copies of the configured assertion header.
TrustedProxy,
/// Assertions are attached by the client to the same protected HTTP
/// request as its NIP-98 proof.
ClientAttached,
}
/// Provider-neutral NIP-FI enrollment mode.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum NipFiEnrollmentMode {
/// First enrollment requires an assertion key attestation.
AttestedKey,
/// Binding creation requires a separate privileged transition.
Provisioned,
/// First valid use may create the binding under explicit TOFU policy.
Tofu,
}
/// Provider-neutral NIP-FI discovery object.
///
/// Construction makes the delegation bound invariant unrepresentable:
/// delegation is `true` exactly when a positive finite maximum is present.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct NipFiDiscovery {
transports: Vec<NipFiTransportProfile>,
enrollment: NipFiEnrollmentMode,
delegation: bool,
#[serde(skip_serializing_if = "Option::is_none")]
delegated_lease_max_seconds: Option<NonZeroU64>,
}
impl NipFiDiscovery {
/// Validate provider-neutral discovery configuration.
pub fn new(
mut transports: Vec<NipFiTransportProfile>,
enrollment: NipFiEnrollmentMode,
delegated_lease_max_seconds: Option<NonZeroU64>,
) -> Result<Self, NipFiDiscoveryError> {
if transports.is_empty() {
return Err(NipFiDiscoveryError::NoTransport);
}
transports.sort_unstable();
let original_len = transports.len();
transports.dedup();
if transports.len() != original_len {
return Err(NipFiDiscoveryError::DuplicateTransport);
}
Ok(Self {
transports,
enrollment,
delegation: delegated_lease_max_seconds.is_some(),
delegated_lease_max_seconds,
})
}
fn includes_transport(&self, transport: NipFiTransportProfile) -> bool {
self.transports.contains(&transport)
}
}
/// Complete-stack conformance input supplied by the release/conformance lane.
///
/// A source may return `true` only when every applicable NIP-FI row passed
/// against the same reviewed implementation revision. Trusted-proxy support
/// additionally requires deployment evidence for origin isolation and inbound
/// header stripping; synthetic code tests alone are insufficient.
pub trait CompleteNipFiRuntimeConformance: Send + Sync {
/// Exact reviewed implementation revision used for every applicable row.
fn reviewed_implementation_revision(&self) -> &str;
/// Whether every applicable row passed at the reviewed revision.
fn all_applicable_rows_passed_at_same_revision(&self) -> bool;
/// Whether trusted-proxy deployment controls and negative tests passed.
fn trusted_proxy_deployment_evidence_passed(&self) -> bool;
}
/// Discovery proven ready by an injected complete-stack conformance source.
///
/// The reviewed revision and evidence are deliberately not serialized into
/// NIP-11. This wrapper has no public field or unchecked constructor.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConformanceReadyNipFiDiscovery(NipFiDiscovery);
impl ConformanceReadyNipFiDiscovery {
/// Gate discovery on complete same-revision runtime and deployment proof.
pub fn from_complete_stack(
discovery: NipFiDiscovery,
conformance: &dyn CompleteNipFiRuntimeConformance,
) -> Result<Self, NipFiDiscoveryError> {
let revision = conformance.reviewed_implementation_revision();
if revision.len() != 40
|| !revision
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
{
return Err(NipFiDiscoveryError::InvalidReviewedRevision);
}
if !conformance.all_applicable_rows_passed_at_same_revision() {
return Err(NipFiDiscoveryError::IncompleteConformance);
}
if discovery.includes_transport(NipFiTransportProfile::TrustedProxy)
&& !conformance.trusted_proxy_deployment_evidence_passed()
{
return Err(NipFiDiscoveryError::MissingTrustedProxyEvidence);
}
Ok(Self(discovery))
}
}
/// Fail-closed NIP-FI discovery construction error.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
pub enum NipFiDiscoveryError {
/// At least one supported transport must be advertised.
#[error("NIP-FI discovery requires at least one transport")]
NoTransport,
/// Each supported transport may appear only once.
#[error("NIP-FI discovery contains a duplicate transport")]
DuplicateTransport,
/// The complete-stack report did not identify one exact Git revision.
#[error("NIP-FI conformance report has an invalid reviewed revision")]
InvalidReviewedRevision,
/// Not every applicable row passed at the same revision.
#[error("NIP-FI complete-stack conformance is incomplete")]
IncompleteConformance,
/// Trusted-proxy origin isolation and header stripping were not proven.
#[error("NIP-FI trusted-proxy deployment evidence is incomplete")]
MissingTrustedProxyEvidence,
}
/// Relay information document served at `GET /` with `Accept: application/nostr+json`.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RelayInfo {
@@ -55,6 +195,10 @@ pub struct RelayInfo {
/// Relay's own signing pubkey (NIP-11 `self` field, NIP-43).
#[serde(rename = "self", skip_serializing_if = "Option::is_none")]
pub relay_self: Option<String>,
/// Provider-neutral NIP-FI capabilities. Omitted until a complete-stack
/// same-revision conformance input explicitly enables discovery.
#[serde(skip_serializing_if = "Option::is_none")]
federated_identity: Option<NipFiDiscovery>,
}
/// Protocol and resource limits advertised in the NIP-11 document.
@@ -85,6 +229,9 @@ pub struct RelayLimitation {
/// NIP-ER: maximum allowed `not_before` horizon in seconds from now.
#[serde(skip_serializing_if = "Option::is_none")]
pub max_not_before_delta: Option<u64>,
/// NIP-FI support. Omitted until complete-stack conformance is proven.
#[serde(skip_serializing_if = "Option::is_none")]
federated_identity: Option<bool>,
}
/// Canonical `RelayLimitation` advertised by this relay.
@@ -116,6 +263,7 @@ fn relay_limitation(max_message_length: usize) -> RelayLimitation {
restricted_writes: true,
due_delivery_mode: Some("push".to_string()),
max_not_before_delta: Some(max_not_before_delta),
federated_identity: None,
}
}
@@ -169,8 +317,25 @@ impl RelayInfo {
limitation: Some(relay_limitation(max_message_length)),
pairing_relay_url: pairing_relay_url.map(str::to_string),
relay_self: relay_self.map(|s| s.to_string()),
federated_identity: None,
}
}
/// Add provider-neutral NIP-FI discovery after complete-stack proof.
///
/// There is intentionally no raw boolean/configuration overload. The
/// normal runtime build path has no readiness input and therefore remains
/// silent until the release/conformance lane supplies this gated value.
pub fn with_conformant_federated_identity(
mut self,
ready: ConformanceReadyNipFiDiscovery,
) -> Self {
if let Some(limitation) = &mut self.limitation {
limitation.federated_identity = Some(true);
}
self.federated_identity = Some(ready.0);
self
}
}
/// Axum handler that returns the NIP-11 relay information document as JSON.
@@ -267,6 +432,9 @@ pub(crate) async fn nip11_document(state: &crate::state::AppState, raw_host: &st
.push("nip-pl".to_string());
info.push = Some(push);
}
if let Some(ready) = state.nip_fi_discovery().cloned() {
info = info.with_conformant_federated_identity(ready);
}
info
}
@@ -350,6 +518,34 @@ const _RELAY_INFO_BUILD_STATIC_INPUT_FENCE: fn(
mod tests {
use super::*;
struct SyntheticConformance {
revision: &'static str,
complete: bool,
trusted_proxy_evidence: bool,
}
impl CompleteNipFiRuntimeConformance for SyntheticConformance {
fn reviewed_implementation_revision(&self) -> &str {
self.revision
}
fn all_applicable_rows_passed_at_same_revision(&self) -> bool {
self.complete
}
fn trusted_proxy_deployment_evidence_passed(&self) -> bool {
self.trusted_proxy_evidence
}
}
fn complete_conformance() -> SyntheticConformance {
SyntheticConformance {
revision: "0123456789abcdef0123456789abcdef01234567",
complete: true,
trusted_proxy_evidence: true,
}
}
#[test]
fn push_descriptor_is_gated_by_gateway_configuration_and_tenant_binding() {
let keys = nostr::Keys::generate();
@@ -402,6 +598,142 @@ mod tests {
assert_eq!(info.software, "https://github.com/block/buzz");
}
#[test]
fn default_discovery_is_silent_until_complete_stack_input_exists() {
let info = RelayInfo::build(None, None, false, DEFAULT_MAX_FRAME_BYTES, None);
let json = serde_json::to_value(info).expect("serialize default NIP-11");
assert!(json.get("federated_identity").is_none());
assert!(json["limitation"].get("federated_identity").is_none());
}
#[test]
fn discovery_rejects_empty_duplicate_and_incomplete_inputs() {
assert_eq!(
NipFiDiscovery::new(Vec::new(), NipFiEnrollmentMode::AttestedKey, None),
Err(NipFiDiscoveryError::NoTransport)
);
assert_eq!(
NipFiDiscovery::new(
vec![
NipFiTransportProfile::ClientAttached,
NipFiTransportProfile::ClientAttached,
],
NipFiEnrollmentMode::Provisioned,
None,
),
Err(NipFiDiscoveryError::DuplicateTransport)
);
let discovery = NipFiDiscovery::new(
vec![NipFiTransportProfile::ClientAttached],
NipFiEnrollmentMode::Provisioned,
None,
)
.expect("synthetic discovery is valid");
for conformance in [
SyntheticConformance {
revision: "not-a-revision",
complete: true,
trusted_proxy_evidence: true,
},
SyntheticConformance {
revision: "0123456789abcdef0123456789abcdef01234567",
complete: false,
trusted_proxy_evidence: true,
},
] {
assert!(ConformanceReadyNipFiDiscovery::from_complete_stack(
discovery.clone(),
&conformance,
)
.is_err());
}
}
#[test]
fn trusted_proxy_advertisement_requires_deployment_evidence() {
let discovery = NipFiDiscovery::new(
vec![NipFiTransportProfile::TrustedProxy],
NipFiEnrollmentMode::AttestedKey,
None,
)
.expect("synthetic discovery is valid");
let conformance = SyntheticConformance {
revision: "0123456789abcdef0123456789abcdef01234567",
complete: true,
trusted_proxy_evidence: false,
};
assert_eq!(
ConformanceReadyNipFiDiscovery::from_complete_stack(discovery, &conformance),
Err(NipFiDiscoveryError::MissingTrustedProxyEvidence)
);
}
#[test]
fn conformant_discovery_is_provider_neutral_and_delegation_bounded() {
let max = NonZeroU64::new(300).expect("synthetic bound is positive");
let discovery = NipFiDiscovery::new(
vec![
NipFiTransportProfile::TrustedProxy,
NipFiTransportProfile::ClientAttached,
],
NipFiEnrollmentMode::AttestedKey,
Some(max),
)
.expect("synthetic discovery is valid");
let ready =
ConformanceReadyNipFiDiscovery::from_complete_stack(discovery, &complete_conformance())
.expect("complete synthetic report enables discovery");
let info = RelayInfo::build(None, None, false, DEFAULT_MAX_FRAME_BYTES, None)
.with_conformant_federated_identity(ready);
let json = serde_json::to_value(info).expect("serialize conformant discovery");
assert_eq!(json["limitation"]["federated_identity"], true);
assert_eq!(
json["federated_identity"],
serde_json::json!({
"transports": ["trusted-proxy", "client-attached"],
"enrollment": "attested-key",
"delegation": true,
"delegated_lease_max_seconds": 300,
})
);
let encoded = json.to_string();
for private in [
"synthetic-issuer",
"synthetic-subject",
"tenant.example",
"private-audience",
"assertion-header-name",
] {
assert!(!encoded.contains(private));
}
}
#[test]
fn discovery_without_delegation_omits_lease_bound() {
let discovery = NipFiDiscovery::new(
vec![NipFiTransportProfile::ClientAttached],
NipFiEnrollmentMode::Tofu,
None,
)
.expect("synthetic discovery is valid");
let ready =
ConformanceReadyNipFiDiscovery::from_complete_stack(discovery, &complete_conformance())
.expect("complete synthetic report enables discovery");
let info = RelayInfo::build(None, None, false, DEFAULT_MAX_FRAME_BYTES, None)
.with_conformant_federated_identity(ready);
let json = serde_json::to_value(info).expect("serialize conformant discovery");
assert_eq!(json["federated_identity"]["delegation"], false);
assert!(json["federated_identity"]
.get("delegated_lease_max_seconds")
.is_none());
assert_eq!(json["supported_nips"], serde_json::json!(SUPPORTED_NIPS));
}
#[test]
fn configured_pairing_relay_is_advertised_and_unset_value_is_omitted() {
let info = RelayInfo::build(
+13 -2
View File
@@ -58,8 +58,13 @@ pub async fn run_matcher(state: Arc<AppState>) {
let mut idle_delay = IDLE_POLL_FLOOR;
let mut last_reap = tokio::time::Instant::now();
loop {
let excluded = state.enforcing_protected_domain_ids();
if last_reap.elapsed() >= REAP_INTERVAL {
match state.db.reap_exhausted_push_matches().await {
match state
.db
.reap_exhausted_push_matches_excluding(&excluded)
.await
{
Ok(reaped) if reaped > 0 => warn!(reaped, "reaped exhausted push match jobs"),
Ok(_) => {}
Err(e) => error!("push match reap failed: {e}"),
@@ -69,7 +74,7 @@ pub async fn run_matcher(state: Arc<AppState>) {
let until = Utc::now() + TimeDelta::seconds(CLAIM_SECS);
match state
.db
.claim_due_push_match_batch(MATCH_BATCH_LIMIT, until)
.claim_due_push_match_batch_excluding(MATCH_BATCH_LIMIT, until, &excluded)
.await
{
Ok(Some(batch)) => {
@@ -321,6 +326,9 @@ pub async fn run_delivery_worker(state: Arc<AppState>) {
Ok(communities) => {
for community in communities {
let community = buzz_core::CommunityId::from_uuid(community.id);
if state.is_protected_enforcing(community) {
continue;
}
let until = Utc::now() + TimeDelta::seconds(CLAIM_SECS);
match state.db.claim_due_push_wakes(community, 16, until).await {
Ok(wakes) => {
@@ -351,6 +359,9 @@ async fn deliver_one(
http: &reqwest::Client,
claimed: buzz_db::push::ClaimedWake,
) {
if state.is_protected_enforcing(claimed.community) {
return;
}
let outcome = match state
.db
.revalidate_push_wake(claimed.community, claimed.id, claimed.claim_id)
+118
View File
@@ -551,6 +551,124 @@ CREATE INDEX idx_identity_lifecycle_operations_principal
CREATE INDEX idx_identity_lifecycle_operations_key
ON identity_lifecycle_operations (community_id, pubkey, created_at);
-- ── Authorization invalidation authority ──────────────────────────────────────
-- Generations and selector floors are durable authority. Cross-node pub/sub
-- carries only a hint that consumers should reconcile from these tables.
CREATE TABLE authorization_invalidation_domains (
community_id UUID NOT NULL REFERENCES communities(id),
generation BIGINT NOT NULL DEFAULT 0 CHECK (generation >= 0),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (community_id)
);
CREATE TABLE authorization_invalidation_receipts (
community_id UUID NOT NULL REFERENCES communities(id),
event_id UUID NOT NULL,
generation BIGINT NOT NULL CHECK (generation > 0),
request_fingerprint BYTEA NOT NULL CHECK (length(request_fingerprint) = 32),
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (community_id, event_id),
UNIQUE (community_id, generation)
);
CREATE TABLE authorization_invalidation_floors (
community_id UUID NOT NULL REFERENCES communities(id),
selector_kind TEXT NOT NULL CHECK (selector_kind IN (
'principal_fingerprint',
'nostr_key',
'binding',
'session',
'domain',
'policy_version',
'delegated_owner'
)),
selector_fingerprint BYTEA NOT NULL CHECK (length(selector_fingerprint) = 32),
generation BIGINT NOT NULL CHECK (generation > 0),
sticky_deny BOOLEAN NOT NULL DEFAULT FALSE,
binding_version_floor BIGINT CHECK (binding_version_floor > 0),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (community_id, selector_kind, selector_fingerprint),
FOREIGN KEY (community_id, generation)
REFERENCES authorization_invalidation_receipts (community_id, generation),
CHECK ((selector_kind = 'binding') = (binding_version_floor IS NOT NULL))
);
CREATE INDEX idx_authorization_invalidation_floors_generation
ON authorization_invalidation_floors (community_id, generation);
-- Transaction-owned protected-operation idempotency. This is commit protocol
-- state, not an authorization decision or operator audit log.
CREATE TABLE authorization_operation_receipts (
community_id UUID NOT NULL REFERENCES communities(id),
operation_id UUID NOT NULL,
operation_kind TEXT NOT NULL CHECK (
length(operation_kind) > 0 AND length(operation_kind) <= 128
),
request_fingerprint BYTEA NOT NULL CHECK (length(request_fingerprint) = 32),
result_version SMALLINT NOT NULL DEFAULT 1 CHECK (result_version > 0),
result_payload BYTEA NOT NULL CHECK (octet_length(result_payload) <= 65536),
lease_expires_at TIMESTAMPTZ NOT NULL,
committed_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY (community_id, operation_id)
);
CREATE INDEX idx_authorization_operation_receipts_committed_at
ON authorization_operation_receipts (community_id, committed_at);
CREATE FUNCTION authorization_operation_expiry_guard() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
IF NEW.lease_expires_at <= clock_timestamp() THEN
RAISE EXCEPTION 'protected operation authorization expired before commit'
USING ERRCODE = 'check_violation';
END IF;
RETURN NULL;
END
$$;
CREATE CONSTRAINT TRIGGER authorization_operation_expiry
AFTER INSERT OR UPDATE OF lease_expires_at
ON authorization_operation_receipts
DEFERRABLE INITIALLY DEFERRED
FOR EACH ROW
EXECUTE FUNCTION authorization_operation_expiry_guard();
CREATE TABLE git_repo_publications (
community_id UUID NOT NULL,
repo_id TEXT NOT NULL,
owner_pubkey TEXT NOT NULL,
manifest_sha256 TEXT NOT NULL CHECK (manifest_sha256 ~ '^[0-9a-f]{64}$'),
publication_version BIGINT NOT NULL CHECK (publication_version > 0),
state TEXT NOT NULL DEFAULT 'active' CHECK (state IN ('active', 'unpublished')),
created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY (community_id, repo_id),
FOREIGN KEY (community_id, repo_id)
REFERENCES git_repo_names (community_id, repo_id)
);
CREATE TABLE media_publications (
community_id UUID NOT NULL REFERENCES communities(id),
sha256 TEXT NOT NULL CHECK (sha256 ~ '^[0-9a-f]{64}$'),
object_key TEXT NOT NULL CHECK (length(object_key) > 0 AND length(object_key) <= 512),
extension TEXT NOT NULL CHECK (extension ~ '^[a-z0-9]{1,8}$'),
mime_type TEXT NOT NULL CHECK (length(mime_type) > 0 AND length(mime_type) <= 255),
object_size BIGINT NOT NULL CHECK (object_size >= 0),
metadata JSONB NOT NULL CHECK (octet_length(metadata::text) <= 16384),
thumbnail_key TEXT CHECK (
thumbnail_key IS NULL OR (length(thumbnail_key) > 0 AND length(thumbnail_key) <= 512)
),
publication_version BIGINT NOT NULL DEFAULT 1 CHECK (publication_version > 0),
state TEXT NOT NULL DEFAULT 'active' CHECK (state IN ('active', 'unpublished')),
created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY (community_id, sha256)
);
CREATE INDEX idx_media_publications_state
ON media_publications (community_id, state);
-- ── Events (partitioned by month on created_at) ──────────────────────────────
-- Conformance: "Channel-less global events and DMs". `community_id` leads the
-- PK and every hot-path index. Partition stays BY RANGE (created_at) — the