diff --git a/crates/buzz-auth/src/evidence_adapter.rs b/crates/buzz-auth/src/evidence_adapter.rs new file mode 100644 index 000000000..5c032b6e3 --- /dev/null +++ b/crates/buzz-auth/src/evidence_adapter.rs @@ -0,0 +1,640 @@ +//! Narrow trusted-workspace adapter for sealed authorization evidence. +//! +//! Wire data cannot deserialize into any type in this module. The adapter is +//! used only after the existing cryptographic verifier, assertion verifier, +//! community policy, or typed binding store has returned success. It repeats +//! exact domain, transport, actor, time, identifier, version, and provenance +//! checks before crossing the crate's sealed evidence boundary. + +use buzz_core::{tenant::TenantContext, CommunityId}; +use nostr::{Event, PublicKey}; +use sha2::{Digest, Sha256}; +use thiserror::Error; +use uuid::Uuid; + +use crate::{ + context::{ + AssertionExpiry, AssertionNotBefore, AssertionTransport, AuthContextError, AuthMethod, + AuthTransport, AuthorizedCommunityAccess, BindingExpiry, BindingSource, BindingVersion, + DelegationExpiry, FederatedPrincipal, VerifiedFederatedAssertion, VerifiedKeyAttestation, + VerifiedNostrProof, VerifiedOperationBinding, VerifiedOperationBindingKind, + VerifiedTransportDelegation, VersionedBindingRef, + }, + nip42::verify_nip42_event, + nip98::verify_nip98_event, + AuthorizationReason, Scope, +}; + +/// Binding lifecycle result returned by an authoritative store adapter. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ActiveBindingResolution { + /// Exact active binding already existed. + Existing, + /// Binding was atomically enrolled during this decision. + Enrolled, +} + +/// Existing-verifier output for transport-wide Nostr delegation. +/// +/// Construction is deliberately explicit and non-serializable. The relay may +/// create it only by translating a successful NIP-OA verifier result whose +/// constraints were proved transport-wide; raw `auth` tag data is not such a +/// result. +pub struct VerifiedDelegationOutput { + owner_pubkey: PublicKey, + delegate_pubkey: PublicKey, + expires_at: Option, + transport_wide: bool, +} + +impl VerifiedDelegationOutput { + /// Translate an existing verifier's exact output. + pub const fn from_workspace_verifier( + owner_pubkey: PublicKey, + delegate_pubkey: PublicKey, + expires_at: Option, + transport_wide: bool, + ) -> Self { + Self { + owner_pubkey, + delegate_pubkey, + expires_at, + transport_wide, + } + } +} + +/// Stateless factory for sealed evidence. +/// +/// This value carries no configuration and cannot select a domain, provider, +/// or capability. Every method requires the server-resolved values again and +/// rejects mismatches in the translated verifier/store output. +#[derive(Debug, Default, Clone, Copy)] +pub struct VerifiedEvidenceAdapter; + +fn operation_binding( + kind: VerifiedOperationBindingKind, + parts: &[&[u8]], +) -> VerifiedOperationBinding { + let mut hasher = Sha256::new(); + hasher.update(b"buzz-auth:verified-operation-binding:v1"); + for part in parts { + hasher.update((part.len() as u64).to_be_bytes()); + hasher.update(part); + } + VerifiedOperationBinding::from_evidence_adapter(kind, hasher.finalize().into()) +} + +impl VerifiedEvidenceAdapter { + /// Create the trusted-workspace adapter. + pub const fn new() -> Self { + Self + } + + /// Attach exact transport-wide delegation output to an already sealed + /// cryptographic proof. + /// + /// This consumes the original proof, rejects replacement of an existing + /// delegation, and rechecks that the verified delegate is the proof actor. + pub fn attach_transport_delegation( + &self, + proof: VerifiedNostrProof, + delegation: VerifiedDelegationOutput, + ) -> Result { + if proof.verified_delegation().is_some() { + return Err(EvidenceAdapterError::DelegationAlreadyPresent); + } + let verified_delegation = self.delegation(proof.actor_pubkey(), delegation)?; + VerifiedNostrProof::from_evidence_adapter( + proof.authorization_domain(), + proof.authorized_transport(), + proof.actor_pubkey(), + proof.proof_method(), + proof.operation_binding(), + Some(verified_delegation), + ) + .map_err(Into::into) + } + + /// Verify NIP-42 and bind its result to relay or audio transport. + pub fn verify_nip42( + &self, + authorization_domain: CommunityId, + transport: AuthTransport, + event: &Event, + expected_challenge: &str, + relay_url: &str, + delegation: Option, + ) -> Result { + if !matches!( + transport, + AuthTransport::RelayWebSocket | AuthTransport::Audio + ) { + return Err(EvidenceAdapterError::TransportMethodMismatch); + } + verify_nip42_event(event, expected_challenge, relay_url)?; + let delegation = delegation + .map(|delegation| self.delegation(event.pubkey, delegation)) + .transpose()?; + let binding = operation_binding( + VerifiedOperationBindingKind::NostrSession, + &[expected_challenge.as_bytes(), relay_url.as_bytes()], + ); + VerifiedNostrProof::from_evidence_adapter( + authorization_domain, + transport, + event.pubkey, + AuthMethod::Nip42, + binding, + delegation, + ) + .map_err(Into::into) + } + + /// Verify NIP-98 and bind it to one exact HTTP transport operation. + // The individual arguments are intentional trust-boundary inputs: folding + // them into an unverified request bag would make it easier to omit an exact + // domain, transport, method, body, or delegation cross-check. + #[allow(clippy::too_many_arguments)] + pub fn verify_nip98( + &self, + authorization_domain: CommunityId, + transport: AuthTransport, + event_json: &str, + expected_url: &str, + expected_method: &str, + body: Option<&[u8]>, + delegation: Option, + ) -> Result { + if !matches!( + transport, + AuthTransport::HttpBridge | AuthTransport::Git | AuthTransport::MediaUpload + ) { + return Err(EvidenceAdapterError::TransportMethodMismatch); + } + let actor = verify_nip98_event(event_json, expected_url, expected_method, body)?; + let delegation = delegation + .map(|delegation| self.delegation(actor, delegation)) + .transpose()?; + let body_presence = [u8::from(body.is_some())]; + let binding = operation_binding( + VerifiedOperationBindingKind::HttpRequest, + &[ + expected_method.as_bytes(), + expected_url.as_bytes(), + &body_presence, + body.unwrap_or_default(), + ], + ); + VerifiedNostrProof::from_evidence_adapter( + authorization_domain, + transport, + actor, + AuthMethod::Nip98, + binding, + delegation, + ) + .map_err(Into::into) + } + + /// Fully verify and bind one exact Blossom upload operation. + pub fn verify_blossom_upload( + &self, + authorization_domain: CommunityId, + event: &Event, + sha256: &str, + server_domain: Option<&str>, + max_age_secs: u64, + ) -> Result { + crate::blossom::verify_blossom_upload_auth(event, sha256, server_domain, max_age_secs)?; + let event_id = event.id.to_bytes(); + let max_age = max_age_secs.to_be_bytes(); + let binding = operation_binding( + VerifiedOperationBindingKind::BlossomUpload, + &[ + &event_id, + sha256.as_bytes(), + server_domain.unwrap_or_default().as_bytes(), + &max_age, + ], + ); + self.blossom_proof( + authorization_domain, + AuthTransport::MediaUpload, + event, + binding, + ) + } + + /// Fully verify and bind one exact Blossom GET or HEAD operation. + pub fn verify_blossom_download( + &self, + authorization_domain: CommunityId, + event: &Event, + sha256: &str, + server_domain: Option<&str>, + max_age_secs: u64, + ) -> Result { + crate::blossom::verify_blossom_get_auth(event, sha256, server_domain, max_age_secs)?; + let event_id = event.id.to_bytes(); + let max_age = max_age_secs.to_be_bytes(); + let binding = operation_binding( + VerifiedOperationBindingKind::BlossomDownload, + &[ + &event_id, + sha256.as_bytes(), + server_domain.unwrap_or_default().as_bytes(), + &max_age, + ], + ); + self.blossom_proof( + authorization_domain, + AuthTransport::MediaDownload, + event, + binding, + ) + } + + fn blossom_proof( + &self, + authorization_domain: CommunityId, + transport: AuthTransport, + event: &Event, + binding: VerifiedOperationBinding, + ) -> Result { + VerifiedNostrProof::from_evidence_adapter( + authorization_domain, + transport, + event.pubkey, + AuthMethod::Blossom, + binding, + None, + ) + .map_err(Into::into) + } + + fn delegation( + &self, + actor: PublicKey, + output: VerifiedDelegationOutput, + ) -> Result { + if !output.transport_wide { + return Err(EvidenceAdapterError::NarrowDelegation); + } + if output.delegate_pubkey != actor { + return Err(EvidenceAdapterError::DelegatedActorMismatch); + } + let expires_at = output.expires_at.map(DelegationExpiry::new).transpose()?; + VerifiedTransportDelegation::new_unrestricted( + output.owner_pubkey, + output.delegate_pubkey, + expires_at, + ) + .map_err(Into::into) + } + + /// Translate validated assertion claims into sealed, exact-bound evidence. + /// + /// `now_unix_seconds` and every claim argument must be copied from the + /// successful configured assertion-verifier result, never decoded again + /// from request data at this boundary. + #[allow(clippy::too_many_arguments)] + pub fn federated_assertion_from_validated_claims( + &self, + authorization_domain: CommunityId, + transport: AuthTransport, + issuer: &str, + subject: &str, + attested_pubkey: Option, + assertion_transport: AssertionTransport, + not_before: Option, + expires_at: u64, + now_unix_seconds: u64, + ) -> Result { + let principal = FederatedPrincipal::new(issuer, subject)?; + let expires_at = AssertionExpiry::new(expires_at)?; + if expires_at.is_expired_at(now_unix_seconds) { + return Err(EvidenceAdapterError::AssertionExpired); + } + let not_before = not_before.map(AssertionNotBefore::new); + if not_before.is_some_and(|bound| bound.is_not_yet_valid_at(now_unix_seconds)) { + return Err(EvidenceAdapterError::AssertionNotYetValid); + } + let key_attestation = attested_pubkey.map(VerifiedKeyAttestation::from_evidence_adapter); + Ok(VerifiedFederatedAssertion::from_evidence_adapter( + authorization_domain, + transport, + principal, + key_attestation, + assertion_transport, + not_before, + expires_at, + )) + } + + /// Translate a typed active binding-store result. + /// + /// The adapter derives the authorization reason from lifecycle resolution + /// and provenance. Callers cannot label an enrolled row as pre-existing or + /// turn a provisioned row into a first-use enrollment. + #[allow(clippy::too_many_arguments)] + pub fn active_binding_from_store( + &self, + authorization_domain: CommunityId, + binding_domain: CommunityId, + binding_id: Uuid, + issuer: &str, + subject: &str, + bound_pubkey: PublicKey, + binding_version: u64, + expires_at: Option, + source: BindingSource, + resolution: ActiveBindingResolution, + assertion: Option<&VerifiedFederatedAssertion>, + ) -> Result { + if binding_domain != authorization_domain { + return Err(EvidenceAdapterError::BindingDomainMismatch); + } + let binding_version = BindingVersion::new(binding_version)?; + let expires_at = expires_at.map(BindingExpiry::new).transpose()?; + let principal = FederatedPrincipal::new(issuer, subject)?; + if let Some(assertion) = assertion { + if assertion.authorization_domain() != authorization_domain + || assertion.principal() != &principal + { + return Err(EvidenceAdapterError::AssertionBindingMismatch); + } + if assertion + .key_attestation() + .is_some_and(|key| key.pubkey() != bound_pubkey) + { + return Err(EvidenceAdapterError::AssertionBindingMismatch); + } + } + let reason = match (resolution, source) { + (ActiveBindingResolution::Existing, _) => AuthorizationReason::ExistingBinding, + (ActiveBindingResolution::Enrolled, BindingSource::AttestedKey) => { + if assertion + .and_then(VerifiedFederatedAssertion::key_attestation) + .is_none_or(|key| key.pubkey() != bound_pubkey) + { + return Err(EvidenceAdapterError::AttestationRequired); + } + AuthorizationReason::EnrolledAttestedKey + } + (ActiveBindingResolution::Enrolled, BindingSource::Tofu) => { + AuthorizationReason::EnrolledTofu + } + (ActiveBindingResolution::Enrolled, BindingSource::Provisioned) => { + return Err(EvidenceAdapterError::InvalidBindingResolution) + } + }; + VersionedBindingRef::from_evidence_adapter( + authorization_domain, + binding_id, + principal, + bound_pubkey, + binding_version, + expires_at, + source, + reason, + ) + .map_err(Into::into) + } + + /// Translate successful local community policy into sealed admission. + pub fn community_access_from_policy( + &self, + tenant: &TenantContext, + resolved_domain: CommunityId, + scopes: Vec, + channel_ids: Option>, + ) -> Result { + if tenant.community() != resolved_domain { + return Err(EvidenceAdapterError::AdmissionDomainMismatch); + } + Ok(AuthorizedCommunityAccess::from_evidence_adapter( + resolved_domain, + scopes, + channel_ids, + )) + } +} + +/// Rejection while translating trusted-workspace verifier/store output. +#[derive(Debug, Error)] +pub enum EvidenceAdapterError { + /// Existing cryptographic proof failed. + #[error(transparent)] + Authentication(#[from] crate::AuthError), + /// Existing full Blossom operation verification failed. + #[error(transparent)] + Blossom(#[from] crate::blossom::BlossomAuthError), + /// Sealed context evidence was inconsistent. + #[error(transparent)] + Context(#[from] AuthContextError), + /// Proof method cannot establish the requested transport. + #[error("verified proof method does not match requested transport")] + TransportMethodMismatch, + /// Delegation output named a different actor. + #[error("verified delegation output does not match authenticated actor")] + DelegatedActorMismatch, + /// Operation-scoped delegation cannot be promoted to transport-wide. + #[error("verified delegation output is narrower than the transport")] + NarrowDelegation, + /// A sealed proof cannot have its verified delegation replaced. + #[error("verified transport delegation is already present")] + DelegationAlreadyPresent, + /// Validated assertion was already expired. + #[error("validated assertion is expired")] + AssertionExpired, + /// Validated assertion is not yet current. + #[error("validated assertion is not yet valid")] + AssertionNotYetValid, + /// Binding store result came from another exact domain. + #[error("active binding output belongs to another authorization domain")] + BindingDomainMismatch, + /// Binding does not match the validated assertion. + #[error("active binding output does not match validated assertion")] + AssertionBindingMismatch, + /// Attested-key enrollment did not carry the exact attested key. + #[error("attested-key enrollment lacks matching attestation")] + AttestationRequired, + /// Store provenance cannot produce the supplied lifecycle result. + #[error("active binding resolution and provenance are inconsistent")] + InvalidBindingResolution, + /// Local policy result belongs to another exact domain. + #[error("community admission belongs to another authorization domain")] + AdmissionDomainMismatch, +} + +#[cfg(test)] +mod tests { + use nostr::{EventBuilder, Keys, Kind, RelayUrl, Tag, Timestamp}; + + use super::*; + + fn domain(value: u128) -> CommunityId { + CommunityId::from_uuid(Uuid::from_u128(value)) + } + + #[test] + fn nip42_factory_reverifies_signature_challenge_transport_and_actor() { + let adapter = VerifiedEvidenceAdapter::new(); + let actor = Keys::generate(); + let owner = Keys::generate(); + let challenge = "challenge"; + let relay = "wss://relay.example"; + let event = EventBuilder::auth(challenge, RelayUrl::parse(relay).expect("relay url")) + .sign_with_keys(&actor) + .expect("auth event"); + + let proof = adapter + .verify_nip42( + domain(1), + AuthTransport::RelayWebSocket, + &event, + challenge, + relay, + Some(VerifiedDelegationOutput::from_workspace_verifier( + owner.public_key(), + actor.public_key(), + None, + true, + )), + ) + .expect("verified NIP-42 proof"); + assert_eq!( + proof.operation_binding().kind(), + VerifiedOperationBindingKind::NostrSession + ); + let substituted = adapter + .verify_nip42( + domain(1), + AuthTransport::RelayWebSocket, + &event, + challenge, + "wss://other.example", + None, + ) + .expect_err("relay URL substitution must fail full verification"); + assert!(matches!( + substituted, + EvidenceAdapterError::Authentication(_) + )); + assert!(matches!( + adapter.verify_nip42( + domain(1), + AuthTransport::Git, + &event, + challenge, + relay, + None, + ), + Err(EvidenceAdapterError::TransportMethodMismatch) + )); + assert!(matches!( + adapter.verify_nip42( + domain(1), + AuthTransport::RelayWebSocket, + &event, + challenge, + relay, + Some(VerifiedDelegationOutput::from_workspace_verifier( + owner.public_key(), + Keys::generate().public_key(), + None, + true, + )), + ), + Err(EvidenceAdapterError::DelegatedActorMismatch) + )); + } + + #[test] + fn assertion_and_binding_factories_reject_cross_domain_and_forged_provenance() { + let adapter = VerifiedEvidenceAdapter::new(); + let actor = Keys::generate().public_key(); + let assertion = adapter + .federated_assertion_from_validated_claims( + domain(1), + AuthTransport::HttpBridge, + "https://issuer.example", + "subject", + Some(actor), + AssertionTransport::TrustedProxy, + None, + 200, + 100, + ) + .expect("assertion"); + assert!(matches!( + adapter.active_binding_from_store( + domain(1), + domain(2), + Uuid::new_v4(), + "https://issuer.example", + "subject", + actor, + 1, + None, + BindingSource::AttestedKey, + ActiveBindingResolution::Existing, + Some(&assertion), + ), + Err(EvidenceAdapterError::BindingDomainMismatch) + )); + assert!(matches!( + adapter.active_binding_from_store( + domain(1), + domain(1), + Uuid::new_v4(), + "https://issuer.example", + "subject", + actor, + 1, + None, + BindingSource::Provisioned, + ActiveBindingResolution::Enrolled, + Some(&assertion), + ), + Err(EvidenceAdapterError::InvalidBindingResolution) + )); + } + + #[test] + fn blossom_factories_reverify_exact_hash_verb_and_server() { + let adapter = VerifiedEvidenceAdapter::new(); + let expiration = (Timestamp::now().as_secs() + 300).to_string(); + let upload_hash = "a".repeat(64); + let substituted_hash = "b".repeat(64); + let upload = EventBuilder::new(Kind::from(24_242), "upload") + .tags([ + Tag::parse(["t", "upload"]).expect("verb"), + Tag::parse(["x", &upload_hash]).expect("hash"), + Tag::parse(["server", "relay.example"]).expect("server"), + Tag::parse(["expiration", &expiration]).expect("expiration"), + ]) + .sign_with_keys(&Keys::generate()) + .expect("event"); + + assert!(adapter + .verify_blossom_upload(domain(1), &upload, &upload_hash, Some("relay.example"), 600,) + .is_ok()); + assert!(adapter + .verify_blossom_upload( + domain(1), + &upload, + &substituted_hash, + Some("relay.example"), + 600, + ) + .is_err()); + assert!(adapter + .verify_blossom_download(domain(1), &upload, &upload_hash, Some("relay.example"), 600,) + .is_err()); + assert!(adapter + .verify_blossom_upload(domain(1), &upload, &upload_hash, Some("other.example"), 600,) + .is_err()); + } +} diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index 5a8f24df8..15ad91119 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -59,6 +59,7 @@ async fn enforce_http_admission( /// /// Returns the authenticated public key and an event ID for replay detection. /// For X-Pubkey dev mode, the event ID is a zero hash (no replay concern). +#[cfg(test)] pub(crate) fn verify_bridge_auth( headers: &HeaderMap, method: &str, @@ -127,6 +128,99 @@ pub(crate) fn verify_bridge_auth_with_options( Err(api_error(StatusCode::UNAUTHORIZED, "missing Nostr auth")) } +/// Verify tenant bridge authentication and retain sealed proof evidence. +/// +/// The development-only `X-Pubkey` path returns no proof and is therefore +/// rejected later if this exact domain is configured for enforcement. +type ProtectedBridgeAuth = ( + nostr::PublicKey, + [u8; 32], + Option, +); + +pub(crate) fn verify_protected_bridge_auth( + headers: &HeaderMap, + method: &str, + url: &str, + body: Option<&[u8]>, + require_auth_token: bool, + require_payload: bool, + authorization_domain: buzz_core::CommunityId, +) -> Result)> { + let (pubkey, event_id) = verify_bridge_auth_with_options( + headers, + method, + url, + body, + require_auth_token, + require_payload, + )?; + let Some(encoded) = headers + .get("authorization") + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.strip_prefix("Nostr ")) + else { + return Ok((pubkey, event_id, None)); + }; + use base64::Engine as _; + let event_json = base64::engine::general_purpose::STANDARD + .decode(encoded) + .ok() + .and_then(|bytes| String::from_utf8(bytes).ok()) + .ok_or_else(|| api_error(StatusCode::UNAUTHORIZED, "invalid Nostr auth"))?; + let proof = buzz_auth::VerifiedEvidenceAdapter::new() + .verify_nip98( + authorization_domain, + buzz_auth::AuthTransport::HttpBridge, + &event_json, + url, + method, + body, + None, + ) + .map_err(|error| { + api_error( + StatusCode::UNAUTHORIZED, + &format!("NIP-98 evidence: {error}"), + ) + })?; + if proof.actor_pubkey() != pubkey { + return Err(api_error( + StatusCode::UNAUTHORIZED, + "NIP-98 evidence actor mismatch", + )); + } + Ok((pubkey, event_id, Some(proof))) +} + +pub(crate) fn retain_bridge_proof( + proof: Option, + auth_tag: Option<&str>, +) -> Result>, (StatusCode, Json)> { + let Some(proof) = proof else { + return Ok(None); + }; + let actor = proof.actor_pubkey(); + let proof = match crate::corporate_identity::verify_unconditional_nip_oa_owner(actor, auth_tag) + { + Some(owner) => buzz_auth::VerifiedEvidenceAdapter::new() + .attach_transport_delegation( + proof, + buzz_auth::VerifiedDelegationOutput::from_workspace_verifier( + owner, actor, None, true, + ), + ) + .map_err(|_| { + api_error( + StatusCode::UNAUTHORIZED, + "NIP-98 delegation evidence mismatch", + ) + })?, + None => proof, + }; + Ok(Some(Arc::new(proof))) +} + /// Corporate identity enrollment must always start from cryptographic proof of /// the Nostr key. The development-only `X-Pubkey` fallback is caller-controlled /// and therefore cannot safely participate in a durable identity binding. @@ -188,28 +282,78 @@ async fn verify_bridge_corporate_identity( headers: &HeaderMap, pubkey: nostr::PublicKey, auth_tag: Option<&str>, -) -> Result)> { - let identity_jwt = crate::corporate_identity::identity_jwt_from_headers( +) -> Result, (StatusCode, Json)> { + let identity_assertion = crate::corporate_identity::identity_assertion_from_headers( + state, + tenant.community(), headers, - &state.config.corporate_identity, - ); - crate::corporate_identity::verify_corporate_identity( + ) + .map_err(crate::corporate_identity::CorporateIdentityError::into_api_error)?; + match crate::corporate_identity::verify_corporate_identity( state, tenant.community(), pubkey, - identity_jwt.as_deref(), + identity_assertion.as_ref(), auth_tag, ) .await - .map_err(|e| e.into_api_error()) + { + Ok(proof) => Ok(Some(proof)), + Err(error) + if crate::authorization_runtime::transport::legacy_identity_lane( + state, + tenant.community(), + ) == crate::authorization_runtime::transport::LegacyIdentityLane::ObserveOnly => + { + tracing::warn!(error = ?error, "observational bridge identity verification unavailable"); + Ok(None) + } + Err(error) => Err(error.into_api_error()), + } +} + +fn seal_bridge_assertion( + state: &AppState, + tenant: &TenantContext, + proof: Option<&crate::corporate_identity::CorporateIdentityProof>, +) -> Result>, (StatusCode, Json)> { + let Some(proof) = proof else { + return Ok(None); + }; + match crate::corporate_identity::current_verified_assertion_for_proof( + state, + proof, + tenant.community(), + buzz_auth::AuthTransport::HttpBridge, + ) { + Ok(assertion) => Ok(assertion.map(Arc::new)), + Err(error) + if crate::authorization_runtime::transport::legacy_identity_lane( + state, + tenant.community(), + ) == crate::authorization_runtime::transport::LegacyIdentityLane::ObserveOnly => + { + tracing::warn!(error = %error, "observational bridge assertion sealing unavailable"); + Ok(None) + } + Err(error) => Err(error.into_api_error()), + } } async fn finalize_bridge_corporate_identity( state: &AppState, tenant: &TenantContext, pubkey: nostr::PublicKey, - proof: crate::corporate_identity::CorporateIdentityProof, + proof: Option, ) -> Result<(), (StatusCode, Json)> { + if crate::authorization_runtime::transport::legacy_identity_lane(state, tenant.community()) + != crate::authorization_runtime::transport::LegacyIdentityLane::Legacy + { + return Ok(()); + } + let Some(proof) = proof else { + return Ok(()); + }; crate::corporate_identity::finalize_corporate_identity(state, tenant.community(), pubkey, proof) .await .map(|_| ()) @@ -449,6 +593,7 @@ async fn handle_channel_window_filter( filter: &nostr::Filter, accessible_channels: &[uuid::Uuid], events: &mut Vec, + release_channels: &mut std::collections::BTreeSet, ) -> Result<(), (StatusCode, Json)> { use buzz_core::kind::{KIND_THREAD_SUMMARY, KIND_WINDOW_BOUNDS}; @@ -461,6 +606,7 @@ async fn handle_channel_window_filter( if !accessible_channels.contains(&ch_id) { return Ok(()); } + release_channels.insert(ch_id); // Composite request cursor: `until` + `before_id`, both or neither. The // window path has no timestamp-only fallback — that ambiguity is the @@ -678,7 +824,7 @@ pub async fn submit_event( })?; let url = nip98_expected_url(&state.config.relay_url, &tenant, "/events"); - let (pubkey, event_id_bytes) = verify_bridge_auth( + let (pubkey, event_id_bytes, verified_proof) = verify_protected_bridge_auth( &headers, "POST", &url, @@ -687,20 +833,27 @@ pub async fn submit_event( state.config.require_auth_token, state.config.corporate_identity.require, ), + false, + tenant.community(), )?; - let pubkey_hex = pubkey.to_hex(); - // Everything after auth — admission, replay, membership, parse, ingest — // runs inside the helper. The thin wrapper here owns the single terminal // attribution line so it fires for every outcome, including admission/ // replay/membership failures that previously returned before any log fired. - let outcome = - submit_event_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + let outcome = submit_event_authed( + &state, + &tenant, + &headers, + &body, + pubkey, + event_id_bytes, + verified_proof, + ) + .await; match &outcome { SubmitOutcome::Ok { accepted, .. } => { tracing::info!( - pubkey = %pubkey_hex, route = "/events", status = 200u16, accepted, @@ -714,7 +867,6 @@ pub async fn submit_event( .. } => { tracing::warn!( - pubkey = %pubkey_hex, route = "/events", status = 400u16, accepted = false, @@ -726,7 +878,6 @@ pub async fn submit_event( } SubmitOutcome::Rejected { kind, reason, .. } => { tracing::warn!( - pubkey = %pubkey_hex, route = "/events", status = 400u16, accepted = false, @@ -737,7 +888,6 @@ pub async fn submit_event( } SubmitOutcome::Err { status, .. } => { tracing::warn!( - pubkey = %pubkey_hex, route = "/events", status = status.as_u16(), accepted = false, @@ -802,6 +952,7 @@ async fn submit_event_authed( body: &[u8], pubkey: nostr::PublicKey, event_id_bytes: [u8; 32], + verified_proof: Option, ) -> SubmitOutcome { // Admission and replay checks fire before body parse — a 429 or replay // reject on a malformed body must still be attributed. @@ -877,6 +1028,59 @@ async fn submit_event_authed( }; } }; + let verified_proof = match retain_bridge_proof(verified_proof, auth_tag) { + Ok(proof) => proof, + Err(response) => { + return SubmitOutcome::Err { + status: response.0, + response, + } + } + }; + let verified_assertion = match seal_bridge_assertion(state, tenant, identity_proof.as_ref()) { + Ok(assertion) => assertion, + Err(response) => { + return SubmitOutcome::Err { + status: response.0, + response, + } + } + }; + let protected_result = match verified_proof.as_ref() { + Some(proof) => { + crate::authorization_runtime::transport::authorize_if_configured( + state, + Arc::clone(proof), + verified_assertion.clone(), + crate::protected_surface::event_ingest_capability(buzz_core::kind::event_kind_u32( + &event, + )), + uuid::Uuid::new_v4(), + "http_events", + ) + .await + } + None => crate::authorization_runtime::transport::authorize_unwired_if_configured( + state, + tenant.community(), + ), + }; + let protected = match protected_result { + Ok(authority) => authority, + Err(error) => { + tracing::warn!(error = %error, "http event protected authorization denied"); + return SubmitOutcome::Err { + status: StatusCode::FORBIDDEN, + response: api_error(StatusCode::FORBIDDEN, "protected authorization denied"), + }; + } + }; + if protected.revalidate().is_err() { + return SubmitOutcome::Err { + status: StatusCode::FORBIDDEN, + response: api_error(StatusCode::FORBIDDEN, "protected authorization expired"), + }; + } if let Err(e) = finalize_bridge_corporate_identity(state, tenant, pubkey, identity_proof).await { return SubmitOutcome::Err { @@ -884,13 +1088,25 @@ async fn submit_event_authed( response: e, }; } - if let Some(owner) = nip_oa_owner { + if let Some(owner) = nip_oa_owner.filter(|_| { + crate::authorization_runtime::transport::legacy_identity_lane(state, tenant.community()) + == crate::authorization_runtime::transport::LegacyIdentityLane::Legacy + }) { + if protected.revalidate().is_err() { + return SubmitOutcome::Err { + status: StatusCode::FORBIDDEN, + response: api_error(StatusCode::FORBIDDEN, "protected authorization expired"), + }; + } super::relay_members::materialize_nip_oa_owner(state, tenant, &pubkey, &owner).await; } let kind_u32 = buzz_core::kind::event_kind_u32(&event); let auth = IngestAuth::Http { pubkey, + owner_pubkey: nip_oa_owner, + verified_proof, + verified_assertion, scopes: buzz_auth::Scope::all_known(), // Pure Nostr: full scopes, channel access via membership auth_method: crate::handlers::ingest::HttpAuthMethod::Nip98, }; @@ -966,7 +1182,7 @@ pub async fn query_events( })?; let url = nip98_expected_url(&state.config.relay_url, &tenant, "/query"); - let (pubkey, event_id_bytes) = verify_bridge_auth( + let (pubkey, event_id_bytes, verified_proof) = verify_protected_bridge_auth( &headers, "POST", &url, @@ -975,31 +1191,32 @@ pub async fn query_events( state.config.require_auth_token, state.config.corporate_identity.require, ), + false, + tenant.community(), )?; - let pubkey_hex = pubkey.to_hex(); - // Admission, replay, membership, and filter execution all run inside the // helper. The single terminal attribution line fires here from the Result // so every outcome — including admission/replay/membership failures that // previously returned before any log — is attributed. - let result = - query_events_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + let result = query_events_authed( + &state, + &tenant, + &headers, + &body, + pubkey, + event_id_bytes, + verified_proof, + ) + .await; match &result { - Ok(Json(Value::Array(events))) => { - tracing::info!( - pubkey = %pubkey_hex, - route = "/query", - status = 200u16, - result_count = events.len(), - "HTTP bridge request" - ); + Ok(Json(Value::Array(_))) => { + tracing::info!(route = "/query", status = 200u16, "HTTP bridge request"); } Ok(_) => { - tracing::info!(pubkey = %pubkey_hex, route = "/query", status = 200u16, "HTTP bridge request"); + tracing::info!(route = "/query", status = 200u16, "HTTP bridge request"); } Err((status, _)) => { tracing::warn!( - pubkey = %pubkey_hex, route = "/query", status = status.as_u16(), "HTTP bridge request" @@ -1019,6 +1236,7 @@ async fn query_events_authed( body: &[u8], pubkey: nostr::PublicKey, event_id_bytes: [u8; 32], + verified_proof: Option, ) -> Result, (StatusCode, Json)> { enforce_http_admission(state, tenant, &pubkey).await?; check_nip98_replay(state, tenant, event_id_bytes).await?; @@ -1071,320 +1289,424 @@ async fn query_events_authed( .get_accessible_channel_ids_cached(tenant.community(), &pubkey_bytes) .await .map_err(|e| internal_error(&format!("channel access lookup: {e}")))?; + let verified_proof = retain_bridge_proof(verified_proof, auth_tag)?; + let verified_assertion = seal_bridge_assertion(state, tenant, identity_proof.as_ref())?; + let protected_result = match verified_proof { + Some(proof) => { + crate::authorization_runtime::transport::authorize_if_configured( + state, + proof, + verified_assertion, + buzz_auth::AuthorizationCapability::CommunityRead, + uuid::Uuid::new_v4(), + "http_query", + ) + .await + } + None => crate::authorization_runtime::transport::authorize_unwired_if_configured( + state, + tenant.community(), + ), + }; + let protected = protected_result + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization denied"))?; + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; finalize_bridge_corporate_identity(state, tenant, pubkey, identity_proof).await?; - if filters.iter().any(|f| f.search.is_some()) { - if has_mixed_search_filters(&filters) { - return Err(api_error( - StatusCode::BAD_REQUEST, - "mixed search and non-search filters not supported", + // Keep the complete post-authorization computation inside one release + // boundary. Every success and every backend-derived failure must pass the + // same final authority check before its response shape becomes observable. + let fetched = async { + if filters.iter().any(|f| f.search.is_some()) { + if has_mixed_search_filters(&filters) { + return Err(api_error( + StatusCode::BAD_REQUEST, + "mixed search and non-search filters not supported", + )); + } + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; + let result = handle_bridge_search( + state, + &raw_filters, + &filters, + &accessible_channels, + tenant, + &authed_pubkey_hex, + &pubkey_bytes, + ) + .await?; + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; + return Ok(result); + } + + if let Some(presence_events) = synthesize_presence(state, tenant, &filters).await { + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; + return Ok(( + Json(Value::Array(presence_events)), + std::collections::BTreeSet::new(), )); } - return handle_bridge_search( - state, - &raw_filters, - &filters, - &accessible_channels, - tenant, - &authed_pubkey_hex, - &pubkey_bytes, - ) - .await; - } - if let Some(presence_events) = synthesize_presence(state, tenant, &filters).await { - return Ok(Json(Value::Array(presence_events))); - } + let mut events: Vec = Vec::new(); + let mut release_channels = std::collections::BTreeSet::new(); + let mut handled: std::collections::HashSet = std::collections::HashSet::new(); - let mut events: Vec = Vec::new(); - let mut handled: std::collections::HashSet = std::collections::HashSet::new(); - - // Channel-window filters (`top_level: true`) — the GUI read-model surface. - // Dispatched first: a window filter is never a feed/thread/catchall query. - for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { - if !extension_flag(raw, "top_level") { - continue; - } - handle_channel_window_filter( - state, - tenant, - raw, - filter, - &accessible_channels, - &mut events, - ) - .await?; - handled.insert(idx); - } - - for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { - if handled.contains(&idx) { - continue; - } - let feed_types = match extract_feed_types(raw) { - Some(t) => t, - None => continue, - }; - - let limit = filter - .limit - .map(|l| (l as i64).min(BRIDGE_FEED_MAX_LIMIT)) - .unwrap_or(20); - let since = filter - .since - .and_then(|s| chrono::DateTime::from_timestamp(s.as_secs() as i64, 0)); - - let mut seen_types = std::collections::HashSet::new(); - let mut seen = std::collections::HashSet::new(); - let mut feed_count = 0i64; - for feed_type in &feed_types { - let canonical = if feed_type == "agent_activity" { - "activity" - } else { - feed_type.as_str() - }; - if !seen_types.insert(canonical) { + // Channel-window filters (`top_level: true`) — the GUI read-model surface. + // Dispatched first: a window filter is never a feed/thread/catchall query. + for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { + if !extension_flag(raw, "top_level") { continue; } - if feed_count >= limit { - break; + handle_channel_window_filter( + state, + tenant, + raw, + filter, + &accessible_channels, + &mut events, + &mut release_channels, + ) + .await?; + handled.insert(idx); + } + + for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { + if handled.contains(&idx) { + continue; } - let remaining = limit - feed_count; - let type_events = match canonical { - "mentions" => state - .db - .query_feed_mentions_routed( - "bridge_feed", - tenant.community(), - &pubkey_bytes, - &accessible_channels, - since, - remaining, - ) - .await - .map_err(|e| internal_error(&format!("feed mentions error: {e}")))?, - "needs_action" => state - .db - .query_feed_needs_action_routed( - "bridge_feed", - tenant.community(), - &pubkey_bytes, - &accessible_channels, - since, - remaining, - ) - .await - .map_err(|e| internal_error(&format!("feed needs_action error: {e}")))?, - "activity" => state - .db - .query_feed_activity_routed( - "bridge_feed", - tenant.community(), - &accessible_channels, - since, - remaining, - ) - .await - .map_err(|e| internal_error(&format!("feed activity error: {e}")))?, - _ => continue, + let feed_types = match extract_feed_types(raw) { + Some(t) => t, + None => continue, }; - for se in type_events { - if !seen.insert(se.event.id) { + + let limit = filter + .limit + .map(|l| (l as i64).min(BRIDGE_FEED_MAX_LIMIT)) + .unwrap_or(20); + let since = filter + .since + .and_then(|s| chrono::DateTime::from_timestamp(s.as_secs() as i64, 0)); + + let mut seen_types = std::collections::HashSet::new(); + let mut seen = std::collections::HashSet::new(); + let mut feed_count = 0i64; + for feed_type in &feed_types { + let canonical = if feed_type == "agent_activity" { + "activity" + } else { + feed_type.as_str() + }; + if !seen_types.insert(canonical) { continue; } + if feed_count >= limit { + break; + } + let remaining = limit - feed_count; + let type_events = match canonical { + "mentions" => state + .db + .query_feed_mentions_routed( + "bridge_feed", + tenant.community(), + &pubkey_bytes, + &accessible_channels, + since, + remaining, + ) + .await + .map_err(|e| internal_error(&format!("feed mentions error: {e}")))?, + "needs_action" => state + .db + .query_feed_needs_action_routed( + "bridge_feed", + tenant.community(), + &pubkey_bytes, + &accessible_channels, + since, + remaining, + ) + .await + .map_err(|e| internal_error(&format!("feed needs_action error: {e}")))?, + "activity" => state + .db + .query_feed_activity_routed( + "bridge_feed", + tenant.community(), + &accessible_channels, + since, + remaining, + ) + .await + .map_err(|e| internal_error(&format!("feed activity error: {e}")))?, + _ => continue, + }; + for se in type_events { + if !seen.insert(se.event.id) { + continue; + } + if !event_in_accessible_channel(&se, &accessible_channels) { + continue; + } + // Defense-in-depth: never deliver a result-gated event (e.g. kind:44200 + // or kind:30622) to a non-owner via the feed path, even though feed SQL + // kind allowlists already exclude these kinds. + if !buzz_core::filter::reader_authorized_for_event( + &se.event, + &authed_pubkey_hex, + ) { + continue; + } + if append_bridge_stored_event(&mut events, &mut release_channels, &se) { + feed_count += 1; + } + } + } + handled.insert(idx); + } + + let e_tag_key = nostr::SingleLetterTag::lowercase(nostr::Alphabet::E); + for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { + if handled.contains(&idx) { + continue; + } + let depth = match extract_depth_limit(raw) { + Some(d) => d, + None => continue, + }; + let e_values = match filter.generic_tags.get(&e_tag_key) { + Some(vs) if vs.len() == 1 => vs, + _ => continue, + }; + let root_hex = match e_values.iter().next() { + Some(h) => h, + None => continue, + }; + let root_bytes = match hex::decode(root_hex) { + Ok(b) if b.len() == 32 => b, + _ => continue, + }; + + if let Some(ch_id) = extract_channel_from_filter(filter) { + if !accessible_channels.contains(&ch_id) { + handled.insert(idx); + continue; + } + } + + let limit = filter + .limit + .unwrap_or(100) + .min(BRIDGE_THREAD_MAX_LIMIT as usize) as u32; + let thread_cursor = extract_thread_cursor(raw); + let thread_replies = state + .db + .get_thread_replies( + tenant.community(), + &root_bytes, + Some(depth), + limit, + thread_cursor.as_deref(), + ) + .await + .map_err(|e| internal_error(&format!("thread query error: {e}")))?; + + for reply in thread_replies { + let se = reply.stored_event; if !event_in_accessible_channel(&se, &accessible_channels) { continue; } // Defense-in-depth: never deliver a result-gated event (e.g. kind:44200 - // or kind:30622) to a non-owner via the feed path, even though feed SQL - // kind allowlists already exclude these kinds. + // or kind:30622) to a non-owner via the thread path, even though + // requires_h_channel_scope already excludes these kinds from thread metadata. if !buzz_core::filter::reader_authorized_for_event(&se.event, &authed_pubkey_hex) { continue; } - if let Ok(v) = serde_json::to_value(&se.event) { - events.push(v); - feed_count += 1; + append_bridge_stored_event(&mut events, &mut release_channels, &se); + } + handled.insert(idx); + } + + // Phase 1 — pure construction + validation, in filter order. Access-scope + // skips and the `before_id` BAD_REQUEST are decided here, before any DB + // work is issued (validation errors are deterministic client mistakes, so + // surfacing them ahead of transient DB errors is strictly more predictable). + let mut catchall_queries: Vec<(usize, buzz_db::EventQuery)> = Vec::new(); + for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { + if handled.contains(&idx) { + continue; + } + + if let Some(ch_id) = extract_channel_from_filter(filter) { + if !accessible_channels.contains(&ch_id) { + continue; } } - } - handled.insert(idx); - } - let e_tag_key = nostr::SingleLetterTag::lowercase(nostr::Alphabet::E); - for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { - if handled.contains(&idx) { - continue; - } - let depth = match extract_depth_limit(raw) { - Some(d) => d, - None => continue, - }; - let e_values = match filter.generic_tags.get(&e_tag_key) { - Some(vs) if vs.len() == 1 => vs, - _ => continue, - }; - let root_hex = match e_values.iter().next() { - Some(h) => h, - None => continue, - }; - let root_bytes = match hex::decode(root_hex) { - Ok(b) if b.len() == 32 => b, - _ => continue, - }; - - if let Some(ch_id) = extract_channel_from_filter(filter) { - if !accessible_channels.contains(&ch_id) { - handled.insert(idx); - continue; - } - } - - let limit = filter - .limit - .unwrap_or(100) - .min(BRIDGE_THREAD_MAX_LIMIT as usize) as u32; - let thread_cursor = extract_thread_cursor(raw); - let thread_replies = state - .db - .get_thread_replies( + let mut query = crate::handlers::req::build_event_query_from_filter( + filter, + &pubkey_bytes, + state, tenant.community(), - &root_bytes, - Some(depth), - limit, - thread_cursor.as_deref(), ) - .await - .map_err(|e| internal_error(&format!("thread query error: {e}")))?; - - for reply in thread_replies { - let se = reply.stored_event; - if !event_in_accessible_channel(&se, &accessible_channels) { - continue; + .await; + crate::handlers::req::apply_access_scope_to_query( + &mut query, + extract_channel_from_filter(filter), + &accessible_channels, + ); + // Shared-gated visibility pushdown: must mirror WS REQ so that a page of + // newer private events does not starve older shared ones off the page. + if crate::handlers::req::filter_can_match_shared_gated_kinds(filter) { + query.shared_gated_reader = Some(pubkey_bytes.clone()); } - // Defense-in-depth: never deliver a result-gated event (e.g. kind:44200 - // or kind:30622) to a non-owner via the thread path, even though - // requires_h_channel_scope already excludes these kinds from thread metadata. - if !buzz_core::filter::reader_authorized_for_event(&se.event, &authed_pubkey_hex) { - continue; - } - if let Ok(v) = serde_json::to_value(&se.event) { - events.push(v); - } - } - handled.insert(idx); - } - // Phase 1 — pure construction + validation, in filter order. Access-scope - // skips and the `before_id` BAD_REQUEST are decided here, before any DB - // work is issued (validation errors are deterministic client mistakes, so - // surfacing them ahead of transient DB errors is strictly more predictable). - let mut catchall_queries: Vec<(usize, buzz_db::EventQuery)> = Vec::new(); - for (idx, (raw, filter)) in raw_filters.iter().zip(filters.iter()).enumerate() { - if handled.contains(&idx) { - continue; - } - - if let Some(ch_id) = extract_channel_from_filter(filter) { - if !accessible_channels.contains(&ch_id) { - continue; - } - } - - let mut query = crate::handlers::req::build_event_query_from_filter( - filter, - &pubkey_bytes, - state, - tenant.community(), - ) - .await; - crate::handlers::req::apply_access_scope_to_query( - &mut query, - extract_channel_from_filter(filter), - &accessible_channels, - ); - // Shared-gated visibility pushdown: must mirror WS REQ so that a page of - // newer private events does not starve older shared ones off the page. - if crate::handlers::req::filter_can_match_shared_gated_kinds(filter) { - query.shared_gated_reader = Some(pubkey_bytes.clone()); - } - - match extract_before_id(raw) { - BeforeId::Malformed => { - return Err(api_error( - StatusCode::BAD_REQUEST, - "before_id must be a 64-char hex event id", - )); - } - BeforeId::Valid(bid) => { - if query.until.is_none() { + match extract_before_id(raw) { + BeforeId::Malformed => { return Err(api_error( StatusCode::BAD_REQUEST, - "before_id requires until to be set", + "before_id must be a 64-char hex event id", )); } - query.before_id = Some(bid); + BeforeId::Valid(bid) => { + if query.until.is_none() { + return Err(api_error( + StatusCode::BAD_REQUEST, + "before_id requires until to be set", + )); + } + query.before_id = Some(bid); + } + BeforeId::Absent => {} } - BeforeId::Absent => {} + + // Honor `page` on non-search general queries so offset paging works for + // the empty-query people directory (kind:0 listing). The FTS path + // (`handle_bridge_search`) has its own `page`/`per_page`; a filter with + // no `search` field lands here instead, where paging would otherwise be + // dropped and the directory would terminate at its first page. Deterministic + // ordering in `query_events` (`created_at DESC, id ASC`) makes offset paging + // stable. `page` defaults to 1 → offset 0, so unrelated general queries are + // unaffected. + if let Some(offset) = extract_page_offset(raw, query.limit) { + query.offset = Some(offset); + } + + catchall_queries.push((idx, query)); } - // Honor `page` on non-search general queries so offset paging works for - // the empty-query people directory (kind:0 listing). The FTS path - // (`handle_bridge_search`) has its own `page`/`per_page`; a filter with - // no `search` field lands here instead, where paging would otherwise be - // dropped and the directory would terminate at its first page. Deterministic - // ordering in `query_events` (`created_at DESC, id ASC`) makes offset paging - // stable. `page` defaults to 1 → offset 0, so unrelated general queries are - // unaffected. - if let Some(offset) = extract_page_offset(raw, query.limit) { - query.offset = Some(offset); - } + // Phase 2 — DB reads, bounded-concurrent, order-preserving (`buffered`). + // Phase 3 consumes results in original filter order, so response ordering + // and error semantics match the previous serial loop. + use futures_util::stream::{self, StreamExt}; + let db = state.db.clone(); + let mut catchall_results = + stream::iter(catchall_queries.into_iter().map(|(idx, query)| { + let db = db.clone(); + async move { (idx, db.query_events_routed("bridge_query", &query).await) } + })) + .buffered(crate::handlers::req::FILTER_QUERY_CONCURRENCY); - catchall_queries.push((idx, query)); - } - - // Phase 2 — DB reads, bounded-concurrent, order-preserving (`buffered`). - // Phase 3 consumes results in original filter order, so response ordering - // and error semantics match the previous serial loop. - use futures_util::stream::{self, StreamExt}; - let db = state.db.clone(); - let mut catchall_results = stream::iter(catchall_queries.into_iter().map(|(idx, query)| { - let db = db.clone(); - async move { (idx, db.query_events_routed("bridge_query", &query).await) } - })) - .buffered(crate::handlers::req::FILTER_QUERY_CONCURRENCY); - - // Phase 3 — post-processing, strictly in filter order. - while let Some((idx, filter_events)) = catchall_results.next().await { - let filter = &filters[idx]; - match filter_events { - Ok(stored_events) => { - for se in stored_events { - if !event_in_accessible_channel(&se, &accessible_channels) { - continue; - } - if !buzz_core::filter::filters_match(std::slice::from_ref(filter), &se) { - continue; - } - // Result-level read auth: never hand a viewer-private snapshot - // (kind:30622) to anyone but its owner, even via kindless `ids`. - // Also enforces author-only kinds (30300/30350) and the persona - // shared-gate (kind:30175 without ["shared","true"]). Single call - // covers all three gated event classes. - if !crate::handlers::req::event_visible_to_reader(&se.event, &pubkey_bytes) { - continue; - } - if let Ok(v) = serde_json::to_value(&se.event) { - events.push(v); + // Phase 3 — post-processing, strictly in filter order. + while let Some((idx, filter_events)) = catchall_results.next().await { + let filter = &filters[idx]; + match filter_events { + Ok(stored_events) => { + for se in stored_events { + if !event_in_accessible_channel(&se, &accessible_channels) { + continue; + } + if !buzz_core::filter::filters_match(std::slice::from_ref(filter), &se) { + continue; + } + // Result-level read auth: never hand a viewer-private snapshot + // (kind:30622) to anyone but its owner, even via kindless `ids`. + // Also enforces author-only kinds (30300/30350) and the persona + // shared-gate (kind:30175 without ["shared","true"]). Single call + // covers all three gated event classes. + if !crate::handlers::req::event_visible_to_reader(&se.event, &pubkey_bytes) + { + continue; + } + append_bridge_stored_event(&mut events, &mut release_channels, &se); } } - } - Err(e) => { - return Err(internal_error(&format!("query error: {e}"))); + Err(e) => { + return Err(internal_error(&format!("query error: {e}"))); + } } } - } - Ok(Json(Value::Array(events))) + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; + Ok((Json(Value::Array(events)), release_channels)) + } + .await; + let fetched = match fetched { + Ok((response, release_channels)) => { + revalidate_bridge_response_channels( + state, + tenant.community(), + &pubkey_bytes, + &release_channels, + ) + .await?; + Ok(response) + } + Err(error) => Err(error), + }; + release_protected_bridge_fetch(fetched, |fetched| protected.release_fetched(fetched))? +} + +/// Revalidate the authoritative stored channel ids retained alongside the +/// response without caches after all asynchronous fetch work. This is the HTTP +/// bridge's post-fetch/pre-emission channel fence. Event tags are intentionally +/// not consulted: a stored channel-scoped row remains fenced even if its signed +/// event omits or malforms an `h` tag. +async fn revalidate_bridge_response_channels( + state: &AppState, + community_id: buzz_core::tenant::CommunityId, + actor: &[u8], + channels: &std::collections::BTreeSet, +) -> Result<(), (StatusCode, Json)> { + for &channel_id in channels { + let allowed = state + .db + .channel_read_authorized(community_id, channel_id, actor) + .await + .map_err(|error| internal_error(&format!("channel release fence: {error}")))?; + if !allowed { + return Err(api_error( + StatusCode::FORBIDDEN, + "restricted: channel access changed before response release", + )); + } + } + Ok(()) +} + +fn append_bridge_stored_event( + events: &mut Vec, + release_channels: &mut std::collections::BTreeSet, + stored: &buzz_core::StoredEvent, +) -> bool { + let Ok(value) = serde_json::to_value(&stored.event) else { + return false; + }; + if let Some(channel_id) = stored.channel_id { + release_channels.insert(channel_id); + } + events.push(value); + true } /// Count events via HTTP bridge (NIP-98 auth). Returns `{"count": N}`. @@ -1414,7 +1736,7 @@ pub async fn count_events( })?; let url = nip98_expected_url(&state.config.relay_url, &tenant, "/count"); - let (pubkey, event_id_bytes) = verify_bridge_auth( + let (pubkey, event_id_bytes, verified_proof) = verify_protected_bridge_auth( &headers, "POST", &url, @@ -1423,29 +1745,29 @@ pub async fn count_events( state.config.require_auth_token, state.config.corporate_identity.require, ), + false, + tenant.community(), )?; - let pubkey_hex = pubkey.to_hex(); - // Admission, replay, membership, and count execution all run inside the // helper. The single terminal attribution line fires here from the Result // so every outcome — including admission/replay/membership failures that // previously returned before any log — is attributed. - let result = - count_events_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + let result = count_events_authed( + &state, + &tenant, + &headers, + &body, + pubkey, + event_id_bytes, + verified_proof, + ) + .await; match &result { - Ok(Json(value)) => { - let count = value.get("count").and_then(Value::as_u64); - tracing::info!( - pubkey = %pubkey_hex, - route = "/count", - status = 200u16, - result_count = count, - "HTTP bridge request" - ); + Ok(Json(_)) => { + tracing::info!(route = "/count", status = 200u16, "HTTP bridge request"); } Err((status, _)) => { tracing::warn!( - pubkey = %pubkey_hex, route = "/count", status = status.as_u16(), "HTTP bridge request" @@ -1465,6 +1787,7 @@ async fn count_events_authed( body: &[u8], pubkey: nostr::PublicKey, event_id_bytes: [u8; 32], + verified_proof: Option, ) -> Result, (StatusCode, Json)> { enforce_http_admission(state, tenant, &pubkey).await?; check_nip98_replay(state, tenant, event_id_bytes).await?; @@ -1509,9 +1832,36 @@ async fn count_events_authed( .get_accessible_channel_ids_cached(tenant.community(), &pubkey_bytes) .await .map_err(|e| internal_error(&format!("channel access lookup: {e}")))?; + let verified_proof = retain_bridge_proof(verified_proof, auth_tag)?; + let verified_assertion = seal_bridge_assertion(state, tenant, identity_proof.as_ref())?; + let protected_result = match verified_proof { + Some(proof) => { + crate::authorization_runtime::transport::authorize_if_configured( + state, + proof, + verified_assertion, + buzz_auth::AuthorizationCapability::CommunityRead, + uuid::Uuid::new_v4(), + "http_count", + ) + .await + } + None => crate::authorization_runtime::transport::authorize_unwired_if_configured( + state, + tenant.community(), + ), + }; + let protected = Arc::new( + protected_result + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization denied"))?, + ); + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; finalize_bridge_corporate_identity(state, tenant, pubkey, identity_proof).await?; let mut total: u64 = 0; + let mut release_channels = std::collections::BTreeSet::new(); for filter in &filters { let needs_author_only_filtering = crate::handlers::req::filter_can_match_author_only_kinds(filter); @@ -1535,6 +1885,7 @@ async fn count_events_authed( if !accessible_channels.contains(&ch_id) { continue; // Skip filters targeting inaccessible channels. } + release_channels.insert(ch_id); // Channel is accessible — count with pushability check. let mut query = crate::handlers::req::build_event_query_from_filter( filter, @@ -1559,7 +1910,10 @@ async fn count_events_authed( && !needs_result_gated_filtering && !needs_shared_gate_filtering { - match state.db.count_events_routed("bridge_count", &query).await { + let fetched = state.db.count_events_routed("bridge_count", &query).await; + match release_protected_bridge_fetch(fetched, |fetched| { + protected.release_fetched(fetched) + })? { Ok(n) => total += n as u64, Err(e) => { return Err(internal_error(&format!("count error: {e}"))); @@ -1569,11 +1923,13 @@ async fn count_events_authed( // Fallback: query + post-filter for non-pushable constraints. let mut q = query; crate::handlers::req::apply_count_fallback_limit(&mut q); - match state + let fetched = state .db .query_events_routed_bounded("bridge_count_fallback", &q) - .await - { + .await; + match release_protected_bridge_fetch(fetched, |fetched| { + protected.release_fetched(fetched) + })? { Ok(stored_events) => { if crate::handlers::req::count_fallback_exceeded(stored_events.len()) { metrics::counter!("buzz_count_fallback_rejections_total").increment(1); @@ -1604,6 +1960,7 @@ async fn count_events_authed( } else { // No channel filter — use SQL-level channel_ids pushdown to count // only events in accessible channels (+ global events). + release_channels.extend(accessible_channels.iter().copied()); let mut query = crate::handlers::req::build_event_query_from_filter( filter, &pubkey_bytes, @@ -1630,7 +1987,10 @@ async fn count_events_authed( && !needs_shared_gate_filtering { query.limit = None; - match state.db.count_events_routed("bridge_count", &query).await { + let fetched = state.db.count_events_routed("bridge_count", &query).await; + match release_protected_bridge_fetch(fetched, |fetched| { + protected.release_fetched(fetched) + })? { Ok(n) => total += n as u64, Err(e) => { return Err(internal_error(&format!("count error: {e}"))); @@ -1639,11 +1999,13 @@ async fn count_events_authed( } else { // Fallback: query a bounded candidate set + post-filter. crate::handlers::req::apply_count_fallback_limit(&mut query); - match state + let fetched = state .db .query_events_routed_bounded("bridge_count_fallback", &query) - .await - { + .await; + match release_protected_bridge_fetch(fetched, |fetched| { + protected.release_fetched(fetched) + })? { Ok(stored_events) => { if crate::handlers::req::count_fallback_exceeded(stored_events.len()) { metrics::counter!("buzz_count_fallback_rejections_total").increment(1); @@ -1674,6 +2036,20 @@ async fn count_events_authed( } } + if !crate::connection::release_channel_set_read_authority( + state.db.clone(), + tenant.community(), + release_channels.into_iter().collect(), + pubkey_bytes, + Some(protected), + ) + .await + { + return Err(api_error( + StatusCode::FORBIDDEN, + "restricted: channel access changed before response release", + )); + } Ok(Json(serde_json::json!({ "count": total }))) } @@ -1723,7 +2099,7 @@ async fn handle_bridge_search( tenant: &buzz_core::tenant::TenantContext, reader_pubkey_hex: &str, pubkey_bytes: &[u8], -) -> Result, (StatusCode, Json)> { +) -> Result<(Json, std::collections::BTreeSet), (StatusCode, Json)> { // Bridge always includes global (channel-less) events — same as WS with // full scopes. `None` means no accessible channels and no global access → // empty result set (the caller short-circuits exactly as the WS door EOSEs). @@ -1732,10 +2108,16 @@ async fn handle_bridge_search( true, // include_global ) { Some(scope) => scope, - None => return Ok(Json(Value::Array(Vec::new()))), + None => { + return Ok(( + Json(Value::Array(Vec::new())), + std::collections::BTreeSet::new(), + )); + } }; let mut events: Vec = Vec::new(); + let mut release_channels = std::collections::BTreeSet::new(); let mut seen_ids: std::collections::HashSet<[u8; 32]> = std::collections::HashSet::new(); for (raw, filter) in raw_filters.iter().zip(filters) { @@ -1848,13 +2230,11 @@ async fn handle_bridge_search( if !seen_ids.insert(*id_array) { continue; } - if let Ok(v) = serde_json::to_value(&stored.event) { - events.push(v); - } + append_bridge_stored_event(&mut events, &mut release_channels, stored); } } - Ok(Json(Value::Array(events))) + Ok((Json(Value::Array(events)), release_channels)) } /// Query parameters for the webhook trigger endpoint. @@ -1894,6 +2274,19 @@ pub async fn workflow_webhook( .map_err(|_| not_found("workflow not found"))?; let community_id = tenant.community(); + // A webhook secret authenticates the trigger but carries no protected + // authority into the workflow run or its delayed actions. Enforce must + // therefore stop before even the run row exists until a transaction-owning + // workflow executor can validate authority at every durable commit. + let mode = state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(community_id)); + crate::protected_surface::require_effect_permit( + mode, + crate::protected_surface::EffectSurfaceId::WorkflowBackgroundExecution, + ) + .map_err(|_| not_found("workflow not found"))?; + let workflow = state .db .get_workflow(community_id, id) @@ -2078,12 +2471,43 @@ async fn synthesize_presence( all_pubkeys.dedup(); // Look up Redis. - let presence_map = state + let stored_presence = state .pubsub .get_presence_bulk(tenant, &all_pubkeys) .await .unwrap_or_default(); + let verifier = crate::authorization_runtime::ephemeral::AuthorityTokenVerifier::new( + state.db.clone(), + state.relay_keypair.secret_key().as_secret_bytes(), + ); + let mut presence_map = std::collections::HashMap::new(); + for (pubkey_hex, stored_status) in stored_presence { + match crate::authorization_runtime::ephemeral::decode_presence(&stored_status) { + Ok(Some(protected)) => { + let Ok(pubkey) = nostr::PublicKey::from_hex(&pubkey_hex) else { + continue; + }; + if verifier + .verify_actor_context( + tenant.community(), + protected.context_id, + pubkey.to_bytes(), + &protected.authority, + ) + .await + .is_ok() + { + presence_map.insert(pubkey_hex, protected.status); + } + } + Ok(None) if !state.is_protected_enforcing(tenant.community()) => { + presence_map.insert(pubkey_hex, stored_status); + } + Ok(None) | Err(_) => {} + } + } + if presence_map.is_empty() { return Some(Vec::new()); } @@ -2139,7 +2563,7 @@ async fn authorize_moderation_read( headers: &HeaderMap, path: &str, raw_query: Option<&str>, -) -> Result)> { +) -> Result)> { let raw_host = headers .get(axum::http::header::HOST) .and_then(|v| v.to_str().ok()) @@ -2158,7 +2582,7 @@ async fn authorize_moderation_read( _ => path.to_string(), }; let url = nip98_expected_url(&state.config.relay_url, &tenant, &path_with_query); - let (pubkey, event_id_bytes) = verify_bridge_auth( + let (pubkey, event_id_bytes, verified_proof) = verify_protected_bridge_auth( headers, "GET", &url, @@ -2167,6 +2591,8 @@ async fn authorize_moderation_read( state.config.require_auth_token, state.config.corporate_identity.require, ), + false, + tenant.community(), )?; check_nip98_replay(state, &tenant, event_id_bytes).await?; let pubkey_bytes = pubkey.to_bytes().to_vec(); @@ -2190,14 +2616,91 @@ async fn authorize_moderation_read( "restricted: moderator access required", ) })?; + let verified_proof = retain_bridge_proof(verified_proof, auth_tag)?; + let verified_assertion = seal_bridge_assertion(state, &tenant, identity_proof.as_ref())?; + let protected_result = match verified_proof { + Some(proof) => { + crate::authorization_runtime::transport::authorize_if_configured( + state, + proof, + verified_assertion, + buzz_auth::AuthorizationCapability::Moderate, + uuid::Uuid::new_v4(), + "http_moderation_read", + ) + .await + } + None => crate::authorization_runtime::transport::authorize_unwired_if_configured( + state, + tenant.community(), + ), + }; + let protected = protected_result + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization denied"))?; + protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; finalize_bridge_corporate_identity(state, &tenant, pubkey, identity_proof).await?; - Ok(tenant) + Ok(ModerationReadAuthorization { + tenant, + protected, + actor_pubkey: pubkey_bytes, + }) +} + +/// Retains exact protected authority through the moderation fetch boundary. +struct ModerationReadAuthorization { + tenant: TenantContext, + protected: crate::authorization_runtime::transport::ProtectedAuthorization, + actor_pubkey: Vec, +} + +/// Revalidate both provider-neutral protected authority and the local +/// moderation role after the fetch transaction releases its locks and before +/// the response is returned to the transport. +async fn release_moderation_read( + state: &Arc, + authorization: &ModerationReadAuthorization, +) -> Result<(), (StatusCode, Json)> { + authorization + .protected + .revalidate() + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired"))?; + crate::handlers::moderation_authz::authorize_moderation_action( + &authorization.tenant, + state, + &authorization.actor_pubkey, + None, + crate::handlers::moderation_authz::ModerationTarget::None, + crate::handlers::moderation_authz::ModerationAction::ViewQueue, + ) + .await + .map(|_| ()) + .map_err(|_| { + api_error( + StatusCode::FORBIDDEN, + "restricted: moderator access changed before response release", + ) + }) } /// Cap on rows returned by a single moderation read. const MODERATION_READ_LIMIT: i64 = 500; +/// Releases a bridge fetch result only after the retained authority check. +/// +/// Keeping the complete `Result` inside the release boundary ensures that +/// authority loss takes precedence over both successful rows and backend +/// failures; neither response shape is observable after the lease is stale. +fn release_protected_bridge_fetch( + fetched: Result, + release: impl FnOnce(Result) -> Result, R>, +) -> Result, (StatusCode, Json)> { + release(fetched) + .map_err(|_| api_error(StatusCode::FORBIDDEN, "protected authorization expired")) +} + /// Optional `?status=` and `?limit=` query for moderation reads. #[derive(serde::Deserialize, Default)] pub struct ModerationReadQuery { @@ -2219,23 +2722,51 @@ pub async fn moderation_reports( RawQuery(raw_query): RawQuery, Query(q): Query, ) -> Result, (StatusCode, Json)> { - let tenant = authorize_moderation_read( + let authorization = authorize_moderation_read( &state, &headers, "/moderation/reports", raw_query.as_deref(), ) .await?; - let rows = state + let mut transaction = state .db - .list_moderation_reports( - tenant.community(), - q.status.as_deref(), - clamp_limit(q.limit), - ) + .begin_transaction() .await - .map_err(|e| internal_error(&format!("list reports: {e}")))?; - Ok(Json(Value::Array(rows.iter().map(report_json).collect()))) + .map_err(|error| internal_error(&format!("begin moderation read: {error}")))?; + crate::handlers::moderation_authz::authorize_moderation_action_tx( + &mut transaction, + &authorization.tenant, + &authorization.actor_pubkey, + None, + crate::handlers::moderation_authz::ModerationTarget::None, + crate::handlers::moderation_authz::ModerationAction::ViewQueue, + ) + .await + .map_err(|_| { + api_error( + StatusCode::FORBIDDEN, + "restricted: moderator access required", + ) + })?; + let fetched = buzz_db::moderation::list_reports_tx( + &mut transaction, + authorization.tenant.community(), + q.status.as_deref(), + clamp_limit(q.limit), + ) + .await; + let rows = release_protected_bridge_fetch(fetched, |fetched| { + authorization.protected.release_fetched(fetched) + })? + .map_err(|e| internal_error(&format!("list reports: {e}")))?; + let response = Json(Value::Array(rows.iter().map(report_json).collect())); + transaction + .commit() + .await + .map_err(|error| internal_error(&format!("finish moderation read: {error}")))?; + release_moderation_read(&state, &authorization).await?; + Ok(response) } /// `GET /moderation/audit` — the moderation audit log (NIP-98 + mod-authz). @@ -2245,15 +2776,46 @@ pub async fn moderation_audit( RawQuery(raw_query): RawQuery, Query(q): Query, ) -> Result, (StatusCode, Json)> { - let tenant = + let authorization = authorize_moderation_read(&state, &headers, "/moderation/audit", raw_query.as_deref()) .await?; - let rows = state + let mut transaction = state .db - .list_moderation_actions(tenant.community(), clamp_limit(q.limit)) + .begin_transaction() .await - .map_err(|e| internal_error(&format!("list actions: {e}")))?; - Ok(Json(Value::Array(rows.iter().map(action_json).collect()))) + .map_err(|error| internal_error(&format!("begin moderation read: {error}")))?; + crate::handlers::moderation_authz::authorize_moderation_action_tx( + &mut transaction, + &authorization.tenant, + &authorization.actor_pubkey, + None, + crate::handlers::moderation_authz::ModerationTarget::None, + crate::handlers::moderation_authz::ModerationAction::ViewQueue, + ) + .await + .map_err(|_| { + api_error( + StatusCode::FORBIDDEN, + "restricted: moderator access required", + ) + })?; + let fetched = buzz_db::moderation::list_actions_tx( + &mut transaction, + authorization.tenant.community(), + clamp_limit(q.limit), + ) + .await; + let rows = release_protected_bridge_fetch(fetched, |fetched| { + authorization.protected.release_fetched(fetched) + })? + .map_err(|e| internal_error(&format!("list actions: {e}")))?; + let response = Json(Value::Array(rows.iter().map(action_json).collect())); + transaction + .commit() + .await + .map_err(|error| internal_error(&format!("finish moderation read: {error}")))?; + release_moderation_read(&state, &authorization).await?; + Ok(response) } /// `GET /moderation/restricted` — currently banned/timed-out members. @@ -2261,14 +2823,42 @@ pub async fn moderation_restricted( State(state): State>, headers: HeaderMap, ) -> Result, (StatusCode, Json)> { - let tenant = + let authorization = authorize_moderation_read(&state, &headers, "/moderation/restricted", None).await?; - let rows = state + let mut transaction = state .db - .list_community_restrictions(tenant.community()) + .begin_transaction() .await - .map_err(|e| internal_error(&format!("list restrictions: {e}")))?; - Ok(Json(Value::Array(rows.iter().map(ban_json).collect()))) + .map_err(|error| internal_error(&format!("begin moderation read: {error}")))?; + crate::handlers::moderation_authz::authorize_moderation_action_tx( + &mut transaction, + &authorization.tenant, + &authorization.actor_pubkey, + None, + crate::handlers::moderation_authz::ModerationTarget::None, + crate::handlers::moderation_authz::ModerationAction::ViewQueue, + ) + .await + .map_err(|_| { + api_error( + StatusCode::FORBIDDEN, + "restricted: moderator access required", + ) + })?; + let fetched = + buzz_db::moderation::list_restricted_tx(&mut transaction, authorization.tenant.community()) + .await; + let rows = release_protected_bridge_fetch(fetched, |fetched| { + authorization.protected.release_fetched(fetched) + })? + .map_err(|e| internal_error(&format!("list restrictions: {e}")))?; + let response = Json(Value::Array(rows.iter().map(ban_json).collect())); + transaction + .commit() + .await + .map_err(|error| internal_error(&format!("finish moderation read: {error}")))?; + release_moderation_read(&state, &authorization).await?; + Ok(response) } fn report_json(r: &buzz_db::moderation::ReportRecord) -> Value { @@ -2351,6 +2941,61 @@ mod tests { .to_bytes() } + #[test] + fn bridge_release_retains_stored_channel_without_trusting_event_tags() { + let event = EventBuilder::new(Kind::TextNote, "channel-scoped") + .sign_with_keys(&Keys::generate()) + .expect("sign event without h tag"); + let channel_id = uuid::Uuid::new_v4(); + let stored = buzz_core::StoredEvent::new(event, Some(channel_id)); + let mut events = Vec::new(); + let mut release_channels = std::collections::BTreeSet::new(); + + assert!(append_bridge_stored_event( + &mut events, + &mut release_channels, + &stored, + )); + assert_eq!(events.len(), 1); + assert_eq!( + release_channels.into_iter().collect::>(), + [channel_id] + ); + } + + #[test] + fn observational_http_modes_use_read_only_identity_lane() { + use crate::authorization_runtime::{ + finalization::AuthorizationMode, + transport::{legacy_identity_lane_for_mode, LegacyIdentityLane}, + }; + + for mode in [AuthorizationMode::Shadow, AuthorizationMode::VerifyOnly] { + assert_eq!( + legacy_identity_lane_for_mode(Some(mode)), + LegacyIdentityLane::ObserveOnly + ); + } + assert_eq!( + legacy_identity_lane_for_mode(Some(AuthorizationMode::Enforce)), + LegacyIdentityLane::ProtectedEnforce + ); + } + + #[test] + fn protected_bridge_fetch_release_fences_success_and_backend_error_outcomes() { + let released = release_protected_bridge_fetch::(Ok(7), Ok) + .expect("current authority releases fetched rows") + .expect("successful fetch remains successful"); + assert_eq!(released, 7); + + for fetched in [Ok(7), Err("database unavailable")] { + let (status, _) = release_protected_bridge_fetch::(fetched, |_| Err(())) + .expect_err("authority loss must hide every fetched outcome"); + assert_eq!(status, StatusCode::FORBIDDEN); + } + } + #[test] fn corporate_identity_disables_x_pubkey_bridge_fallback() { let keys = Keys::generate(); @@ -3153,13 +3798,6 @@ mod tests { assert_eq!(extract_page_offset(&raw, None), None); } - /// Offsets are sized from the *clamped* limit the DB will honor, not from - /// what the client asked for. `filter_to_query_params` clamps an absent or - /// over-ceiling `limit` to `DEFAULT_MAX_PAGE_LIMIT` (guarded in - /// `handlers::req::tests::req_filter_limit_clamps_to_advertised_nip11_max_limit`) - /// and that clamped value is what arrives here — so page N starts exactly - /// N-1 full pages in. Sizing from an unclamped limit would step past rows - /// the previous page never returned. #[test] fn extract_page_offset_sizes_pages_from_clamped_limit() { let clamped = buzz_db::DEFAULT_MAX_PAGE_LIMIT; @@ -3488,7 +4126,10 @@ mod tests { require_corporate_identity: bool, ) -> Option> { let mut config = crate::config::Config::from_env().ok()?; - config.database_url = TEST_DB_URL.to_string(); + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| TEST_DB_URL.to_owned()); + config.database_url = database_url.clone(); // Use the real local Redis so enforce_http_admission can pass. config.redis_url = std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); @@ -3502,7 +4143,7 @@ mod tests { config.corporate_identity.audience = "buzz-relay".to_string(); } - let pool = sqlx::PgPool::connect(TEST_DB_URL).await.ok()?; + let pool = sqlx::PgPool::connect(&database_url).await.ok()?; let db = buzz_db::Db::from_pool(pool.clone()); let redis_pool = deadpool_redis::Config::from_url(&config.redis_url) .create_pool(Some(deadpool_redis::Runtime::Tokio1)) @@ -3871,8 +4512,8 @@ mod tests { "expected exactly 1 attribution line for invalid-JSON arm, got {n};\nlog:\n{log}" ); assert!( - log.contains(&pubkey_hex[..16]), - "attribution line must carry the pubkey;\nlog:\n{log}" + !log.contains(&pubkey_hex[..16]), + "runtime log exposed pubkey;\nlog:\n{log}" ); } @@ -3931,8 +4572,8 @@ mod tests { "expected exactly 1 attribution line for IngestError::Rejected arm, got {n};\nlog:\n{log}" ); assert!( - log.contains(&pubkey_hex[..16]), - "attribution line must carry the pubkey;\nlog:\n{log}" + !log.contains(&pubkey_hex[..16]), + "runtime log exposed pubkey;\nlog:\n{log}" ); } } diff --git a/crates/buzz-relay/src/authorization_runtime/transport.rs b/crates/buzz-relay/src/authorization_runtime/transport.rs new file mode 100644 index 000000000..d497b4fb4 --- /dev/null +++ b/crates/buzz-relay/src/authorization_runtime/transport.rs @@ -0,0 +1,1420 @@ +//! Provider-neutral protected-transport authorization. +//! +//! Transport handlers request one portable capability at the point of use. +//! An injected resolver owns provider evaluation, binding/admission ordering, +//! and finalization. This module accepts only a finalized access context as +//! authority, binds it back to the exact operation, and retains a guard for +//! mandatory pre-commit or pre-emission revalidation. + +use std::{collections::HashMap, fmt, sync::Arc}; + +use async_trait::async_trait; +use buzz_auth::{ + AuthContext, AuthTransport, AuthorizationCapability, AuthorizationLease, + AuthorizationLeaseValidator, AuthorizationProfileId, BindingVersion, LeaseUseRequirement, + LeaseValidationError, PolicyVersion, SharedAuthorizationClock, VerificationOnlyDisposition, + VerifiedFederatedAssertion, VerifiedNostrProof, +}; +use buzz_core::CommunityId; +use thiserror::Error; +use tokio_util::sync::CancellationToken; +use uuid::Uuid; + +use super::finalization::{AuthorizationMode, EnrollmentDisposition}; + +/// Compatibility policy for the inherited corporate-identity lane. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum LegacyIdentityLane { + /// Runtime absent or explicitly Off: run the inherited verifier and + /// finalizer without changing pre-O4 behavior. + Legacy, + /// Shadow or VerifyOnly: cryptographic verification and read-only policy + /// evaluation are allowed, but identity state and public projections are + /// immutable. This lane can never enroll, reactivate, strengthen, update, + /// or retire a binding. + ObserveOnly, + /// Enforce: the protected resolver owns admission/finalization and legacy + /// projection must not run. + ProtectedEnforce, +} + +/// Resolve the centralized identity compatibility policy for one exact domain. +pub fn legacy_identity_lane( + state: &crate::state::AppState, + authorization_domain: CommunityId, +) -> LegacyIdentityLane { + legacy_identity_lane_for_mode( + state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(authorization_domain)), + ) +} + +/// Resolve compatibility behavior from an already selected exact-domain mode. +pub const fn legacy_identity_lane_for_mode(mode: Option) -> LegacyIdentityLane { + match mode { + None | Some(AuthorizationMode::Off) => LegacyIdentityLane::Legacy, + Some(AuthorizationMode::Shadow) | Some(AuthorizationMode::VerifyOnly) => { + LegacyIdentityLane::ObserveOnly + } + Some(AuthorizationMode::Enforce) => LegacyIdentityLane::ProtectedEnforce, + } +} + +/// Server-owned activation for one exact authorization domain. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct DomainTransportPolicy { + authorization_domain: CommunityId, + mode: AuthorizationMode, +} + +impl DomainTransportPolicy { + /// Construct policy from immutable server configuration. + pub const fn from_server_configuration( + authorization_domain: CommunityId, + mode: AuthorizationMode, + ) -> Self { + Self { + authorization_domain, + mode, + } + } + + /// Exact configured domain. + pub const fn authorization_domain(self) -> CommunityId { + self.authorization_domain + } + + /// Exact configured activation mode. + pub const fn mode(self) -> AuthorizationMode { + self.mode + } +} + +/// Exact protected operation presented to the configured resolver. +#[derive(Clone)] +pub struct ProtectedOperationRequest { + verified_proof: Arc, + capability: AuthorizationCapability, + correlation_id: Uuid, + session_id: Option, + surface: &'static str, + cancellation: Option, + verified_assertion: Option>, + enrollment_assertion: Option>, +} + +impl ProtectedOperationRequest { + /// Build one exact operation after transport authentication. + pub fn new( + verified_proof: Arc, + verified_assertion: Option>, + capability: AuthorizationCapability, + correlation_id: Uuid, + surface: &'static str, + ) -> Result { + Self::new_with_cancellation( + verified_proof, + verified_assertion, + capability, + correlation_id, + surface, + None, + None, + ) + } + + pub(crate) fn new_with_cancellation( + verified_proof: Arc, + verified_assertion: Option>, + capability: AuthorizationCapability, + correlation_id: Uuid, + surface: &'static str, + session_id: Option, + cancellation: Option, + ) -> Result { + if correlation_id.is_nil() { + return Err(ProtectedTransportError::InvalidCorrelationId); + } + if session_id == Some(Uuid::nil()) { + return Err(ProtectedTransportError::InvalidSessionId); + } + if surface.is_empty() { + return Err(ProtectedTransportError::InvalidSurface); + } + if !crate::protected_surface::protected_operation_matches( + surface, + verified_proof.authorized_transport(), + capability, + ) { + return Err(ProtectedTransportError::SurfaceCapabilityMismatch); + } + Ok(Self { + verified_proof, + capability, + correlation_id, + session_id, + surface, + cancellation, + verified_assertion, + enrollment_assertion: None, + }) + } + + fn new_enrollment( + verified_proof: Arc, + assertion: Arc, + correlation_id: Uuid, + surface: &'static str, + ) -> Result { + let mut request = Self::new( + verified_proof, + Some(Arc::clone(&assertion)), + AuthorizationCapability::InviteClaim, + correlation_id, + surface, + )?; + request.enrollment_assertion = Some(assertion); + Ok(request) + } + + /// Exact server-resolved authorization domain. + pub fn authorization_domain(&self) -> CommunityId { + self.verified_proof.authorization_domain() + } + + /// Transport carrying this operation. + pub fn transport(&self) -> AuthTransport { + self.verified_proof.authorized_transport() + } + + /// Authenticated actor. + pub fn actor_pubkey(&self) -> nostr::PublicKey { + self.verified_proof.actor_pubkey() + } + + /// Verified Nostr owner for delegated authority, when present. + pub fn owner_pubkey(&self) -> Option { + self.verified_proof + .verified_delegation() + .map(buzz_auth::VerifiedTransportDelegation::owner_pubkey) + } + + /// Sealed verifier evidence for this exact transport request or session. + pub fn verified_proof(&self) -> &Arc { + &self.verified_proof + } + + /// Direct verified assertion retained only for first-enrollment resolution. + pub fn enrollment_assertion(&self) -> Option<&Arc> { + self.enrollment_assertion.as_ref() + } + + /// Current direct assertion verified by the transport identity adapter. + pub fn verified_assertion(&self) -> Option<&Arc> { + self.verified_assertion.as_ref() + } + + /// Exact dynamically selected capability. + pub const fn capability(&self) -> AuthorizationCapability { + self.capability + } + + /// Correlation identifier for this decision. + pub const fn correlation_id(&self) -> Uuid { + self.correlation_id + } + + /// Stable server-owned session identity for long-lived transports. + pub const fn session_id(&self) -> Option { + self.session_id + } + + /// Stable low-cardinality surface name for resolver telemetry. + pub const fn surface(&self) -> &'static str { + self.surface + } + + /// Server-owned cancellation for a leased session, when applicable. + pub fn cancellation(&self) -> Option { + self.cancellation.clone() + } +} + +impl fmt::Debug for ProtectedOperationRequest { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProtectedOperationRequest") + .field("authorization_domain", &"[redacted]") + .field("transport", &self.transport()) + .field("actor_pubkey", &"[redacted]") + .field("owner_pubkey", &"[redacted]") + .field("capability", &"[redacted]") + .field("correlation_id", &"[redacted]") + .field("surface", &self.surface) + .finish() + } +} + +/// Resolver failure that carries only a stable, non-sensitive code. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +#[error("protected authorization resolver denied the operation ({code})")] +pub struct ProtectedResolutionError { + code: &'static str, +} + +impl ProtectedResolutionError { + /// Construct a sanitized resolver denial. + pub const fn new(code: &'static str) -> Self { + Self { code } + } + + /// Stable denial code suitable for metrics and protocol mapping. + pub const fn code(self) -> &'static str { + self.code + } +} + +/// Exact-domain resolver implemented by the authorization adapter. +/// +/// Direct and delegated actors intentionally enter through the same method. +/// The request's verified owner is evidence, never a separate bypass path. +#[async_trait] +pub trait ProtectedAuthorizationResolver: Send + Sync { + /// Resolve and fully finalize current authority for one operation. + async fn resolve( + &self, + request: &ProtectedOperationRequest, + ) -> Result; + + /// Run a read-only observation without producing authority or mutating + /// identity, binding, membership, projection, or lease state. + async fn observe( + &self, + _request: &ProtectedOperationRequest, + ) -> Result<(), ProtectedResolutionError> { + Err(ProtectedResolutionError::new( + "observational_provider_unavailable", + )) + } + + /// Resolve display-only current binding status without granting access or + /// mutating binding, membership, projection, or lease state. + async fn present( + &self, + _request: &ProtectedOperationRequest, + ) -> Result { + Err(ProtectedResolutionError::new( + "client_status_presentation_unavailable", + )) + } +} + +/// Display-only result retained with its exact invalidation observer. +pub struct ProtectedStatusResolution { + disposition: VerificationOnlyDisposition, + observer: Arc, + evaluation_generation: u64, +} + +impl ProtectedStatusResolution { + /// Couple one current status to the invalidation fence captured for it. + pub fn new( + disposition: VerificationOnlyDisposition, + observer: Arc, + evaluation_generation: u64, + ) -> Self { + Self { + disposition, + observer, + evaluation_generation, + } + } + + /// Consume the display result for dedicated delivery and reconciliation. + pub(crate) fn into_parts( + self, + ) -> ( + VerificationOnlyDisposition, + Arc, + u64, + ) { + (self.disposition, self.observer, self.evaluation_generation) + } +} + +impl fmt::Debug for ProtectedStatusResolution { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProtectedStatusResolution") + .field("disposition", &"[redacted]") + .field("observer", &"[registered]") + .field("evaluation_generation", &"[redacted]") + .finish() + } +} + +/// Finalized resolver output. +/// +/// Fields are private so enforcing access cannot be constructed without the +/// exact per-decision observer registered during finalization. +pub struct ProtectedResolution { + kind: ProtectedResolutionKind, +} + +enum ProtectedResolutionKind { + Access { + context: Box, + observer: Arc, + }, + VerificationOnly(VerificationOnlyDisposition), + Enrollment { + disposition: EnrollmentDisposition, + observer: Arc, + }, +} + +impl ProtectedResolution { + /// Couple finalized access with its exact read-only invalidation observer. + /// + /// Enforcing access has no observer-free constructor: + /// + /// ```compile_fail + /// use buzz_auth::AuthContext; + /// use buzz_relay::authorization_runtime::transport::ProtectedResolution; + /// + /// fn invalid(context: Box) -> ProtectedResolution { + /// ProtectedResolution::access(context) + /// } + /// ``` + pub fn access(context: Box, observer: Arc) -> Self { + Self { + kind: ProtectedResolutionKind::Access { context, observer }, + } + } + + /// Preserve a display-only result without manufacturing access state. + pub fn verification_only(disposition: VerificationOnlyDisposition) -> Self { + Self { + kind: ProtectedResolutionKind::VerificationOnly(disposition), + } + } + + /// Couple staged direct enrollment with its exact invalidation observer. + pub fn enrollment( + disposition: EnrollmentDisposition, + observer: Arc, + ) -> Self { + Self { + kind: ProtectedResolutionKind::Enrollment { + disposition, + observer, + }, + } + } +} + +impl fmt::Debug for ProtectedResolution { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProtectedResolution") + .field("kind", &"[redacted]") + .finish() + } +} + +/// Current trusted state for every dependency that can invalidate a lease. +#[derive(Clone, PartialEq, Eq)] +pub struct LeaseCurrentState { + binding_version: BindingVersion, + profile_id: AuthorizationProfileId, + policy_version: PolicyVersion, +} + +impl LeaseCurrentState { + /// Construct state from the trusted finalization/invalidation runtime. + pub const fn from_trusted_runtime( + binding_version: BindingVersion, + profile_id: AuthorizationProfileId, + policy_version: PolicyVersion, + ) -> Self { + Self { + binding_version, + profile_id, + policy_version, + } + } +} + +impl fmt::Debug for LeaseCurrentState { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LeaseCurrentState") + .field("binding_version", &"[redacted]") + .field("profile_id", &"[redacted]") + .field("policy_version", &"[redacted]") + .finish() + } +} + +/// Exact per-decision observer used by a retained lease guard. +/// +/// Implementations retain the registered invalidation fence and session +/// dependencies for this exact finalization. They must check those dependencies +/// and read current binding/profile/policy state before returning. The observer +/// is intentionally read-only; transport adoption cannot publish invalidations. +pub trait LeaseCurrentStateObserver: Send + Sync { + /// Validate the registered fence/dependencies and return current versions. + fn observe_current(&self) -> Result; + + /// Return the captured durable generation and exact selector dependencies + /// required by a transaction-owned mutation commit. + /// + /// Read-only observers may retain the default fail-closed implementation. + fn observe_commit_fence( + &self, + ) -> Result { + Err(LeaseCurrentStateError::Unavailable) + } +} + +/// Fail-closed current-state result. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum LeaseCurrentStateError { + /// Binding, profile, or policy state changed. + #[error("protected authorization state is stale")] + Stale, + /// Current state could not be read. + #[error("protected authorization state is unavailable")] + Unavailable, +} + +/// Optional protected-transport runtime. +/// +/// Domains absent from this map, and domains explicitly configured `Off`, use +/// legacy behavior. An enforcing domain never falls back after a resolver or +/// lease failure. +pub struct ProtectedTransportRuntime { + domains: HashMap, + resolver: Arc, + validator: AuthorizationLeaseValidator, +} + +impl ProtectedTransportRuntime { + /// Build an exact-domain runtime, rejecting duplicate configuration. + pub fn new( + policies: impl IntoIterator, + resolver: Arc, + clock: SharedAuthorizationClock, + ) -> Result { + let mut domains = HashMap::new(); + for policy in policies { + if domains + .insert(policy.authorization_domain, policy.mode) + .is_some() + { + return Err(ProtectedTransportError::AmbiguousDomainPolicy); + } + } + Ok(Self { + domains, + resolver, + validator: AuthorizationLeaseValidator::new(clock), + }) + } + + /// Return the exact configured mode, if this runtime owns the domain. + pub fn mode_for_domain(&self, domain: CommunityId) -> Option { + self.domains.get(&domain).copied() + } + + /// Exact domains whose background effects must remain unavailable. + pub fn enforcing_domains(&self) -> Vec { + self.domains + .iter() + .filter_map(|(domain, mode)| (*mode == AuthorizationMode::Enforce).then_some(*domain)) + .collect() + } + + /// Resolve one protected operation or explicitly preserve legacy behavior. + pub async fn authorize( + &self, + request: &ProtectedOperationRequest, + ) -> Result { + let Some(mode) = self.mode_for_domain(request.authorization_domain()) else { + return Ok(ProtectedAuthorization::Legacy); + }; + match mode { + AuthorizationMode::Off => Ok(ProtectedAuthorization::Legacy), + // Transport adoption has no mutation-free decision-only resolver. + // Preserve legacy access in both observational modes and leave + // their status evaluation to the lane that owns that API. + AuthorizationMode::Shadow | AuthorizationMode::VerifyOnly => { + // Observation failures are telemetry only. These modes must + // never alter the inherited access result. + let _ = self.resolver.observe(request).await; + Ok(ProtectedAuthorization::Legacy) + } + AuthorizationMode::Enforce => { + let resolution = self + .resolver + .resolve(request) + .await + .map_err(ProtectedTransportError::Resolution)?; + let (context, observer) = match resolution.kind { + ProtectedResolutionKind::Access { context, observer } => (context, observer), + ProtectedResolutionKind::VerificationOnly(_status) => { + return Err(ProtectedTransportError::VerificationOnlyCannotGrant) + } + ProtectedResolutionKind::Enrollment { .. } => { + return Err(ProtectedTransportError::EnrollmentCannotGrantAccess) + } + }; + let context: Arc = Arc::from(context); + let authority = ProtectedOperationAuthority { + context, + capability: request.capability(), + validator: self.validator.clone(), + observer, + }; + authority.validate_exact_request(request)?; + authority.revalidate()?; + Ok(ProtectedAuthorization::Access(authority)) + } + } + } + + /// Resolve a dedicated client presentation only in VerifyOnly or Enforce. + /// Off and Shadow remain byte-for-byte non-presenting legacy behavior. + pub async fn present_status( + &self, + request: &ProtectedOperationRequest, + ) -> Result, ProtectedTransportError> { + match self.mode_for_domain(request.authorization_domain()) { + Some(AuthorizationMode::VerifyOnly | AuthorizationMode::Enforce) => self + .resolver + .present(request) + .await + .map(Some) + .map_err(ProtectedTransportError::Resolution), + None | Some(AuthorizationMode::Off | AuthorizationMode::Shadow) => Ok(None), + } + } + + /// Resolve direct first-enrollment evidence without requiring an active binding. + pub async fn authorize_enrollment( + &self, + request: &ProtectedOperationRequest, + ) -> Result { + let Some(mode) = self.mode_for_domain(request.authorization_domain()) else { + return Ok(ProtectedEnrollmentAuthorization::Legacy); + }; + match mode { + AuthorizationMode::Off | AuthorizationMode::Shadow | AuthorizationMode::VerifyOnly => { + Ok(ProtectedEnrollmentAuthorization::Legacy) + } + AuthorizationMode::Enforce => { + if request.capability() != AuthorizationCapability::InviteClaim + || request.enrollment_assertion().is_none() + || request.owner_pubkey().is_some() + { + return Err(ProtectedTransportError::EnrollmentEvidenceRequired); + } + let resolution = self + .resolver + .resolve(request) + .await + .map_err(ProtectedTransportError::Resolution)?; + let (disposition, observer) = match resolution.kind { + ProtectedResolutionKind::Enrollment { + disposition, + observer, + } => (disposition, observer), + _ => return Err(ProtectedTransportError::EnrollmentEvidenceRequired), + }; + let authority = ProtectedEnrollmentAuthority { + disposition, + observer, + validator: self.validator.clone(), + }; + authority.validate_exact_request(request)?; + authority.revalidate()?; + Ok(ProtectedEnrollmentAuthorization::Enrollment(authority)) + } + } + } +} + +/// Consult the optional runtime for one authenticated relay operation. +/// +/// This is the common transport seam used by HTTP, WebSocket, Git, media, and +/// audio handlers. Runtime absence is the only implicit legacy case. +pub async fn authorize_if_configured( + state: &crate::state::AppState, + verified_proof: Arc, + verified_assertion: Option>, + capability: AuthorizationCapability, + correlation_id: Uuid, + surface: &'static str, +) -> Result { + let Some(runtime) = state.protected_transport() else { + return Ok(ProtectedAuthorization::Legacy); + }; + let request = ProtectedOperationRequest::new( + verified_proof, + verified_assertion, + capability, + correlation_id, + surface, + )?; + runtime.authorize(&request).await +} + +/// Consult the runtime for a leased session and expose only its server-owned +/// cancellation token to the resolver's invalidation registration. +#[allow(clippy::too_many_arguments)] +pub async fn authorize_session_if_configured( + state: &crate::state::AppState, + verified_proof: Arc, + verified_assertion: Option>, + capability: AuthorizationCapability, + correlation_id: Uuid, + surface: &'static str, + session_id: Uuid, + cancellation: CancellationToken, +) -> Result { + let Some(runtime) = state.protected_transport() else { + return Ok(ProtectedAuthorization::Legacy); + }; + let request = ProtectedOperationRequest::new_with_cancellation( + verified_proof, + verified_assertion, + capability, + correlation_id, + surface, + Some(session_id), + Some(cancellation), + )?; + let authority = runtime.authorize(&request).await?; + state + .conn_manager + .retain_protected_session_authority(session_id, &authority); + Ok(authority) +} + +/// Resolve staged direct authority for atomic invite enrollment. +pub async fn authorize_enrollment_if_configured( + state: &crate::state::AppState, + verified_proof: Arc, + assertion: Arc, + correlation_id: Uuid, +) -> Result { + let Some(runtime) = state.protected_transport() else { + return Ok(ProtectedEnrollmentAuthorization::Legacy); + }; + let request = ProtectedOperationRequest::new_enrollment( + verified_proof, + assertion, + correlation_id, + "invite.claim", + )?; + runtime.authorize_enrollment(&request).await +} + +/// Preserve legacy behavior only when an unwired surface cannot enter an +/// enforcing exact-domain runtime. +/// +/// This helper never constructs a resolver request from caller-supplied key +/// primitives. It exists solely to make missing typed evidence fail closed on +/// enforcing domains while retaining the configured Off/Shadow/VerifyOnly +/// behavior for legacy-only authentication paths. +pub fn authorize_unwired_if_configured( + state: &crate::state::AppState, + authorization_domain: CommunityId, +) -> Result { + authorize_unwired_for_mode( + state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(authorization_domain)), + ) +} + +/// Deny an unproved mutation in Enforce before any backend-visible work starts. +/// +/// Alternate helpers such as secret-triggered workflow execution have no +/// request proof that a protected resolver can finalize. Until a transaction- +/// owning executor carries durable authority into their commit, they must be +/// unavailable in Enforce while preserving every legacy lane. +pub fn require_unwired_atomic_mutation_if_configured( + state: &crate::state::AppState, + authorization_domain: CommunityId, +) -> Result<(), ProtectedTransportError> { + authorize_unwired_if_configured(state, authorization_domain)?.require_atomic_mutation() +} + +fn authorize_unwired_for_mode( + mode: Option, +) -> Result { + match mode { + None + | Some(AuthorizationMode::Off) + | Some(AuthorizationMode::Shadow) + | Some(AuthorizationMode::VerifyOnly) => Ok(ProtectedAuthorization::Legacy), + Some(AuthorizationMode::Enforce) => Err(ProtectedTransportError::MissingVerifiedProof), + } +} + +impl fmt::Debug for ProtectedTransportRuntime { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProtectedTransportRuntime") + .field("domains", &"[redacted]") + .field("resolver", &"[configured]") + .field("validator", &self.validator) + .finish() + } +} + +/// Result of consulting the direct first-enrollment runtime. +#[must_use] +pub enum ProtectedEnrollmentAuthorization { + /// Runtime absent or configured non-enforcing. + Legacy, + /// Enforcing staged enrollment with no binding or membership yet. + Enrollment(ProtectedEnrollmentAuthority), +} + +impl ProtectedEnrollmentAuthorization { + /// Seal the enrollment decision for transaction-owned execution. + pub fn seal_postgres_enrollment( + &self, + operation_id: super::executor::ProtectedOperationId, + operation_kind: &'static str, + request_fingerprint: [u8; 32], + ) -> Result, ProtectedTransportError> { + match self { + Self::Legacy => Ok(None), + Self::Enrollment(authority) => super::executor::SealedEnrollmentPermit::from_authority( + authority, + operation_id, + operation_kind, + request_fingerprint, + ) + .map(Some), + } + } + + /// Whether this decision is enforcing. + pub const fn is_enforcing(&self) -> bool { + matches!(self, Self::Enrollment(_)) + } +} + +/// Staged direct authority that cannot be consumed as ordinary access. +pub struct ProtectedEnrollmentAuthority { + disposition: EnrollmentDisposition, + observer: Arc, + validator: AuthorizationLeaseValidator, +} + +impl ProtectedEnrollmentAuthority { + pub(super) const fn disposition(&self) -> &EnrollmentDisposition { + &self.disposition + } + + pub(super) fn observer(&self) -> &dyn LeaseCurrentStateObserver { + self.observer.as_ref() + } + + fn validate_exact_request( + &self, + request: &ProtectedOperationRequest, + ) -> Result<(), ProtectedTransportError> { + if self.disposition.authorization_domain() != request.authorization_domain() + || self.disposition.actor_pubkey() != request.actor_pubkey() + || self.disposition.correlation_id() != request.correlation_id() + || request.capability() != AuthorizationCapability::InviteClaim + { + return Err(ProtectedTransportError::FinalizedContextMismatch); + } + Ok(()) + } + + /// Recheck time and the captured invalidation fence before sealing. + pub fn revalidate(&self) -> Result<(), ProtectedTransportError> { + self.validator + .seconds_until(self.disposition.expires_at())?; + self.observer.observe_commit_fence()?; + Ok(()) + } +} + +impl fmt::Debug for ProtectedEnrollmentAuthority { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("ProtectedEnrollmentAuthority([redacted])") + } +} + +/// Result of consulting the optional exact-domain runtime. +#[derive(Clone)] +#[must_use] +pub enum ProtectedAuthorization { + /// Runtime absent for this exact domain or configured non-enforcing. + Legacy, + /// Enforcing access with a retained operation guard. + Access(ProtectedOperationAuthority), +} + +impl ProtectedAuthorization { + /// Revalidate enforcing authority; legacy mode is a no-op. + pub fn revalidate(&self) -> Result<(), ProtectedTransportError> { + match self { + Self::Legacy => Ok(()), + Self::Access(authority) => authority.revalidate(), + } + } + + /// Release already-fetched protected data only after a final authority check. + /// + /// Callers place this at the last synchronous boundary before constructing + /// a response, queue entry, or stream item. Keeping the fetched value owned + /// by this method makes the post-fetch/pre-emission ordering explicit. + pub fn release_fetched(&self, value: T) -> Result { + self.revalidate()?; + Ok(value) + } + + /// Require an atomic mutation coordinator for an enforcing operation. + /// + /// Legacy lanes retain their exact behavior. Enforcing callers must select + /// the transaction/CAS executor for their registered surface; an adjacent + /// authorization check is never treated as an atomic fence. + pub fn require_atomic_mutation(&self) -> Result<(), ProtectedTransportError> { + require_atomic_mutation_executor(matches!(self, Self::Access(_))) + } + + /// Seal an enforcing mutation for transaction-owned PostgreSQL execution. + /// + /// Legacy lanes return `None` and retain their existing mutation path. + /// Enforcing lanes return a permit that cannot be constructed from request + /// fields or an adjacent authorization check. + pub fn seal_postgres_mutation( + &self, + operation_id: super::executor::ProtectedOperationId, + operation_kind: &'static str, + request_fingerprint: [u8; 32], + ) -> Result, ProtectedTransportError> { + match self { + Self::Legacy => Ok(None), + Self::Access(authority) => super::executor::SealedOperationPermit::from_authority( + authority, + operation_id, + operation_kind, + request_fingerprint, + ) + .map(Some), + } + } + + /// Seal opaque sender authority for one ephemeral event delivered across + /// the trusted relay Redis plane. The resulting claim still requires a + /// relay signature and writer-database revalidation on the receiving pod. + pub(crate) fn seal_ephemeral_delivery( + &self, + event_id: [u8; 32], + ) -> Result, ProtectedTransportError> { + match self { + Self::Legacy => Ok(None), + Self::Access(authority) => { + super::executor::EphemeralAuthorityClaim::from_authority(authority, event_id) + .map(Some) + } + } + } + + /// Whether this value carries enforcing authority rather than legacy mode. + pub const fn is_enforcing(&self) -> bool { + matches!(self, Self::Access(_)) + } + + /// Hard lease expiry for enforcing authority. + pub fn expires_at(&self) -> Option { + match self { + Self::Legacy => None, + Self::Access(authority) => Some(authority.expires_at()), + } + } + + /// Delay to hard expiry using the same injected clock as lease validation. + pub fn expiry_delay(&self) -> Result, ProtectedTransportError> { + match self { + Self::Legacy => Ok(None), + Self::Access(authority) => authority.expiry_delay().map(Some), + } + } +} + +fn require_atomic_mutation_executor(enforcing: bool) -> Result<(), ProtectedTransportError> { + if enforcing { + Err(ProtectedTransportError::AtomicMutationFenceUnavailable) + } else { + Ok(()) + } +} + +impl fmt::Debug for ProtectedAuthorization { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Legacy => formatter.write_str("ProtectedAuthorization::Legacy"), + Self::Access(_) => formatter.write_str("ProtectedAuthorization::Access([redacted])"), + } + } +} + +/// Retained authority for request, stream, and pre-commit checkpoints. +#[derive(Clone)] +pub struct ProtectedOperationAuthority { + context: Arc, + capability: AuthorizationCapability, + validator: AuthorizationLeaseValidator, + observer: Arc, +} + +impl ProtectedOperationAuthority { + pub(super) fn context(&self) -> &AuthContext { + &self.context + } + + pub(super) fn observer(&self) -> &dyn LeaseCurrentStateObserver { + self.observer.as_ref() + } + + fn lease(&self) -> Result<&AuthorizationLease, ProtectedTransportError> { + self.context + .authorization_lease() + .ok_or(ProtectedTransportError::MissingAccessLease) + } + + fn requirement( + &self, + current: &LeaseCurrentState, + ) -> Result { + let lease = self.lease()?; + Ok(LeaseUseRequirement { + context_version: lease.context_version(), + lease_version: lease.lease_version(), + authorization_domain: lease.authorization_domain(), + transport: lease.transport(), + actor_pubkey: lease.actor_pubkey(), + binding_id: lease.binding_id(), + binding_version: current.binding_version, + profile_id: current.profile_id.clone(), + policy_version: current.policy_version.clone(), + capability: self.capability, + }) + } + + fn validate_exact_request( + &self, + request: &ProtectedOperationRequest, + ) -> Result<(), ProtectedTransportError> { + if self.context.tenant().community() != request.authorization_domain() + || self.context.transport() != request.transport() + || self.context.pubkey() != request.actor_pubkey() + || self.context.agent_owner_pubkey() != request.owner_pubkey() + || !exact_correlation_matches(self.context.correlation_id(), request.correlation_id()) + { + return Err(ProtectedTransportError::FinalizedContextMismatch); + } + Ok(()) + } + + /// Revalidate current state, exact capability, typed versions, and expiry. + pub fn revalidate(&self) -> Result<(), ProtectedTransportError> { + let current = self.observer.observe_current()?; + self.context + .authorize_lease_use(&self.validator, &self.requirement(¤t)?)?; + Ok(()) + } + + /// Exact operation capability retained by this authority. + pub const fn capability(&self) -> AuthorizationCapability { + self.capability + } + + /// Hard conservative expiry after which revalidation fails. + pub fn expires_at(&self) -> u64 { + self.context + .authorization_lease() + .map_or(0, AuthorizationLease::expires_at) + } + + fn expiry_delay(&self) -> Result { + let expires_at = self.expires_at(); + let seconds = self.validator.seconds_until(expires_at)?; + Ok(std::time::Duration::from_secs(seconds)) + } +} + +fn exact_correlation_matches(finalized: Uuid, requested: Uuid) -> bool { + finalized == requested +} + +impl fmt::Debug for ProtectedOperationAuthority { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProtectedOperationAuthority") + .field("context", &"[redacted]") + .field("capability", &"[redacted]") + .field("validator", &self.validator) + .field("observer", &"[configured]") + .finish() + } +} + +/// Fail-closed protected-transport error. +#[derive(Debug, Error)] +pub enum ProtectedTransportError { + /// Duplicate exact-domain activation is ambiguous. + #[error("protected authorization policy is ambiguous for this domain")] + AmbiguousDomainPolicy, + /// Correlation IDs must be non-nil. + #[error("protected authorization correlation id is invalid")] + InvalidCorrelationId, + /// Session IDs for long-lived transports must be non-nil. + #[error("protected authorization session id is invalid")] + InvalidSessionId, + /// Surface labels must be stable non-empty server constants. + #[error("protected authorization surface is invalid")] + InvalidSurface, + /// The surface, sealed proof transport, and capability do not match. + #[error("protected authorization surface does not match proof and capability")] + SurfaceCapabilityMismatch, + /// An enforcing surface did not retain sealed verifier evidence. + #[error("protected authorization requires verified transport evidence")] + MissingVerifiedProof, + /// Resolver denied or could not evaluate current policy. + #[error(transparent)] + Resolution(#[from] ProtectedResolutionError), + /// Display verification is deliberately not access authority. + #[error("verification-only authorization cannot grant protected access")] + VerificationOnlyCannotGrant, + /// A staged enrollment result was presented as ordinary access. + #[error("first-enrollment authorization cannot grant ordinary protected access")] + EnrollmentCannotGrantAccess, + /// Invite enrollment lacked matching direct assertion/provider evidence. + #[error("protected invite enrollment requires direct verified evidence")] + EnrollmentEvidenceRequired, + /// Resolver returned a context for another exact operation. + #[error("finalized authorization context does not match protected operation")] + FinalizedContextMismatch, + /// Enforcing disposition lacked its bounded lease. + #[error("enforcing authorization context is missing its access lease")] + MissingAccessLease, + /// Current binding/profile/policy observation failed closed. + #[error(transparent)] + CurrentState(#[from] LeaseCurrentStateError), + /// Typed lease validation failed. + #[error(transparent)] + Lease(#[from] LeaseValidationError), + /// No transaction- or CAS-coupled mutation coordinator owns this commit. + #[error("protected mutation requires an atomic authorization fence")] + AtomicMutationFenceUnavailable, +} + +#[cfg(test)] +mod tests { + use std::sync::atomic::{AtomicUsize, Ordering}; + + use buzz_auth::{ + AuthorizationClock, AuthorizationClockError, AuthorizationTime, VerifiedEvidenceAdapter, + }; + use nostr::{EventBuilder, Keys, RelayUrl}; + + use super::*; + + struct UnavailableResolver; + + #[async_trait] + impl ProtectedAuthorizationResolver for UnavailableResolver { + async fn resolve( + &self, + _request: &ProtectedOperationRequest, + ) -> Result { + Err(ProtectedResolutionError::new("synthetic_unavailable")) + } + } + + struct PanicResolver; + + #[async_trait] + impl ProtectedAuthorizationResolver for PanicResolver { + async fn resolve( + &self, + _request: &ProtectedOperationRequest, + ) -> Result { + panic!("non-enforcing mode must not call a finalizing resolver") + } + } + + struct ObservingResolver(AtomicUsize); + + #[async_trait] + impl ProtectedAuthorizationResolver for ObservingResolver { + async fn resolve( + &self, + _request: &ProtectedOperationRequest, + ) -> Result { + panic!("observational mode must not finalize authority") + } + + async fn observe( + &self, + _request: &ProtectedOperationRequest, + ) -> Result<(), ProtectedResolutionError> { + self.0.fetch_add(1, Ordering::SeqCst); + Err(ProtectedResolutionError::new("synthetic_denial")) + } + } + + struct FixedClock; + + impl AuthorizationClock for FixedClock { + fn now(&self) -> Result { + Ok(AuthorizationTime::from_unix_seconds(100)) + } + } + + fn domain(value: u128) -> CommunityId { + CommunityId::from_uuid(Uuid::from_u128(value)) + } + + fn runtime( + policies: impl IntoIterator, + ) -> Result { + ProtectedTransportRuntime::new( + policies, + Arc::new(UnavailableResolver), + Arc::new(FixedClock), + ) + } + + fn request(domain: CommunityId) -> ProtectedOperationRequest { + let keys = Keys::generate(); + let challenge = "protected-request-test"; + let relay_url = "wss://relay.example"; + let event = EventBuilder::auth(challenge, RelayUrl::parse(relay_url).expect("relay URL")) + .sign_with_keys(&keys) + .expect("signed NIP-42 event"); + let proof = VerifiedEvidenceAdapter::new() + .verify_nip42( + domain, + AuthTransport::RelayWebSocket, + &event, + challenge, + relay_url, + None, + ) + .expect("verified NIP-42 proof"); + ProtectedOperationRequest::new( + Arc::new(proof), + None, + AuthorizationCapability::CommunityRead, + Uuid::new_v4(), + "ws_req", + ) + .expect("valid request") + } + + #[test] + fn request_retains_sealed_proof_and_rejects_surface_substitution() { + let request = request(domain(1)); + let retained = Arc::clone(request.verified_proof()); + assert!(Arc::ptr_eq(&retained, request.verified_proof())); + assert_eq!(request.authorization_domain(), domain(1)); + assert_eq!(request.transport(), AuthTransport::RelayWebSocket); + + assert!(matches!( + ProtectedOperationRequest::new( + retained, + None, + AuthorizationCapability::MediaWrite, + Uuid::new_v4(), + "media.upload", + ), + Err(ProtectedTransportError::SurfaceCapabilityMismatch) + )); + } + + #[test] + fn finalized_context_cannot_be_replayed_across_correlations() { + let finalized = Uuid::from_u128(1); + assert!(exact_correlation_matches(finalized, finalized)); + assert!(!exact_correlation_matches(finalized, Uuid::from_u128(2))); + } + + #[tokio::test] + async fn absent_and_off_domains_preserve_legacy_behavior() { + let configured = domain(1); + let absent = domain(2); + let runtime = runtime([DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Off, + )]) + .expect("runtime"); + assert!(matches!( + runtime.authorize(&request(configured)).await, + Ok(ProtectedAuthorization::Legacy) + )); + assert!(matches!( + runtime.authorize(&request(absent)).await, + Ok(ProtectedAuthorization::Legacy) + )); + } + + #[tokio::test] + async fn enforcing_domain_never_falls_back_on_resolver_failure() { + let configured = domain(1); + let runtime = runtime([DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Enforce, + )]) + .expect("runtime"); + assert!(matches!( + runtime.authorize(&request(configured)).await, + Err(ProtectedTransportError::Resolution(_)) + )); + } + + #[tokio::test] + async fn verify_only_preserves_legacy_access_without_finalizing() { + let configured = domain(1); + let runtime = ProtectedTransportRuntime::new( + [DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::VerifyOnly, + )], + Arc::new(PanicResolver), + Arc::new(FixedClock), + ) + .expect("runtime"); + assert!(matches!( + runtime.authorize(&request(configured)).await, + Ok(ProtectedAuthorization::Legacy) + )); + } + + #[tokio::test] + async fn shadow_preserves_legacy_access_without_finalizing() { + let configured = domain(1); + let runtime = ProtectedTransportRuntime::new( + [DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Shadow, + )], + Arc::new(PanicResolver), + Arc::new(FixedClock), + ) + .expect("runtime"); + assert!(matches!( + runtime.authorize(&request(configured)).await, + Ok(ProtectedAuthorization::Legacy) + )); + } + + #[tokio::test] + async fn observational_denial_is_read_only_and_never_changes_access() { + let configured = domain(1); + let resolver = Arc::new(ObservingResolver(AtomicUsize::new(0))); + let runtime = ProtectedTransportRuntime::new( + [DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Shadow, + )], + resolver.clone(), + Arc::new(FixedClock), + ) + .expect("runtime"); + assert!(matches!( + runtime.authorize(&request(configured)).await, + Ok(ProtectedAuthorization::Legacy) + )); + assert_eq!(resolver.0.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn client_status_is_absent_in_off_and_shadow_and_fail_closed_when_unavailable() { + let configured = domain(1); + for mode in [AuthorizationMode::Off, AuthorizationMode::Shadow] { + let runtime = runtime([DomainTransportPolicy::from_server_configuration( + configured, mode, + )]) + .expect("runtime"); + assert!(runtime + .present_status(&request(configured)) + .await + .expect("non-presenting mode") + .is_none()); + } + + for mode in [AuthorizationMode::VerifyOnly, AuthorizationMode::Enforce] { + let runtime = runtime([DomainTransportPolicy::from_server_configuration( + configured, mode, + )]) + .expect("runtime"); + assert!(matches!( + runtime.present_status(&request(configured)).await, + Err(ProtectedTransportError::Resolution(_)) + )); + } + } + + #[test] + fn duplicate_exact_domain_policy_is_rejected() { + let configured = domain(1); + assert!(matches!( + runtime([ + DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Off, + ), + DomainTransportPolicy::from_server_configuration( + configured, + AuthorizationMode::Enforce, + ), + ]), + Err(ProtectedTransportError::AmbiguousDomainPolicy) + )); + } + + #[test] + fn atomic_mutations_are_unavailable_only_for_enforcing_authority() { + assert!(require_atomic_mutation_executor(false).is_ok()); + assert!(matches!( + require_atomic_mutation_executor(true), + Err(ProtectedTransportError::AtomicMutationFenceUnavailable) + )); + assert!(ProtectedAuthorization::Legacy + .require_atomic_mutation() + .is_ok()); + } + + #[test] + fn unwired_mutations_preserve_legacy_modes_and_deny_enforce() { + for mode in [ + None, + Some(AuthorizationMode::Off), + Some(AuthorizationMode::Shadow), + Some(AuthorizationMode::VerifyOnly), + ] { + assert!(matches!( + authorize_unwired_for_mode(mode), + Ok(ProtectedAuthorization::Legacy) + )); + } + assert!(matches!( + authorize_unwired_for_mode(Some(AuthorizationMode::Enforce)), + Err(ProtectedTransportError::MissingVerifiedProof) + )); + } +}