From df4efac2f54eb12af1546df2d23ea98d5e507347 Mon Sep 17 00:00:00 2001 From: Cea Stapleton Cordasco <261786559+cea-block@users.noreply.github.com> Date: Tue, 4 Aug 2026 11:43:41 -0500 Subject: [PATCH] fix(auth): enforce server-owned enrollment policy Canonical O4 slice: 08_enrollment_policy Final-tree source: e45553ea65b66631771c16908513f983d3c0b881 Signed-off-by: Cea Stapleton Cordasco <261786559+cea-block@users.noreply.github.com> --- crates/buzz-auth/src/context/binding.rs | 39 +- crates/buzz-auth/src/context/reason.rs | 12 + crates/buzz-auth/src/finalization.rs | 1016 +++++++++++++ crates/buzz-db/src/archived_identities.rs | 47 +- crates/buzz-db/src/channel.rs | 1332 +++++++++++++++-- crates/buzz-db/src/relay_invite.rs | 123 +- crates/buzz-db/src/relay_members.rs | 132 +- crates/buzz-db/src/user.rs | 107 ++ crates/buzz-relay/src/api/invites.rs | 488 +++++- .../src/handlers/identity_archive.rs | 159 ++ 10 files changed, 3208 insertions(+), 247 deletions(-) create mode 100644 crates/buzz-auth/src/finalization.rs diff --git a/crates/buzz-auth/src/context/binding.rs b/crates/buzz-auth/src/context/binding.rs index e9ca0e9c9..9c0727dae 100644 --- a/crates/buzz-auth/src/context/binding.rs +++ b/crates/buzz-auth/src/context/binding.rs @@ -185,7 +185,6 @@ impl ResolvedFederatedPolicy { pub(crate) const fn from_authoritative_resolution(stamp: FederatedPolicyStamp) -> Self { Self { stamp } } - #[cfg(test)] pub(crate) fn not_required(authorization_domain: CommunityId) -> Self { Self::from_authoritative_resolution( @@ -601,6 +600,42 @@ impl VersionedBindingRef { } } + pub(crate) fn from_evidence_adapter( + authorization_domain: CommunityId, + binding_id: Uuid, + principal: FederatedPrincipal, + bound_pubkey: PublicKey, + binding_version: BindingVersion, + expires_at: Option, + source: BindingSource, + resolution_reason: AuthorizationReason, + ) -> Result { + if binding_id.is_nil() { + return Err(AuthContextError::InvalidBindingId); + } + let valid_reason = matches!( + (source, resolution_reason), + (_, AuthorizationReason::ExistingBinding) + | ( + BindingSource::AttestedKey, + AuthorizationReason::EnrolledAttestedKey + ) + | (BindingSource::Tofu, AuthorizationReason::EnrolledTofu) + ); + if !valid_reason { + return Err(AuthContextError::InvalidAuthorizationReason); + } + Ok(Self { + authorization_domain, + binding_id, + principal, + bound_pubkey, + binding_version, + expires_at, + source, + resolution_reason, + }) + } /// Build a reference to a binding authoritatively resolved as already active. #[cfg(test)] pub(crate) fn new_existing_active_for_test( @@ -697,7 +732,7 @@ impl VersionedBindingRef { } /// Stable reason proven by the authoritative binding lifecycle result. - pub(super) const fn authorization_reason(&self) -> AuthorizationReason { + pub(crate) const fn authorization_reason(&self) -> AuthorizationReason { self.resolution_reason } } diff --git a/crates/buzz-auth/src/context/reason.rs b/crates/buzz-auth/src/context/reason.rs index fe07e9a68..eba2f4356 100644 --- a/crates/buzz-auth/src/context/reason.rs +++ b/crates/buzz-auth/src/context/reason.rs @@ -154,6 +154,9 @@ pub enum AuthContextError { /// Proof method was not valid for the transport being authorized. #[error("Nostr proof method does not match authorization transport")] TransportProofMismatch, + /// Exact verifier-operation evidence was not valid for the transport. + #[error("Nostr operation proof does not match authorization transport")] + OperationProofMismatch, /// Direct federated authorization was attached to a delegated Nostr actor. #[error("direct federated authorization cannot include a delegated Nostr owner")] DirectAuthorizationHasOwner, @@ -172,6 +175,12 @@ pub enum AuthContextError { /// Delegated owner evidence did not resolve an already-active binding. #[error("delegated federated authorization requires an existing active binding")] DelegatedBindingNotExistingActive, + /// Federated finalization did not consume a current provider decision. + #[error("federated authorization requires a current provider decision")] + ProviderDecisionRequired, + /// An issued lease did not match the context evidence being finalized. + #[error("authorization lease does not match finalized context evidence")] + FinalizedLeaseMismatch, } impl AuthContextError { @@ -214,6 +223,7 @@ impl AuthContextError { Self::OwnerAdmissionPrincipalMismatch => "owner_admission_principal_mismatch", Self::AssertionPrincipalMismatch => "federated_assertion_principal_mismatch", Self::TransportProofMismatch => "nostr_transport_proof_mismatch", + Self::OperationProofMismatch => "nostr_operation_proof_mismatch", Self::DirectAuthorizationHasOwner => "federated_direct_has_owner", Self::DirectBindingKeyMismatch => "federated_direct_key_mismatch", Self::DelegateKeyMismatch => "federated_delegate_key_mismatch", @@ -222,6 +232,8 @@ impl AuthContextError { Self::DelegatedBindingNotExistingActive => { "federated_delegated_binding_not_existing_active" } + Self::ProviderDecisionRequired => "authorization_provider_decision_required", + Self::FinalizedLeaseMismatch => "authorization_finalized_lease_mismatch", } } } diff --git a/crates/buzz-auth/src/finalization.rs b/crates/buzz-auth/src/finalization.rs new file mode 100644 index 000000000..0e7ef2559 --- /dev/null +++ b/crates/buzz-auth/src/finalization.rs @@ -0,0 +1,1016 @@ +//! Federated authorization finalization. +//! +//! This is the only crate-owned path from a validated provider capability +//! snapshot to an access lease. Display-only verification is returned as a +//! separate type that cannot be converted into an [`crate::AuthContext`] or +//! consumed by protected-operation lease validation. + +use std::fmt; + +use buzz_core::CommunityId; +use nostr::PublicKey; +use thiserror::Error; +use uuid::Uuid; + +use crate::{ + context::{ + validate_context_evidence, AuthContext, AuthContextError, AuthContextInput, BindingVersion, + FederatedAuthorization, FederatedIdentityRequirement, ResolvedFederatedPolicy, + VersionedBindingRef, + }, + lease::{ + conservative_expiry, AccessLeasePolicy, AuthorizationClockError, AuthorizationLease, + AuthorizationTime, BindingLeaseBound, LeaseIssueError, LeaseVersion, + SharedAuthorizationClock, VerificationStatusPolicy, + }, + provider::{AuthorizationProfileId, CapabilitySnapshot, DecisionSource, PolicyVersion}, +}; + +/// Finalizer using one centrally injected clock for all authorization time. +#[derive(Clone)] +pub struct AuthorizationFinalizer { + clock: SharedAuthorizationClock, +} + +impl AuthorizationFinalizer { + /// Create a finalizer backed by the supplied central authorization clock. + pub fn new(clock: SharedAuthorizationClock) -> Self { + Self { clock } + } + + /// Read current time from the injected authorization clock. + pub fn now(&self) -> Result { + self.clock.now() + } + + /// Finalize enforcing federated authority and issue one bounded lease. + /// + /// The provider snapshot must be the current allow decision for the exact + /// domain, actor, transport, principal, profile, and correlation ID. The + /// resulting lease expires at the earliest provider/identity bound, + /// binding freshness bound, or configured application maximum, shortened + /// by the explicit conservative clock skew. + #[allow(clippy::too_many_arguments)] + pub fn finalize_access( + &self, + input: AuthContextInput, + federated_policy: ResolvedFederatedPolicy, + authorization: FederatedAuthorization, + snapshot: Box, + expected_profile: &AuthorizationProfileId, + binding_bound: BindingLeaseBound, + lease_policy: AccessLeasePolicy, + lease_version: LeaseVersion, + ) -> Result { + let now = self.now()?; + validate_context_evidence( + &input, + &federated_policy, + &authorization, + now.unix_seconds(), + )?; + let active_binding = validate_provider_evidence( + &input, + &federated_policy, + &authorization, + &snapshot, + expected_profile, + &binding_bound, + now, + )?; + let lease = AuthorizationLease::issue( + lease_version, + snapshot.authorization_domain(), + snapshot.transport(), + snapshot.actor_pubkey(), + snapshot.owner_pubkey(), + active_binding, + binding_bound, + snapshot.profile_id().clone(), + snapshot.policy_version().clone(), + snapshot.capabilities().clone(), + snapshot.effective_until(), + now, + lease_policy, + snapshot.correlation_id(), + )?; + Ok(AuthContext::finalize_v1_with_lease( + input, + federated_policy, + authorization, + lease, + now.unix_seconds(), + )?) + } + + /// Finalize display-only verification without issuing access authority. + /// + /// This path requires a direct active binding for the event-author key. + /// The returned type carries no capability set, access lease, membership, + /// or conversion into an authorized context. + #[allow(clippy::too_many_arguments)] + pub fn finalize_verification_only( + &self, + input: AuthContextInput, + federated_policy: ResolvedFederatedPolicy, + authorization: FederatedAuthorization, + snapshot: Box, + expected_profile: &AuthorizationProfileId, + binding_bound: BindingLeaseBound, + status_policy: VerificationStatusPolicy, + ) -> Result { + let now = self.now()?; + validate_context_evidence( + &input, + &federated_policy, + &authorization, + now.unix_seconds(), + )?; + let active_binding = validate_provider_evidence( + &input, + &federated_policy, + &authorization, + &snapshot, + expected_profile, + &binding_bound, + now, + )?; + if !matches!(authorization, FederatedAuthorization::Direct { .. }) { + return Err(FinalizationError::VerificationRequiresDirectBinding); + } + let expires_at = conservative_expiry( + now, + snapshot.effective_until(), + binding_bound.valid_until(), + status_policy.application_limit(), + status_policy.clock_skew(), + )?; + Ok(VerificationOnlyDisposition { + authorization_domain: snapshot.authorization_domain(), + actor_pubkey: snapshot.actor_pubkey(), + binding_id: active_binding.binding_id(), + binding_version: active_binding.binding_version(), + profile_id: snapshot.profile_id().clone(), + policy_version: snapshot.policy_version().clone(), + correlation_id: snapshot.correlation_id(), + issued_at: now.unix_seconds(), + expires_at, + }) + } +} + +impl fmt::Debug for AuthorizationFinalizer { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("AuthorizationFinalizer") + .field("clock", &"[injected]") + .finish() + } +} + +/// Short-lived, display-only proof of a current direct binding. +/// +/// This type is deliberately not an authorization context or lease. Protected +/// operations accept [`AuthorizationLease`], so a +/// verification-only result cannot grant access even if a caller retains it. +#[must_use] +#[derive(PartialEq, Eq)] +pub struct VerificationOnlyDisposition { + authorization_domain: CommunityId, + actor_pubkey: PublicKey, + binding_id: Uuid, + binding_version: BindingVersion, + profile_id: AuthorizationProfileId, + policy_version: PolicyVersion, + correlation_id: Uuid, + issued_at: u64, + expires_at: u64, +} + +impl VerificationOnlyDisposition { + /// Exact authorization domain represented by the display status. + pub const fn authorization_domain(&self) -> CommunityId { + self.authorization_domain + } + + /// Event-author key whose direct active binding was verified. + pub const fn actor_pubkey(&self) -> PublicKey { + self.actor_pubkey + } + + /// Stable direct binding identifier. + pub const fn binding_id(&self) -> Uuid { + self.binding_id + } + + /// Exact active binding version. + pub const fn binding_version(&self) -> BindingVersion { + self.binding_version + } + + /// Server-resolved provider profile. + pub const fn profile_id(&self) -> &AuthorizationProfileId { + &self.profile_id + } + + /// Current opaque provider policy version. + pub const fn policy_version(&self) -> &PolicyVersion { + &self.policy_version + } + + /// Correlation identifier for the display-only decision. + pub const fn correlation_id(&self) -> Uuid { + self.correlation_id + } + + /// Central issue time in Unix seconds. + pub const fn issued_at(&self) -> u64 { + self.issued_at + } + + /// Conservative display-status expiry in Unix seconds. + pub const fn expires_at(&self) -> u64 { + self.expires_at + } + + /// Check display freshness using the supplied central authorization clock. + /// + /// This check only controls presentation and is not an access decision. + pub fn is_current( + &self, + clock: &dyn crate::AuthorizationClock, + ) -> Result { + let now = clock.now()?; + Ok(now.unix_seconds() >= self.issued_at && now.unix_seconds() < self.expires_at) + } +} + +impl fmt::Debug for VerificationOnlyDisposition { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("VerificationOnlyDisposition") + .field("authorization_domain", &"[redacted]") + .field("actor_pubkey", &"[redacted]") + .field("binding_id", &"[redacted]") + .field("binding_version", &"[redacted]") + .field("profile_id", &"[redacted]") + .field("policy_version", &"[redacted]") + .field("correlation_id", &"[redacted]") + .field("issued_at", &"[redacted]") + .field("expires_at", &"[redacted]") + .finish() + } +} + +fn validate_provider_evidence<'a>( + input: &AuthContextInput, + federated_policy: &ResolvedFederatedPolicy, + authorization: &'a FederatedAuthorization, + snapshot: &CapabilitySnapshot, + expected_profile: &AuthorizationProfileId, + binding_bound: &BindingLeaseBound, + now: AuthorizationTime, +) -> Result<&'a VersionedBindingRef, FinalizationError> { + if !matches!( + federated_policy.requirement(), + FederatedIdentityRequirement::Required(_) + ) { + return Err(FinalizationError::FederatedPolicyRequired); + } + let domain = input.tenant().community(); + let proof = input.nostr_proof(); + if federated_policy.authorization_domain() != domain + || !snapshot.is_bound_to_federated_policy(federated_policy) + || snapshot.authorization_domain() != domain + || snapshot.transport() != proof.authorized_transport() + || snapshot.actor_pubkey() != proof.actor_pubkey() + || snapshot.proof_method() != proof.proof_method() + || snapshot.correlation_id() != input.correlation_id() + || snapshot.profile_id() != expected_profile + { + return Err(FinalizationError::ProviderEvidenceMismatch); + } + if snapshot.issued_at() > now.unix_seconds() + || snapshot.fresh_until() <= now.unix_seconds() + || snapshot.effective_until() <= now.unix_seconds() + { + return Err(FinalizationError::ProviderEvidenceStale); + } + let active_binding = authorization + .active_binding() + .ok_or(FinalizationError::FederatedAuthorizationRequired)?; + if active_binding.binding_id() != binding_bound.binding_id() + || active_binding.binding_version() != binding_bound.binding_version() + { + return Err(FinalizationError::BindingEvidenceMismatch); + } + match authorization { + FederatedAuthorization::NotRequired => { + return Err(FinalizationError::FederatedAuthorizationRequired); + } + FederatedAuthorization::Direct { binding, assertion } => { + if snapshot.decision_source() != DecisionSource::DirectAssertion + || snapshot.owner_pubkey().is_some() + || snapshot.binding_id().is_some() + || snapshot.binding_version().is_some() + || snapshot.principal() != binding.principal() + || snapshot.principal() != assertion.principal() + || snapshot.effective_until() > assertion.expires_at().unix_seconds() + { + return Err(FinalizationError::ProviderEvidenceMismatch); + } + } + FederatedAuthorization::Delegated { owner, admission } => { + if !proof + .verified_delegation() + .is_some_and(|delegation| delegation.capability().is_transport_wide()) + { + return Err(FinalizationError::UnsupportedDelegationScope); + } + if snapshot.decision_source() != DecisionSource::DelegatedOwnerBinding + || snapshot.owner_pubkey() != Some(owner.bound_pubkey()) + || snapshot.binding_id() != Some(owner.binding_id()) + || snapshot.binding_version() != Some(owner.binding_version()) + || snapshot.principal() != owner.principal() + || snapshot.principal() != admission.principal() + || snapshot.fresh_until() != admission.fresh_until().unix_seconds() + { + return Err(FinalizationError::ProviderEvidenceMismatch); + } + if let Some(delegation_expiry) = proof + .verified_delegation() + .and_then(|delegation| delegation.expires_at()) + { + if snapshot.effective_until() > delegation_expiry.unix_seconds() { + return Err(FinalizationError::ProviderEvidenceMismatch); + } + } + } + } + Ok(active_binding) +} + +/// Fail-closed federated finalization error. +#[derive(Debug, Error)] +pub enum FinalizationError { + /// The central authorization clock failed. + #[error(transparent)] + Clock(#[from] AuthorizationClockError), + /// Existing context evidence was invalid. + #[error(transparent)] + Context(#[from] AuthContextError), + /// A bounded lease or display expiry could not be issued. + #[error(transparent)] + Lease(#[from] LeaseIssueError), + /// The server-resolved policy did not require federated authorization. + #[error("federated finalization requires server-resolved federated policy")] + FederatedPolicyRequired, + /// No direct or delegated active binding evidence was supplied. + #[error("federated finalization requires active binding authorization")] + FederatedAuthorizationRequired, + /// The provider snapshot did not match the exact finalization evidence. + #[error("provider decision does not match finalization evidence")] + ProviderEvidenceMismatch, + /// The provider snapshot was future-issued, stale, or expired. + #[error("provider decision is no longer current")] + ProviderEvidenceStale, + /// Binding freshness evidence named another binding or version. + #[error("binding freshness evidence does not match active binding")] + BindingEvidenceMismatch, + /// Display verification was attempted for delegated rather than direct authority. + #[error("verification-only display requires the event author's direct active binding")] + VerificationRequiresDirectBinding, + /// Narrower delegation evidence reached the transport-wide finalizer. + #[error("operation-bound delegation cannot be promoted to transport-wide authority")] + UnsupportedDelegationScope, +} + +impl FinalizationError { + /// Stable audit and metric code. + pub fn code(&self) -> &'static str { + match self { + Self::Clock(_) => "authorization_finalize_001", + Self::Context(error) => error.code(), + Self::Lease(error) => error.code(), + Self::FederatedPolicyRequired => "authorization_finalize_002", + Self::FederatedAuthorizationRequired => "authorization_finalize_003", + Self::ProviderEvidenceMismatch => "authorization_finalize_004", + Self::ProviderEvidenceStale => "authorization_finalize_005", + Self::BindingEvidenceMismatch => "authorization_finalize_006", + Self::VerificationRequiresDirectBinding => "authorization_finalize_007", + Self::UnsupportedDelegationScope => "authorization_finalize_008", + } + } +} + +#[cfg(test)] +mod tests { + use std::sync::{ + atomic::{AtomicU64, Ordering}, + Arc, + }; + use std::time::Duration; + + use nostr::Keys; + + use super::*; + use crate::{ + context::{ + AssertionExpiry, AssertionTransport, AuthContextVersion, AuthMethod, AuthTransport, + AuthorizedCommunityAccess, BindingSource, DelegationExpiry, EnrollmentMode, + VerifiedFederatedAssertion, VerifiedKeyAttestation, VerifiedNostrProof, + VerifiedTransportDelegation, + }, + lease::{ + ApplicationLeaseLimit, AuthorizationClock, AuthorizationClockSkew, + AuthorizationLeaseValidator, LeaseRenewalAction, LeaseRenewalLeadTime, + LeaseUseRequirement, + }, + provider::{ + resolve_authorization, AuthorizationCapability, AuthorizationOutcome, + AuthorizationProvider, AuthorizationProviderFuture, AuthorizationRequest, + CapabilitySet, ProviderAllow, ProviderDecision, ProviderTimeout, + }, + Scope, + }; + + struct FixedClock(AtomicU64); + + impl FixedClock { + fn new(now: u64) -> Self { + Self(AtomicU64::new(now)) + } + + fn set(&self, now: u64) { + self.0.store(now, Ordering::SeqCst); + } + } + + impl AuthorizationClock for FixedClock { + fn now(&self) -> Result { + Ok(AuthorizationTime::from_unix_seconds( + self.0.load(Ordering::SeqCst), + )) + } + } + + struct AllowProvider { + issued_at: u64, + fresh_until: u64, + } + + impl AuthorizationProvider for AllowProvider { + fn authorize<'a>( + &'a self, + request: &'a AuthorizationRequest, + ) -> AuthorizationProviderFuture<'a> { + let allow = ProviderAllow::new( + request.authorization_domain(), + request.principal().clone(), + request.profile_id().clone(), + request.requested_capabilities().clone(), + PolicyVersion::new("policy.synthetic.example") + .expect("synthetic policy version is valid"), + self.issued_at, + self.fresh_until, + ) + .expect("synthetic provider result is valid"); + Box::pin(std::future::ready(ProviderDecision::Allow(allow))) + } + } + + fn domain() -> CommunityId { + CommunityId::from_uuid(Uuid::from_u128(0x100)) + } + + fn profile() -> AuthorizationProfileId { + AuthorizationProfileId::from_server_configuration("profile.synthetic.example") + .expect("synthetic profile is valid") + } + + struct DirectFixture { + input: AuthContextInput, + policy: ResolvedFederatedPolicy, + authorization: FederatedAuthorization, + binding_bound: BindingLeaseBound, + binding_id: Uuid, + binding_version: BindingVersion, + actor_pubkey: PublicKey, + } + + fn direct_fixture(assertion_expiry: u64, binding_expiry: u64) -> DirectFixture { + let actor = Keys::generate().public_key(); + let principal = + crate::FederatedPrincipal::new("https://issuer.synthetic.example", "subject-synthetic") + .expect("synthetic principal is valid"); + let proof = VerifiedNostrProof::new( + domain(), + AuthTransport::RelayWebSocket, + actor, + AuthMethod::Nip42, + None, + ) + .expect("synthetic proof is valid"); + let assertion = VerifiedFederatedAssertion::new( + domain(), + AuthTransport::RelayWebSocket, + principal.clone(), + Some(VerifiedKeyAttestation::new(actor)), + AssertionTransport::TrustedProxy, + None, + AssertionExpiry::new(assertion_expiry).expect("synthetic expiry is valid"), + ); + let binding_id = Uuid::from_u128(0x200); + let binding_version = BindingVersion::new(7).expect("synthetic version is valid"); + let binding = VersionedBindingRef::new_existing_active_for_test( + domain(), + binding_id, + principal, + actor, + binding_version, + BindingSource::AttestedKey, + ) + .expect("synthetic binding is valid"); + let binding_bound = BindingLeaseBound::new(&binding, binding_expiry) + .expect("synthetic binding bound is valid"); + let tenant = + buzz_core::tenant::TenantContext::resolved(domain(), "relay.synthetic.example"); + DirectFixture { + input: AuthContextInput::new( + tenant, + Uuid::from_u128(0x300), + proof, + AuthorizedCommunityAccess::new(domain(), vec![Scope::MessagesRead], None), + ), + policy: ResolvedFederatedPolicy::server_resolved_required( + domain(), + EnrollmentMode::AttestedKey, + ), + authorization: FederatedAuthorization::Direct { binding, assertion }, + binding_bound, + binding_id, + binding_version, + actor_pubkey: actor, + } + } + + async fn direct_snapshot( + fixture: &DirectFixture, + profile: AuthorizationProfileId, + provider_fresh_until: u64, + ) -> Box { + let FederatedAuthorization::Direct { assertion, .. } = &fixture.authorization else { + panic!("direct fixture must contain direct authorization"); + }; + let request = AuthorizationRequest::direct( + fixture.input.nostr_proof(), + assertion, + profile, + CapabilitySet::single(AuthorizationCapability::CommunityRead), + fixture.input.correlation_id(), + 1_000, + ) + .expect("synthetic request is valid"); + let outcome = resolve_authorization( + &AllowProvider { + issued_at: 999, + fresh_until: provider_fresh_until, + }, + &request, + 1_000, + ProviderTimeout::new(Duration::from_secs(1)).expect("synthetic timeout is valid"), + ) + .await; + match outcome { + AuthorizationOutcome::Allow(snapshot) => snapshot, + other => panic!("synthetic provider must allow, got {other:?}"), + } + } + + fn access_policy(limit_seconds: u64) -> AccessLeasePolicy { + AccessLeasePolicy::new( + ApplicationLeaseLimit::from_seconds(limit_seconds) + .expect("synthetic application limit is valid"), + AuthorizationClockSkew::from_seconds(5).expect("synthetic skew is valid"), + ) + } + + #[tokio::test] + async fn direct_lease_carries_binding_and_earliest_application_expiry() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context + .authorization_lease() + .expect("enforcing context carries a lease"); + assert_eq!(lease.binding_id(), fixture.binding_id); + assert_eq!(lease.binding_version(), fixture.binding_version); + assert_eq!(lease.expires_at(), 1_095); + assert_eq!( + lease.capabilities().as_slice(), + &[AuthorizationCapability::CommunityRead] + ); + } + + #[test] + fn earliest_expiry_includes_provider_binding_application_and_skew() { + let now = AuthorizationTime::from_unix_seconds(1_000); + let limit = ApplicationLeaseLimit::from_seconds(500).expect("synthetic limit is valid"); + let skew = AuthorizationClockSkew::from_seconds(5).expect("synthetic skew is valid"); + assert_eq!( + conservative_expiry(now, 1_100, 1_200, limit, skew), + Ok(1_095) + ); + assert_eq!( + conservative_expiry(now, 1_300, 1_080, limit, skew), + Ok(1_075) + ); + let short_limit = + ApplicationLeaseLimit::from_seconds(60).expect("synthetic limit is valid"); + assert_eq!( + conservative_expiry(now, 1_300, 1_200, short_limit, skew), + Ok(1_055) + ); + } + + #[tokio::test] + async fn operation_guard_is_per_capability_and_revalidates_at_commit() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let actor = fixture.actor_pubkey; + let binding_id = fixture.binding_id; + let binding_version = fixture.binding_version; + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + let validator = AuthorizationLeaseValidator::new(clock.clone()); + let requirement = LeaseUseRequirement { + context_version: AuthContextVersion::V1, + lease_version: LeaseVersion::INITIAL, + authorization_domain: domain(), + transport: AuthTransport::RelayWebSocket, + actor_pubkey: actor, + binding_id, + binding_version, + profile_id: lease.profile_id().clone(), + policy_version: lease.policy_version().clone(), + capability: AuthorizationCapability::CommunityRead, + }; + let guard = context + .operation_guard(&validator, requirement) + .expect("exact capability creates an operation guard"); + guard + .revalidate() + .expect("guard remains valid before commit boundary"); + clock.set(lease.expires_at()); + assert!(matches!( + guard.revalidate(), + Err(crate::LeaseValidationError::Expired) + )); + } + + #[tokio::test] + async fn same_version_different_binding_is_denied() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let actor = fixture.actor_pubkey; + let binding_version = fixture.binding_version; + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + let requirement = LeaseUseRequirement { + context_version: AuthContextVersion::V1, + lease_version: LeaseVersion::INITIAL, + authorization_domain: domain(), + transport: AuthTransport::RelayWebSocket, + actor_pubkey: actor, + binding_id: Uuid::from_u128(0x201), + binding_version, + profile_id: lease.profile_id().clone(), + policy_version: lease.policy_version().clone(), + capability: AuthorizationCapability::CommunityRead, + }; + + assert_eq!( + AuthorizationLeaseValidator::new(clock).authorize(lease, &requirement), + Err(crate::LeaseValidationError::BindingIdMismatch) + ); + } + + #[tokio::test] + async fn same_policy_version_different_profile_is_denied() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let actor = fixture.actor_pubkey; + let binding_id = fixture.binding_id; + let binding_version = fixture.binding_version; + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + let requirement = LeaseUseRequirement { + context_version: AuthContextVersion::V1, + lease_version: LeaseVersion::INITIAL, + authorization_domain: domain(), + transport: AuthTransport::RelayWebSocket, + actor_pubkey: actor, + binding_id, + binding_version, + profile_id: AuthorizationProfileId::from_server_configuration( + "other-profile.synthetic.example", + ) + .expect("synthetic profile is valid"), + policy_version: lease.policy_version().clone(), + capability: AuthorizationCapability::CommunityRead, + }; + + assert_eq!( + AuthorizationLeaseValidator::new(clock).authorize(lease, &requirement), + Err(crate::LeaseValidationError::AuthorizationProfileMismatch) + ); + } + + #[tokio::test] + async fn typed_version_policy_and_capability_mismatches_fail_closed() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let actor = fixture.actor_pubkey; + let binding_id = fixture.binding_id; + let binding_version = fixture.binding_version; + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + let validator = AuthorizationLeaseValidator::new(clock); + let mut requirement = LeaseUseRequirement { + context_version: AuthContextVersion::V1, + lease_version: LeaseVersion::new(2).expect("synthetic version is valid"), + authorization_domain: domain(), + transport: AuthTransport::RelayWebSocket, + actor_pubkey: actor, + binding_id, + binding_version, + profile_id: lease.profile_id().clone(), + policy_version: lease.policy_version().clone(), + capability: AuthorizationCapability::CommunityRead, + }; + assert!(matches!( + validator.authorize(lease, &requirement), + Err(crate::LeaseValidationError::LeaseVersionMismatch) + )); + requirement.lease_version = LeaseVersion::INITIAL; + requirement.policy_version = + PolicyVersion::new("changed.synthetic.example").expect("synthetic version is valid"); + assert!(matches!( + validator.authorize(lease, &requirement), + Err(crate::LeaseValidationError::PolicyVersionMismatch) + )); + requirement.policy_version = lease.policy_version().clone(); + requirement.capability = AuthorizationCapability::CommunityWrite; + assert!(matches!( + validator.authorize(lease, &requirement), + Err(crate::LeaseValidationError::MissingCapability) + )); + } + + #[tokio::test] + async fn verification_only_is_short_lived_and_has_no_access_context() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let status = finalizer + .finalize_verification_only( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + VerificationStatusPolicy::new( + ApplicationLeaseLimit::from_seconds(30) + .expect("synthetic status limit is valid"), + AuthorizationClockSkew::from_seconds(5).expect("synthetic skew is valid"), + ), + ) + .expect("complete current direct evidence produces display status"); + assert_eq!(status.expires_at(), 1_025); + assert!(status + .is_current(clock.as_ref()) + .expect("clock is available")); + } + + #[tokio::test] + async fn renewal_hooks_expose_renew_and_hard_expiry_boundaries() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock.clone()); + let expected_profile = profile(); + let fixture = direct_fixture(1_500, 1_400); + let snapshot = direct_snapshot(&fixture, expected_profile.clone(), 1_300).await; + let context = finalizer + .finalize_access( + fixture.input, + fixture.policy, + fixture.authorization, + snapshot, + &expected_profile, + fixture.binding_bound, + access_policy(100), + LeaseVersion::INITIAL, + ) + .expect("complete current evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + let schedule = lease.renewal_schedule( + LeaseRenewalLeadTime::from_seconds(20).expect("synthetic lead is valid"), + ); + let validator = AuthorizationLeaseValidator::new(clock.clone()); + assert_eq!(schedule.renew_at(), 1_075); + assert_eq!(schedule.expires_at(), 1_095); + assert_eq!( + validator + .renewal_action(schedule) + .expect("clock is available"), + LeaseRenewalAction::Current + ); + clock.set(schedule.renew_at()); + assert_eq!( + validator + .renewal_action(schedule) + .expect("clock is available"), + LeaseRenewalAction::RenewNow + ); + clock.set(schedule.expires_at()); + assert_eq!( + validator + .renewal_action(schedule) + .expect("clock is available"), + LeaseRenewalAction::Expired + ); + } + + #[tokio::test] + async fn delegated_lease_is_bounded_by_delegation_and_owner_binding() { + let clock = Arc::new(FixedClock::new(1_000)); + let finalizer = AuthorizationFinalizer::new(clock); + let expected_profile = profile(); + let owner = Keys::generate().public_key(); + let delegate = Keys::generate().public_key(); + let principal = crate::FederatedPrincipal::new( + "https://issuer.synthetic.example", + "owner-subject-synthetic", + ) + .expect("synthetic principal is valid"); + let proof = VerifiedNostrProof::new( + domain(), + AuthTransport::RelayWebSocket, + delegate, + AuthMethod::Nip42, + Some( + VerifiedTransportDelegation::new_unrestricted( + owner, + delegate, + Some( + DelegationExpiry::new(1_050).expect("synthetic delegation expiry is valid"), + ), + ) + .expect("synthetic delegation is valid"), + ), + ) + .expect("synthetic proof is valid"); + let binding_id = Uuid::from_u128(0x400); + let binding_version = BindingVersion::new(9).expect("synthetic version is valid"); + let owner_binding = VersionedBindingRef::new_existing_active_for_test( + domain(), + binding_id, + principal.clone(), + owner, + binding_version, + BindingSource::Provisioned, + ) + .expect("synthetic owner binding is valid"); + let request = AuthorizationRequest::delegated( + &proof, + &owner_binding, + expected_profile.clone(), + CapabilitySet::single(AuthorizationCapability::CommunityWrite), + Uuid::from_u128(0x500), + 1_000, + ) + .expect("synthetic delegated request is valid"); + let snapshot = match resolve_authorization( + &AllowProvider { + issued_at: 999, + fresh_until: 1_300, + }, + &request, + 1_000, + ProviderTimeout::new(Duration::from_secs(1)).expect("synthetic timeout is valid"), + ) + .await + { + AuthorizationOutcome::Allow(snapshot) => snapshot, + other => panic!("synthetic provider must allow, got {other:?}"), + }; + assert_eq!(snapshot.effective_until(), 1_050); + let binding_bound = BindingLeaseBound::new(&owner_binding, 1_400) + .expect("synthetic binding bound is valid"); + let admission = snapshot + .verified_owner_admission(&owner_binding) + .expect("delegated snapshot matches the exact owner binding"); + let input = AuthContextInput::new( + buzz_core::tenant::TenantContext::resolved(domain(), "relay.synthetic.example"), + Uuid::from_u128(0x500), + proof, + AuthorizedCommunityAccess::new(domain(), vec![Scope::MessagesWrite], None), + ); + let authorization = FederatedAuthorization::Delegated { + owner: owner_binding, + admission, + }; + let context = finalizer + .finalize_access( + input, + ResolvedFederatedPolicy::server_resolved_required( + domain(), + EnrollmentMode::Provisioned, + ), + authorization, + snapshot, + &expected_profile, + binding_bound, + access_policy(500), + LeaseVersion::INITIAL, + ) + .expect("delegated evidence finalizes"); + let lease = context.authorization_lease().expect("lease is present"); + assert_eq!(lease.owner_pubkey(), Some(owner)); + assert_eq!(lease.binding_id(), binding_id); + assert_eq!(lease.binding_version(), binding_version); + assert_eq!(lease.expires_at(), 1_045); + } +} diff --git a/crates/buzz-db/src/archived_identities.rs b/crates/buzz-db/src/archived_identities.rs index 941c0fc73..296ab89a6 100644 --- a/crates/buzz-db/src/archived_identities.rs +++ b/crates/buzz-db/src/archived_identities.rs @@ -7,7 +7,7 @@ use buzz_core::CommunityId; use chrono::{DateTime, Utc}; -use sqlx::{PgPool, Row as _}; +use sqlx::{PgPool, Postgres, Row as _, Transaction}; use crate::error::Result; @@ -76,6 +76,36 @@ pub async fn archive( Ok(result.rows_affected() > 0) } +/// Transaction-owned identity archive mutation. +#[allow(clippy::too_many_arguments)] +pub async fn archive_tx( + transaction: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &str, + consent_path: &str, + actor: &str, + reason: Option<&str>, + replaced_by: Option<&str>, + request_event_id: &str, +) -> Result { + let result = sqlx::query( + "INSERT INTO archived_identities \ + (community_id, pubkey, consent_path, actor, reason, replaced_by, request_event_id) \ + VALUES ($1, $2, $3, $4, $5, $6, $7) \ + ON CONFLICT (community_id, pubkey) DO NOTHING", + ) + .bind(community_id.as_uuid()) + .bind(pubkey) + .bind(consent_path) + .bind(actor) + .bind(reason) + .bind(replaced_by) + .bind(request_event_id) + .execute(&mut **transaction) + .await?; + Ok(result.rows_affected() > 0) +} + /// Unarchives an identity from `community_id`. /// /// Returns `true` if a row was deleted, `false` if the identity was not archived @@ -91,6 +121,21 @@ pub async fn unarchive(pool: &PgPool, community_id: CommunityId, pubkey: &str) - Ok(result.rows_affected() > 0) } +/// Transaction-owned identity unarchive mutation. +pub async fn unarchive_tx( + transaction: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &str, +) -> Result { + let result = + sqlx::query("DELETE FROM archived_identities WHERE community_id = $1 AND pubkey = $2") + .bind(community_id.as_uuid()) + .bind(pubkey) + .execute(&mut **transaction) + .await?; + Ok(result.rows_affected() > 0) +} + /// Returns all identities archived in `community_id`, ordered by archive time ascending. pub async fn list_archived( pool: &PgPool, diff --git a/crates/buzz-db/src/channel.rs b/crates/buzz-db/src/channel.rs index 13fe05280..fe30dcc84 100644 --- a/crates/buzz-db/src/channel.rs +++ b/crates/buzz-db/src/channel.rs @@ -183,6 +183,36 @@ pub async fn create_channel_with_id( description: Option<&str>, created_by: &[u8], ttl_seconds: Option, +) -> Result<(ChannelRecord, bool)> { + let mut tx = pool.begin().await?; + let result = create_channel_with_id_tx( + &mut tx, + community_id, + channel_id, + name, + channel_type, + visibility, + description, + created_by, + ttl_seconds, + ) + .await?; + tx.commit().await?; + Ok(result) +} + +/// Transaction-aware variant of [`create_channel_with_id`]. +#[allow(clippy::too_many_arguments)] +pub async fn create_channel_with_id_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + name: &str, + channel_type: ChannelType, + visibility: ChannelVisibility, + description: Option<&str>, + created_by: &[u8], + ttl_seconds: Option, ) -> Result<(ChannelRecord, bool)> { if created_by.len() != 32 { return Err(DbError::InvalidData(format!( @@ -202,8 +232,6 @@ pub async fn create_channel_with_id( return Err(DbError::InvalidData("channel name is required".into())); } - let mut tx = pool.begin().await?; - let rows_affected = sqlx::query( r#" INSERT INTO channels (id, community_id, name, channel_type, visibility, description, created_by, ttl_seconds, ttl_deadline) @@ -220,7 +248,7 @@ pub async fn create_channel_with_id( .bind(description) .bind(created_by) .bind(ttl_seconds) - .execute(&mut *tx) + .execute(&mut **tx) .await? .rows_affected(); @@ -242,7 +270,7 @@ pub async fn create_channel_with_id( .bind(channel_id) .bind(created_by) .bind(created_by) - .execute(&mut *tx) + .execute(&mut **tx) .await?; } @@ -260,11 +288,10 @@ pub async fn create_channel_with_id( ) .bind(community_id.as_uuid()) .bind(channel_id) - .fetch_one(&mut *tx) + .fetch_one(&mut **tx) .await?; let record = row_to_channel_record(row)?; - tx.commit().await?; Ok((record, was_created)) } @@ -351,7 +378,7 @@ const CHANNEL_MEMBERSHIP_LOCK_NAMESPACE: &str = "buzz_channel_membership:"; /// Take the per-channel membership lock. MUST be the first statement in the /// transaction that then reads roles/owner counts and writes membership, so the /// whole check-then-write sequence is atomic against a concurrent one. -async fn acquire_channel_membership_lock( +pub(crate) async fn acquire_channel_membership_lock( tx: &mut Transaction<'_, Postgres>, community_id: CommunityId, channel_id: Uuid, @@ -386,14 +413,7 @@ pub async fn add_member( role: MemberRole, invited_by: Option<&[u8]>, ) -> Result { - validate_member_pubkey(pubkey)?; - let mut tx = pool.begin().await?; - - // First statement: serialize the whole role-check / owner-count / upsert - // sequence against concurrent membership writes on this channel. - acquire_channel_membership_lock(&mut tx, community_id, channel_id).await?; - let record = add_member_tx(&mut tx, community_id, channel_id, pubkey, role, invited_by).await?; tx.commit().await?; Ok(record) @@ -418,7 +438,7 @@ pub enum ChannelAdmissionOutcome { IdentityBindingRequired, } -/// Add a channel member and optional corporate identity binding in one transaction. +/// Add a channel member and optional relay-verified identity binding atomically. pub async fn add_member_with_identity( pool: &PgPool, community_id: CommunityId, @@ -436,11 +456,11 @@ pub async fn add_member_with_identity( } let mut tx = pool.begin().await?; - // Keep this first: every channel membership writer shares this lock order. acquire_channel_membership_lock(&mut tx, community_id, channel_id).await?; - let member = - match add_member_tx(&mut tx, community_id, channel_id, pubkey, role, invited_by).await { + match add_member_after_lock_tx(&mut tx, community_id, channel_id, pubkey, role, invited_by) + .await + { Ok(member) => member, Err(error) => { tx.rollback().await?; @@ -491,7 +511,24 @@ fn validate_member_pubkey(pubkey: &[u8]) -> Result<()> { Ok(()) } -async fn add_member_tx( +/// Transaction-aware variant of [`add_member`]. +/// +/// This function owns the channel membership lock. Callers that already hold +/// it use the private after-lock helper so the lock order remains exact. +pub async fn add_member_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + pubkey: &[u8], + role: MemberRole, + invited_by: Option<&[u8]>, +) -> Result { + validate_member_pubkey(pubkey)?; + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + add_member_after_lock_tx(tx, community_id, channel_id, pubkey, role, invited_by).await +} + +async fn add_member_after_lock_tx( tx: &mut Transaction<'_, Postgres>, community_id: CommunityId, channel_id: Uuid, @@ -647,39 +684,54 @@ async fn add_member_tx( /// actor could commit after their role was read and this removal would proceed on /// a stale elevated role. /// -/// The `is_agent_owner` lookup deliberately runs *before* the transaction opens: -/// it borrows a second connection from `pool`, and issuing it while holding the -/// lock could deadlock against ourselves on a small pool. That is safe because -/// `agent_owner_pubkey` is immutable — [`crate::user::set_agent_owner`] only -/// updates it when it `IS NULL` (first-mint-wins), so its value cannot change -/// under us and needs no serialization. +/// The immutable agent-owner relationship is read on the caller-owned +/// transaction connection so this operation also composes with a sealed +/// authorization transaction without borrowing a second pool connection. pub async fn remove_member( pool: &PgPool, community_id: CommunityId, channel_id: Uuid, pubkey: &[u8], actor_pubkey: &[u8], +) -> Result<()> { + let mut tx = pool.begin().await?; + remove_member_tx(&mut tx, community_id, channel_id, pubkey, actor_pubkey).await?; + tx.commit().await?; + Ok(()) +} + +/// Transaction-aware variant of [`remove_member`]. +pub async fn remove_member_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + pubkey: &[u8], + actor_pubkey: &[u8], ) -> Result<()> { let is_self_remove = pubkey == actor_pubkey; - // Immutable, and must not be queried while holding the lock (second pool - // connection). Resolved up front so every *mutable* authorization read can - // sit behind the serialization point below. let actor_is_agent_owner = if is_self_remove { false } else { - crate::user::is_agent_owner(pool, community_id, pubkey, actor_pubkey).await? + sqlx::query_scalar::<_, bool>( + "SELECT agent_owner_pubkey = $3 FROM users \ + WHERE community_id = $1 AND pubkey = $2 AND agent_owner_pubkey IS NOT NULL", + ) + .bind(community_id.as_uuid()) + .bind(pubkey) + .bind(actor_pubkey) + .fetch_optional(&mut **tx) + .await? + .unwrap_or(false) }; - let mut tx = pool.begin().await?; - // First statement: serialize the actor-role check, the last-owner count and // the UPDATE against concurrent membership writes on this channel (same key // as `add_member`). - acquire_channel_membership_lock(&mut tx, community_id, channel_id).await?; + acquire_channel_membership_lock(tx, community_id, channel_id).await?; if !is_self_remove { - let actor_role_str = get_active_role_tx(&mut tx, community_id, channel_id, actor_pubkey) + let actor_role_str = get_active_role_tx(tx, community_id, channel_id, actor_pubkey) .await? .ok_or_else(|| DbError::AccessDenied("actor is not an active member".to_string()))?; let actor_role: MemberRole = actor_role_str.parse().map_err(|_| { @@ -695,7 +747,7 @@ pub async fn remove_member( // Defense-in-depth: prevent removing the last owner regardless of caller. // Callers (REST handlers, NIP-29 handlers) also check this, but the DB // layer enforces it as the final safety net. - let target_role = get_active_role_tx(&mut tx, community_id, channel_id, pubkey).await?; + let target_role = get_active_role_tx(tx, community_id, channel_id, pubkey).await?; if target_role.as_deref() == Some("owner") { let row = sqlx::query( "SELECT COUNT(*) as cnt FROM channel_members \ @@ -703,7 +755,7 @@ pub async fn remove_member( ) .bind(community_id.as_uuid()) .bind(channel_id) - .fetch_one(&mut *tx) + .fetch_one(&mut **tx) .await?; let owner_count: i64 = row.try_get("cnt")?; if owner_count <= 1 { @@ -724,14 +776,13 @@ pub async fn remove_member( .bind(community_id.as_uuid()) .bind(channel_id) .bind(pubkey) - .execute(&mut *tx) + .execute(&mut **tx) .await?; if result.rows_affected() == 0 { return Err(DbError::MemberNotFound(channel_id)); } - tx.commit().await?; Ok(()) } @@ -873,6 +924,84 @@ pub async fn get_accessible_channel_ids( .collect() } +/// Revalidate one actor's current read access to a channel in one database +/// statement. Open channels are readable by any authenticated relay actor; +/// private channels require an active membership. Deleted channels deny. +/// +/// This intentionally bypasses application caches. Callers use it at an +/// outbound release boundary after asynchronous fetch or queueing work. +pub async fn channel_read_authorized( + pool: &PgPool, + community_id: CommunityId, + channel_id: Uuid, + actor: &[u8], +) -> Result { + let allowed = sqlx::query_scalar::<_, bool>( + r#" + SELECT c.visibility::text <> 'private' + OR EXISTS ( + SELECT 1 + FROM channel_members cm + WHERE cm.community_id = c.community_id + AND cm.channel_id = c.id + AND cm.pubkey = $3 + AND cm.removed_at IS NULL + ) + FROM channels c + WHERE c.community_id = $1 + AND c.id = $2 + AND c.deleted_at IS NULL + "#, + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .bind(actor) + .fetch_optional(pool) + .await?; + Ok(allowed.unwrap_or(false)) +} + +/// Revalidate uncached read access to an entire channel set in one database +/// statement. Aggregate disclosures use this at their final release boundary +/// so authority for an earlier channel cannot go stale while later channels +/// are checked one at a time. +pub async fn channel_set_read_authorized( + pool: &PgPool, + community_id: CommunityId, + channel_ids: &[Uuid], + actor: &[u8], +) -> Result { + if channel_ids.is_empty() { + return Ok(true); + } + let allowed = sqlx::query_scalar::<_, bool>( + r#" + SELECT COUNT(DISTINCT c.id) = cardinality($2::uuid[]) + FROM channels c + WHERE c.community_id = $1 + AND c.id = ANY($2::uuid[]) + AND c.deleted_at IS NULL + AND ( + c.visibility::text <> 'private' + OR EXISTS ( + SELECT 1 + FROM channel_members cm + WHERE cm.community_id = c.community_id + AND cm.channel_id = c.id + AND cm.pubkey = $3 + AND cm.removed_at IS NULL + ) + ) + "#, + ) + .bind(community_id.as_uuid()) + .bind(channel_ids) + .bind(actor) + .fetch_one(pool) + .await?; + Ok(allowed) +} + /// Lists channels in a community, optionally filtered by visibility string. pub async fn list_channels( pool: &PgPool, @@ -1251,6 +1380,486 @@ pub struct ChannelUpdate { pub ttl_seconds: Option>, } +/// Transaction-owned NIP-29 channel mutation selected by the relay after +/// protocol-shape validation. Authorization is rechecked from locked rows in +/// [`apply_nip29_mutation_tx`]. +pub enum Nip29Mutation { + /// Create a channel and bootstrap the actor as its owner. + Create { + /// Stable client- or event-derived channel identifier. + channel_id: Uuid, + /// Canonical display name. + name: String, + /// Channel type. + channel_type: ChannelType, + /// Initial visibility. + visibility: ChannelVisibility, + /// Optional description. + description: Option, + /// Optional ephemeral lifetime. + ttl_seconds: Option, + }, + /// Add a member or change an active member's role. + PutUser { + /// Channel identifier. + channel_id: Uuid, + /// Target member key. + target: Vec, + /// Explicit role, or preserve/default when absent. + role: Option, + }, + /// Remove a member. + RemoveUser { + /// Channel identifier. + channel_id: Uuid, + /// Target member key. + target: Vec, + }, + /// Atomically edit channel metadata. + EditMetadata { + /// Channel identifier. + channel_id: Uuid, + /// Durable metadata columns. + updates: ChannelUpdate, + /// Optional topic replacement. + topic: Option, + /// Optional purpose replacement. + purpose: Option, + /// Optional archive transition. + archived: Option, + }, + /// Soft-delete a group and its relay-authored discovery rows. + DeleteGroup { + /// Channel identifier. + channel_id: Uuid, + /// Relay key used to scope discovery cleanup. + relay_pubkey: Vec, + }, + /// Join an open channel without changing an existing role. + Join { + /// Channel identifier. + channel_id: Uuid, + }, + /// Leave a channel without implicit membership creation. + Leave { + /// Channel identifier. + channel_id: Uuid, + }, +} + +/// Durable result of a transaction-owned NIP-29 projection. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct Nip29MutationOutcome { + /// Affected channel. + pub channel_id: Uuid, + /// Whether protected business state changed. + pub changed: bool, + /// Whether membership visibility changed and caches must be invalidated. + pub membership_changed: bool, + /// Whether channel visibility or lifecycle caches must be invalidated. + pub channel_changed: bool, +} + +async fn actor_owns_active_owner_agent_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + actor: &[u8], +) -> Result { + Ok(sqlx::query_scalar::<_, bool>( + "SELECT EXISTS( \ + SELECT 1 FROM channel_members cm \ + JOIN users u ON u.community_id = cm.community_id AND u.pubkey = cm.pubkey \ + WHERE cm.community_id = $1 AND cm.channel_id = $2 \ + AND cm.role = 'owner' AND cm.removed_at IS NULL \ + AND u.agent_owner_pubkey = $3 \ + )", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .bind(actor) + .fetch_one(&mut **tx) + .await?) +} + +async fn update_channel_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + mut updates: ChannelUpdate, +) -> Result<()> { + if let Some(name) = updates.name.as_mut() { + *name = buzz_core::channel::canonical_channel_name(name).to_owned(); + if name.is_empty() { + return Err(DbError::InvalidData("channel name is required".into())); + } + } + if updates.name.is_none() + && updates.description.is_none() + && updates.visibility.is_none() + && updates.ttl_seconds.is_none() + { + return Ok(()); + } + if updates.ttl_seconds.is_some() { + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1, 0))") + .bind(format!( + "buzz_channel_ttl:{}:{}", + community_id.as_uuid(), + channel_id + )) + .execute(&mut **tx) + .await?; + } + let result = sqlx::query( + "UPDATE channels SET \ + name = COALESCE($1, name), \ + description = COALESCE($2, description), \ + visibility = COALESCE($3::channel_visibility, visibility), \ + ttl_seconds = CASE WHEN $4 THEN $5 ELSE ttl_seconds END, \ + ttl_deadline = CASE WHEN $4 THEN CASE WHEN $5 IS NULL THEN NULL \ + ELSE NOW() + ($5 || ' seconds')::interval END ELSE ttl_deadline END, \ + updated_at = NOW() \ + WHERE community_id = $6 AND id = $7 AND deleted_at IS NULL", + ) + .bind(updates.name) + .bind(updates.description) + .bind(updates.visibility) + .bind(updates.ttl_seconds.is_some()) + .bind(updates.ttl_seconds.flatten()) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await?; + if result.rows_affected() == 0 { + return Err(DbError::ChannelNotFound(channel_id)); + } + Ok(()) +} + +/// Apply one NIP-29 durable projection on the same transaction that owns the +/// sealed authorization permit and event receipt. +pub async fn apply_nip29_mutation_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + actor: &[u8], + mutation: Nip29Mutation, +) -> Result { + if actor.len() != 32 { + return Err(DbError::InvalidData("actor pubkey must be 32 bytes".into())); + } + match mutation { + Nip29Mutation::Create { + channel_id, + name, + channel_type, + visibility, + description, + ttl_seconds, + } => { + let (_, changed) = create_channel_with_id_tx( + tx, + community_id, + channel_id, + &name, + channel_type, + visibility, + description.as_deref(), + actor, + ttl_seconds, + ) + .await?; + if !changed { + return Err(DbError::InvalidData("channel already exists".into())); + } + Ok(Nip29MutationOutcome { + channel_id, + changed, + membership_changed: changed, + channel_changed: changed, + }) + } + Nip29Mutation::PutUser { + channel_id, + target, + role, + } => { + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + let channel = get_channel_tx(tx, community_id, channel_id).await?; + let existing = get_active_role_tx(tx, community_id, channel_id, &target).await?; + let effective_role = match (role, existing.as_deref()) { + (Some(role), _) => role, + (None, Some(role)) => role.parse().map_err(|_| { + DbError::InvalidData(format!("invalid role in database: {role}")) + })?, + (None, None) => MemberRole::Member, + }; + if target != actor { + // Establish a row-level serialization point even when the + // target has never published a profile. Without this insert, + // `FOR SHARE` below cannot lock a missing row and a concurrent + // first profile could commit a restrictive policy before this + // membership transaction commits. + crate::user::ensure_user_tx(tx, community_id, &target).await?; + let policy = sqlx::query( + "SELECT channel_add_policy::text AS policy, agent_owner_pubkey \ + FROM users WHERE community_id = $1 AND pubkey = $2 \ + FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(&target) + .fetch_optional(&mut **tx) + .await?; + if let Some(policy) = policy { + let value: String = policy.try_get("policy")?; + let owner: Option> = policy.try_get("agent_owner_pubkey")?; + match value.as_str() { + "owner_only" if owner.as_deref() != Some(actor) => { + return Err(DbError::AccessDenied( + "only the agent owner may add this member".into(), + )); + } + "nobody" => { + return Err(DbError::AccessDenied( + "this member has disabled external channel additions".into(), + )); + } + _ => {} + } + } + } + let before = existing; + // The lock is reentrant for this transaction; `add_member_tx` + // retains the complete role and last-owner checks. + add_member_tx( + tx, + community_id, + channel_id, + &target, + effective_role, + Some(actor), + ) + .await?; + let changed = before.as_deref() != Some(effective_role.as_str()); + Ok(Nip29MutationOutcome { + channel_id, + changed, + membership_changed: changed, + channel_changed: channel.visibility == "open" && before.is_none(), + }) + } + Nip29Mutation::RemoveUser { channel_id, target } => { + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + get_channel_tx(tx, community_id, channel_id).await?; + if target != actor + && get_active_role_tx(tx, community_id, channel_id, actor) + .await? + .is_none() + { + return Err(DbError::AccessDenied( + "actor is not an active member".into(), + )); + } + remove_member_tx(tx, community_id, channel_id, &target, actor).await?; + Ok(Nip29MutationOutcome { + channel_id, + changed: true, + membership_changed: true, + channel_changed: true, + }) + } + Nip29Mutation::EditMetadata { + channel_id, + updates, + topic, + purpose, + archived, + } => { + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + sqlx::query("SELECT 1 FROM channels WHERE community_id = $1 AND id = $2 FOR UPDATE") + .bind(community_id.as_uuid()) + .bind(channel_id) + .fetch_optional(&mut **tx) + .await? + .ok_or(DbError::ChannelNotFound(channel_id))?; + let privileged = updates.name.is_some() + || updates.description.is_some() + || updates.visibility.is_some() + || updates.ttl_seconds.is_some() + || archived.is_some(); + let role = get_active_role_tx(tx, community_id, channel_id, actor).await?; + if privileged { + let elevated = role + .as_deref() + .and_then(|role| role.parse::().ok()) + .is_some_and(|role| role.is_elevated()); + if !elevated + && !actor_owns_active_owner_agent_tx(tx, community_id, channel_id, actor) + .await? + { + return Err(DbError::AccessDenied( + "actor is not authorized to edit channel metadata".into(), + )); + } + } else if (topic.is_some() || purpose.is_some()) && role.is_none() { + return Err(DbError::AccessDenied( + "actor is not an active member".into(), + )); + } + update_channel_tx(tx, community_id, channel_id, updates).await?; + if let Some(topic) = topic { + sqlx::query( + "UPDATE channels SET topic = $1, topic_set_by = $2, topic_set_at = NOW() \ + WHERE community_id = $3 AND id = $4 AND deleted_at IS NULL", + ) + .bind(topic) + .bind(actor) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await?; + } + if let Some(purpose) = purpose { + sqlx::query( + "UPDATE channels SET purpose = $1, purpose_set_by = $2, purpose_set_at = NOW() \ + WHERE community_id = $3 AND id = $4 AND deleted_at IS NULL", + ) + .bind(purpose) + .bind(actor) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await?; + } + if let Some(archived) = archived { + let result = if archived { + sqlx::query( + "UPDATE channels SET archived_at = NOW() WHERE community_id = $1 \ + AND id = $2 AND deleted_at IS NULL AND archived_at IS NULL", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await? + } else { + sqlx::query( + "UPDATE channels SET archived_at = NULL, ttl_deadline = CASE \ + WHEN ttl_seconds IS NOT NULL THEN NOW() + (ttl_seconds || ' seconds')::interval \ + ELSE ttl_deadline END WHERE community_id = $1 AND id = $2 \ + AND deleted_at IS NULL AND archived_at IS NOT NULL", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await? + }; + if result.rows_affected() == 0 { + return Err(DbError::AccessDenied( + "channel archive state did not permit the transition".into(), + )); + } + } + Ok(Nip29MutationOutcome { + channel_id, + changed: true, + membership_changed: false, + channel_changed: true, + }) + } + Nip29Mutation::DeleteGroup { + channel_id, + relay_pubkey, + } => { + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + sqlx::query("SELECT 1 FROM channels WHERE community_id = $1 AND id = $2 FOR UPDATE") + .bind(community_id.as_uuid()) + .bind(channel_id) + .fetch_optional(&mut **tx) + .await? + .ok_or(DbError::ChannelNotFound(channel_id))?; + let owner = get_active_role_tx(tx, community_id, channel_id, actor) + .await? + .as_deref() + == Some("owner"); + if !owner + && !actor_owns_active_owner_agent_tx(tx, community_id, channel_id, actor).await? + { + return Err(DbError::AccessDenied( + "only an owner may delete a group".into(), + )); + } + let changed = sqlx::query( + "UPDATE channels SET deleted_at = NOW() WHERE community_id = $1 \ + AND id = $2 AND deleted_at IS NULL", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .execute(&mut **tx) + .await? + .rows_affected() + > 0; + sqlx::query( + "UPDATE events SET deleted_at = NOW() WHERE community_id = $1 \ + AND channel_id = $2 AND pubkey = $3 AND deleted_at IS NULL \ + AND kind IN (39000, 39001, 39002)", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .bind(relay_pubkey) + .execute(&mut **tx) + .await?; + Ok(Nip29MutationOutcome { + channel_id, + changed, + membership_changed: changed, + channel_changed: changed, + }) + } + Nip29Mutation::Join { channel_id } => { + acquire_channel_membership_lock(tx, community_id, channel_id).await?; + let channel = get_channel_tx(tx, community_id, channel_id).await?; + if channel.visibility != "open" { + return Err(DbError::AccessDenied("channel is private".into())); + } + if get_active_role_tx(tx, community_id, channel_id, actor) + .await? + .is_some() + { + return Ok(Nip29MutationOutcome { + channel_id, + changed: false, + membership_changed: false, + channel_changed: false, + }); + } + add_member_tx( + tx, + community_id, + channel_id, + actor, + MemberRole::Member, + None, + ) + .await?; + Ok(Nip29MutationOutcome { + channel_id, + changed: true, + membership_changed: true, + channel_changed: true, + }) + } + Nip29Mutation::Leave { channel_id } => { + remove_member_tx(tx, community_id, channel_id, actor, actor).await?; + Ok(Nip29MutationOutcome { + channel_id, + changed: true, + membership_changed: true, + channel_changed: true, + }) + } + } +} + /// Updates channel metadata dynamically. /// /// At least one field must be provided; returns `InvalidData` otherwise. @@ -1587,16 +2196,92 @@ pub async fn get_member_role( Ok(row.map(|r| r.try_get("role")).transpose()?) } +/// Get the active role on the caller's transaction snapshot. +/// +/// Permission decisions that combine an event-owned binding with membership +/// use this together with [`crate::event::query_events_tx`] after selecting a +/// repeatable-read transaction, so both facts come from one database snapshot. +pub async fn get_member_role_tx( + transaction: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + pubkey: &[u8], +) -> Result> { + let row = sqlx::query( + "SELECT cm.role::text AS role FROM channel_members cm \ + JOIN channels c ON cm.community_id = c.community_id AND cm.channel_id = c.id AND c.deleted_at IS NULL \ + WHERE cm.community_id = $1 AND cm.channel_id = $2 AND cm.pubkey = $3 AND cm.removed_at IS NULL", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .bind(pubkey) + .fetch_optional(&mut **transaction) + .await?; + Ok(row.map(|r| r.try_get("role")).transpose()?) +} + +/// Lock and revalidate the ordinary member-or-open channel write predicate in +/// the caller's authorization transaction. +pub async fn require_channel_write_authority_tx( + transaction: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + channel_id: Uuid, + actor: &[u8], +) -> Result<()> { + let row = sqlx::query( + "SELECT visibility::text AS visibility, archived_at FROM channels \ + WHERE community_id = $1 AND id = $2 AND deleted_at IS NULL FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .fetch_optional(&mut **transaction) + .await? + .ok_or(DbError::ChannelNotFound(channel_id))?; + let archived_at: Option> = row.try_get("archived_at")?; + if archived_at.is_some() { + return Err(DbError::AccessDenied("channel is archived".into())); + } + let role: Option = sqlx::query_scalar( + "SELECT role::text FROM channel_members \ + WHERE community_id = $1 AND channel_id = $2 AND pubkey = $3 \ + AND removed_at IS NULL FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(channel_id) + .bind(actor) + .fetch_optional(&mut **transaction) + .await?; + let visibility: String = row.try_get("visibility")?; + if role.is_none() && visibility != "open" { + return Err(DbError::AccessDenied( + "actor is not a channel member".into(), + )); + } + Ok(()) +} + /// Archive ephemeral channels whose TTL deadline has passed. /// /// Returns the `(community_id, host, channel_id)` list that was archived. Idempotent — the /// `archived_at IS NULL` guard prevents double-archiving even if called /// concurrently from multiple relay pods. pub async fn reap_expired_ephemeral_channels(pool: &PgPool) -> Result> { + reap_expired_ephemeral_channels_excluding(pool, &[]).await +} + +/// Archive expired ephemeral channels except exact protected domains. +/// +/// The exclusion predicate is part of the `UPDATE`, so an Enforce row cannot +/// be claimed and mutated between an application-side mode check and commit. +pub async fn reap_expired_ephemeral_channels_excluding( + pool: &PgPool, + excluded_communities: &[Uuid], +) -> Result> { let rows = sqlx::query( "UPDATE channels AS ch SET archived_at = NOW() \ FROM communities AS c \ WHERE ch.community_id = c.id \ + AND NOT (ch.community_id = ANY($1::uuid[])) \ AND ch.ttl_seconds IS NOT NULL \ AND ch.ttl_deadline < NOW() \ AND ch.archived_at IS NULL \ @@ -1604,6 +2289,7 @@ pub async fn reap_expired_ephemeral_channels(pool: &PgPool) -> Result Result PgPool { - let database_url = - std::env::var("BUZZ_TEST_DATABASE_URL").unwrap_or_else(|_| 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()); PgPool::connect(&database_url) .await .expect("connect to test DB") @@ -1714,15 +2404,6 @@ mod tests { } } - async fn trusted_assertion_count(pool: &PgPool, community: CommunityId) -> i64 { - sqlx::query_scalar("SELECT COUNT(*) FROM events WHERE community_id = $1 AND kind = $2") - .bind(community.as_uuid()) - .bind(buzz_core::kind::KIND_USER_TRUSTED_ASSERTION as i32) - .fetch_one(pool) - .await - .expect("trusted assertion count") - } - async fn active_membership_count( pool: &PgPool, community: CommunityId, @@ -1749,7 +2430,6 @@ mod tests { let community_id = make_test_community(&pool).await; let community = CommunityId::from_uuid(community_id); let owner = random_pubkey(); - let non_member_inviter = random_pubkey(); let joiner = random_pubkey(); let channel = create_test_channel( &pool, @@ -1771,7 +2451,7 @@ mod tests { channel.id, &joiner, MemberRole::Member, - Some(&non_member_inviter), + Some(&random_pubkey()), Some(&identity), ) .await @@ -1789,7 +2469,6 @@ mod tests { .expect("binding lookup") .is_none() ); - assert_eq!(trusted_assertion_count(&pool, community).await, 0); } #[tokio::test] @@ -1851,7 +2530,6 @@ mod tests { active_membership_count(&pool, community, channel.id, &joiner).await, 0 ); - assert_eq!(trusted_assertion_count(&pool, community).await, 0); } #[tokio::test] @@ -1893,15 +2571,6 @@ mod tests { active_membership_count(&pool, community, channel.id, &joiner).await, 0 ); - assert!( - crate::identity_binding::get_active_identity_binding_by_pubkey( - &pool, community, &joiner, - ) - .await - .expect("binding lookup") - .is_none() - ); - assert_eq!(trusted_assertion_count(&pool, community).await, 0); } #[tokio::test] @@ -1997,88 +2666,6 @@ mod tests { assert_eq!(binding_count, 1); } - #[tokio::test] - #[ignore = "requires Postgres"] - async fn existing_member_and_non_corporate_paths_remain_idempotent() { - let pool = setup_pool().await; - let community_id = make_test_community(&pool).await; - let community = CommunityId::from_uuid(community_id); - let owner = random_pubkey(); - let existing_member = random_pubkey(); - let non_corporate_joiner = random_pubkey(); - let channel = create_test_channel( - &pool, - community_id, - "unchanged-admission-paths", - ChannelType::Stream, - ChannelVisibility::Private, - None, - &owner, - Some(3600), - ) - .await - .expect("create private huddle"); - - add_member( - &pool, - community, - channel.id, - &existing_member, - MemberRole::Member, - Some(&owner), - ) - .await - .expect("existing member add"); - add_member( - &pool, - community, - channel.id, - &existing_member, - MemberRole::Member, - Some(&owner), - ) - .await - .expect("existing member retry"); - assert_eq!( - active_membership_count(&pool, community, channel.id, &existing_member).await, - 1 - ); - - let outcome = add_member_with_identity( - &pool, - community, - channel.id, - &non_corporate_joiner, - MemberRole::Member, - Some(&owner), - None, - ) - .await - .expect("non-corporate admission"); - assert!(matches!( - outcome, - ChannelAdmissionOutcome::Joined { - identity_binding: None, - .. - } - )); - assert_eq!( - active_membership_count(&pool, community, channel.id, &non_corporate_joiner).await, - 1 - ); - assert!( - crate::identity_binding::get_active_identity_binding_by_pubkey( - &pool, - community, - &non_corporate_joiner, - ) - .await - .expect("binding lookup") - .is_none() - ); - assert_eq!(trusted_assertion_count(&pool, community).await, 0); - } - async fn insert_channel_with_id( pool: &PgPool, community_id: Uuid, @@ -2359,6 +2946,14 @@ mod tests { .await .expect("expire channel"); + let excluded = reap_expired_ephemeral_channels_excluding(&pool, &[community_id]) + .await + .expect("run excluded reaper"); + assert!( + !excluded.iter().any(|row| row.channel_id == channel.id), + "an excluded protected domain must remain untouched" + ); + let reaped = reap_expired_ephemeral_channels(&pool) .await .expect("run reaper"); @@ -3161,4 +3756,453 @@ mod tests { .expect("read role after restore"); assert_eq!(restored.as_deref(), Some("owner")); } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn nip29_create_is_owned_by_the_callers_transaction() { + let pool = setup_pool().await; + let community = CommunityId::from_uuid(make_test_community(&pool).await); + let actor = random_pubkey(); + let channel_id = Uuid::new_v4(); + let mut tx = pool.begin().await.expect("begin caller transaction"); + let outcome = apply_nip29_mutation_tx( + &mut tx, + community, + &actor, + Nip29Mutation::Create { + channel_id, + name: "sealed-channel".into(), + channel_type: ChannelType::Stream, + visibility: ChannelVisibility::Private, + description: None, + ttl_seconds: None, + }, + ) + .await + .expect("create projection"); + assert!(outcome.changed); + tx.rollback().await.expect("authorization rollback"); + + assert!(matches!( + get_channel(&pool, community, channel_id).await, + Err(DbError::ChannelNotFound(_)) + )); + let membership_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM channel_members WHERE community_id = $1 AND channel_id = $2", + ) + .bind(community.as_uuid()) + .bind(channel_id) + .fetch_one(&pool) + .await + .expect("count rolled-back members"); + assert_eq!(membership_count, 0); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn nip29_join_retry_is_idempotent_and_never_changes_an_existing_role() { + let pool = setup_pool().await; + let community_id = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_id); + let owner = random_pubkey(); + let member = random_pubkey(); + let channel = create_test_channel( + &pool, + community_id, + "sealed-join", + ChannelType::Stream, + ChannelVisibility::Open, + None, + &owner, + None, + ) + .await + .expect("create channel"); + + let mut tx = pool.begin().await.expect("begin first join"); + let first = apply_nip29_mutation_tx( + &mut tx, + community, + &member, + Nip29Mutation::Join { + channel_id: channel.id, + }, + ) + .await + .expect("first join"); + assert!(first.changed); + tx.commit().await.expect("commit first join"); + + let mut retry = pool.begin().await.expect("begin retry"); + let repeated = apply_nip29_mutation_tx( + &mut retry, + community, + &member, + Nip29Mutation::Join { + channel_id: channel.id, + }, + ) + .await + .expect("retry join"); + assert!(!repeated.changed); + retry.commit().await.expect("commit retry"); + + let members = get_members(&pool, community, channel.id) + .await + .expect("members"); + assert_eq!( + members + .iter() + .filter(|entry| entry.pubkey == member) + .count(), + 1 + ); + assert_eq!( + members + .iter() + .find(|entry| entry.pubkey == owner) + .map(|entry| entry.role.as_str()), + Some("owner") + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn nip29_put_user_serializes_with_target_policy_updates() { + let pool = setup_pool().await; + let community_id = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_id); + let actor = random_pubkey(); + let target = random_pubkey(); + ensure_user(&pool, community, &actor) + .await + .expect("ensure actor"); + ensure_user(&pool, community, &target) + .await + .expect("ensure target"); + set_channel_add_policy(&pool, community, &target, "anyone") + .await + .expect("allow external additions"); + let channel = create_test_channel( + &pool, + community_id, + "sealed-target-policy-race", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &actor, + None, + ) + .await + .expect("create channel"); + + let mut operation = pool.begin().await.expect("begin protected put-user"); + apply_nip29_mutation_tx( + &mut operation, + community, + &actor, + Nip29Mutation::PutUser { + channel_id: channel.id, + target: target.clone(), + role: Some(MemberRole::Member), + }, + ) + .await + .expect("authorize and stage target membership"); + + let update_pool = pool.clone(); + let update_target = target.clone(); + let mut policy_update = tokio::spawn(async move { + set_channel_add_policy(&update_pool, community, &update_target, "nobody").await + }); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(750), &mut policy_update) + .await + .is_err(), + "target policy update must wait for the transaction that authorized the addition" + ); + + operation + .commit() + .await + .expect("commit authorized addition before policy update"); + tokio::time::timeout(std::time::Duration::from_secs(10), policy_update) + .await + .expect("policy update proceeds after authorization transaction") + .expect("policy update task") + .expect("policy update succeeds"); + assert!( + is_member(&pool, community, channel.id, &target) + .await + .expect("membership after serialized commit"), + "the authorization transaction won the serialization order" + ); + + let denied_target = random_pubkey(); + ensure_user(&pool, community, &denied_target) + .await + .expect("ensure denied target"); + set_channel_add_policy(&pool, community, &denied_target, "nobody") + .await + .expect("deny external additions first"); + let mut denied = pool.begin().await.expect("begin denied put-user"); + let result = apply_nip29_mutation_tx( + &mut denied, + community, + &actor, + Nip29Mutation::PutUser { + channel_id: channel.id, + target: denied_target.clone(), + role: Some(MemberRole::Member), + }, + ) + .await; + assert!(matches!(result, Err(DbError::AccessDenied(_)))); + denied.rollback().await.expect("rollback denied put-user"); + assert!( + !is_member(&pool, community, channel.id, &denied_target) + .await + .expect("denied target membership"), + "a policy update that commits first must deny without membership" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn nip29_put_user_serializes_with_first_target_profile() { + let pool = setup_pool().await; + let community_id = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_id); + let actor = random_pubkey(); + ensure_user(&pool, community, &actor) + .await + .expect("ensure actor"); + let channel = create_test_channel( + &pool, + community_id, + "sealed-first-profile-race", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &actor, + None, + ) + .await + .expect("create channel"); + + // PutUser wins: its create-or-conflict establishes the target row and + // retains that row through membership commit. The first restrictive + // profile must wait and therefore takes effect only afterward. + let target = random_pubkey(); + let mut operation = pool.begin().await.expect("begin protected put-user"); + apply_nip29_mutation_tx( + &mut operation, + community, + &actor, + Nip29Mutation::PutUser { + channel_id: channel.id, + target: target.clone(), + role: Some(MemberRole::Member), + }, + ) + .await + .expect("stage absent target membership"); + + let update_pool = pool.clone(); + let update_target = target.clone(); + let mut first_profile = tokio::spawn(async move { + let mut profile_tx = update_pool.begin().await.expect("begin first profile"); + ensure_user_tx(&mut profile_tx, community, &update_target) + .await + .expect("create or observe target"); + set_channel_add_policy_tx(&mut profile_tx, community, &update_target, "nobody") + .await + .expect("set first profile policy"); + profile_tx.commit().await.expect("commit first profile") + }); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(750), &mut first_profile) + .await + .is_err(), + "first target profile must wait for the earlier PutUser transaction" + ); + operation.commit().await.expect("commit absent-target add"); + tokio::time::timeout(std::time::Duration::from_secs(10), first_profile) + .await + .expect("first profile proceeds after membership commit") + .expect("first profile task"); + assert!(is_member(&pool, community, channel.id, &target) + .await + .expect("membership after PutUser-first order")); + + // Profile wins: hold its newly inserted `nobody` row open. PutUser must + // wait at create-or-conflict, then re-read the committed restriction + // and deny before adding membership. + let denied_target = random_pubkey(); + let mut profile_tx = pool.begin().await.expect("begin winning profile"); + ensure_user_tx(&mut profile_tx, community, &denied_target) + .await + .expect("stage first target profile"); + set_channel_add_policy_tx(&mut profile_tx, community, &denied_target, "nobody") + .await + .expect("stage restrictive policy"); + + let put_pool = pool.clone(); + let put_actor = actor.clone(); + let put_target = denied_target.clone(); + let channel_id = channel.id; + let mut put_user = tokio::spawn(async move { + let mut put_tx = put_pool.begin().await.expect("begin waiting put-user"); + let result = apply_nip29_mutation_tx( + &mut put_tx, + community, + &put_actor, + Nip29Mutation::PutUser { + channel_id, + target: put_target, + role: Some(MemberRole::Member), + }, + ) + .await; + match result { + Ok(outcome) => { + put_tx.commit().await.expect("commit unexpected add"); + Ok(outcome) + } + Err(error) => { + put_tx.rollback().await.expect("rollback denied add"); + Err(error) + } + } + }); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(750), &mut put_user) + .await + .is_err(), + "PutUser must wait for the earlier first-profile transaction" + ); + profile_tx + .commit() + .await + .expect("commit restrictive profile"); + let result = tokio::time::timeout(std::time::Duration::from_secs(10), put_user) + .await + .expect("PutUser proceeds after first profile commit") + .expect("PutUser task"); + assert!(matches!(result, Err(DbError::AccessDenied(_)))); + assert!( + !is_member(&pool, community, channel.id, &denied_target) + .await + .expect("membership after profile-first order"), + "a restrictive first profile must deny without membership" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn member_removed_after_precheck_before_commit_denies_event() { + let pool = setup_pool().await; + let community_id = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_id); + let owner = random_pubkey(); + let member = random_pubkey(); + let channel = create_test_channel( + &pool, + community_id, + "sealed-write-race", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &owner, + None, + ) + .await + .expect("create channel"); + sqlx::query( + "INSERT INTO channel_members (community_id, channel_id, pubkey, role, invited_by) \ + VALUES ($1, $2, $3, 'member', $4)", + ) + .bind(community_id) + .bind(channel.id) + .bind(&member) + .bind(&owner) + .execute(&pool) + .await + .expect("add member"); + + let mut preflight = pool.begin().await.expect("begin preflight"); + require_channel_write_authority_tx(&mut preflight, community, channel.id, &member) + .await + .expect("member passes preflight"); + preflight.rollback().await.expect("release preflight locks"); + + sqlx::query( + "UPDATE channel_members SET removed_at = NOW() \ + WHERE community_id = $1 AND channel_id = $2 AND pubkey = $3", + ) + .bind(community_id) + .bind(channel.id) + .bind(&member) + .execute(&pool) + .await + .expect("remove member between boundaries"); + + let mut operation = pool.begin().await.expect("begin operation"); + let result = + require_channel_write_authority_tx(&mut operation, community, channel.id, &member) + .await; + assert!(matches!(result, Err(DbError::AccessDenied(_)))); + operation + .rollback() + .await + .expect("rollback denied operation"); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn open_to_private_after_precheck_denies_nonmember() { + let pool = setup_pool().await; + let community_id = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_id); + let owner = random_pubkey(); + let nonmember = random_pubkey(); + let channel = create_test_channel( + &pool, + community_id, + "sealed-visibility-race", + ChannelType::Stream, + ChannelVisibility::Open, + None, + &owner, + None, + ) + .await + .expect("create channel"); + + let mut preflight = pool.begin().await.expect("begin preflight"); + require_channel_write_authority_tx(&mut preflight, community, channel.id, &nonmember) + .await + .expect("open channel passes preflight"); + preflight.rollback().await.expect("release preflight locks"); + + sqlx::query( + "UPDATE channels SET visibility = 'private'::channel_visibility \ + WHERE community_id = $1 AND id = $2", + ) + .bind(community_id) + .bind(channel.id) + .execute(&pool) + .await + .expect("make channel private between boundaries"); + + let mut operation = pool.begin().await.expect("begin operation"); + let result = + require_channel_write_authority_tx(&mut operation, community, channel.id, &nonmember) + .await; + assert!(matches!(result, Err(DbError::AccessDenied(_)))); + operation + .rollback() + .await + .expect("rollback denied operation"); + } } diff --git a/crates/buzz-db/src/relay_invite.rs b/crates/buzz-db/src/relay_invite.rs index 7189a2381..a79c9afce 100644 --- a/crates/buzz-db/src/relay_invite.rs +++ b/crates/buzz-db/src/relay_invite.rs @@ -22,7 +22,8 @@ use buzz_core::invite::{ V2_SECRET_LEN, }; use chrono::{DateTime, Utc}; -use sqlx::{PgPool, Row as _}; +use sqlx::{PgPool, Postgres, Row as _, Transaction}; +use uuid::Uuid; use crate::error::Result; use crate::identity_binding::{BindIdentityResult, IdentityBindingConflict, IdentityBindingInput}; @@ -113,6 +114,21 @@ pub async fn mint_relay_invite( created_by: &str, ttl_secs: u64, max_uses: Option, +) -> Result { + let mut transaction = pool.begin().await?; + let invite = + mint_relay_invite_tx(&mut transaction, community, created_by, ttl_secs, max_uses).await?; + transaction.commit().await?; + Ok(invite) +} + +/// Mint an invite inside a caller-owned authorization transaction. +pub async fn mint_relay_invite_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + created_by: &str, + ttl_secs: u64, + max_uses: Option, ) -> Result { validate_mint_inputs(ttl_secs, max_uses)?; @@ -133,7 +149,7 @@ pub async fn mint_relay_invite( .bind(max_uses) .bind(expires_at) .bind(created_by) - .fetch_one(pool) + .fetch_one(&mut **transaction) .await?; let invite_id: uuid::Uuid = row.try_get("id")?; @@ -147,6 +163,32 @@ pub async fn mint_relay_invite( }) } +/// Lock and validate the relay role that may mint an invite. +/// +/// Enforcing callers invoke this inside the same transaction that writes the +/// invite and its authorization receipt. Legacy callers retain their existing +/// authorization flow. +pub async fn validate_relay_invite_minter_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + created_by: &str, +) -> Result<()> { + let role: Option = sqlx::query_scalar( + "SELECT role FROM relay_members \ + WHERE community_id = $1 AND pubkey = $2 FOR SHARE", + ) + .bind(community.as_uuid()) + .bind(created_by) + .fetch_optional(&mut **transaction) + .await?; + if !matches!(role.as_deref(), Some("owner" | "admin")) { + return Err(crate::error::DbError::InvalidData( + "invite mint authority changed before commit".into(), + )); + } + Ok(()) +} + fn log_claim_outcome( community: CommunityId, invite_id: Option, @@ -174,17 +216,28 @@ const RETENTION_SWEEP_BATCH_SIZE: i64 = 1_000; /// expiry index makes old rows drain first without turning cleanup into an /// unbounded transaction. pub async fn reap_expired_relay_invites(pool: &PgPool, cutoff: DateTime) -> Result { + reap_expired_relay_invites_excluding(pool, cutoff, &[]).await +} + +/// Delete expired invites outside exact protected Enforce domains. +pub async fn reap_expired_relay_invites_excluding( + pool: &PgPool, + cutoff: DateTime, + excluded_communities: &[Uuid], +) -> Result { let result = sqlx::query( "DELETE FROM relay_invites \ WHERE (community_id, id) IN (\ SELECT community_id, id FROM relay_invites \ WHERE expires_at < $1 \ + AND NOT (community_id = ANY($3::uuid[])) \ ORDER BY expires_at \ LIMIT $2\ )", ) .bind(cutoff) .bind(RETENTION_SWEEP_BATCH_SIZE) + .bind(excluded_communities) .execute(pool) .await?; @@ -217,11 +270,41 @@ pub async fn claim_relay_invite_with_identity( claimer_pubkey: &str, policy_version: Option<&str>, identity: Option<&IdentityBindingInput<'_>>, +) -> Result { + let mut transaction = pool.begin().await?; + let outcome = claim_relay_invite_with_identity_tx( + &mut transaction, + community, + token_hash, + claimer_pubkey, + policy_version, + identity, + ) + .await?; + if matches!( + outcome, + ClaimOutcome::Joined { .. } | ClaimOutcome::AlreadyMember { .. } + ) { + transaction.commit().await?; + } else { + transaction.rollback().await?; + } + Ok(outcome) +} + +/// Stage invite consumption, binding, membership, and policy evidence inside +/// a caller-owned authorization transaction. +pub async fn claim_relay_invite_with_identity_tx( + tx: &mut Transaction<'_, Postgres>, + community: CommunityId, + token_hash: &[u8; 32], + claimer_pubkey: &str, + policy_version: Option<&str>, + identity: Option<&IdentityBindingInput<'_>>, ) -> Result { crate::identity_binding::validate_membership_identity_key(claimer_pubkey, identity)?; - let mut tx = pool.begin().await?; sqlx::query("SET LOCAL lock_timeout = '3s'") - .execute(&mut *tx) + .execute(&mut **tx) .await?; // 2. SELECT FOR UPDATE — lock the invite row for the duration of this txn. @@ -233,12 +316,11 @@ pub async fn claim_relay_invite_with_identity( ) .bind(community.as_uuid()) .bind(token_hash) - .fetch_optional(&mut *tx) + .fetch_optional(&mut **tx) .await?; // 3. No matching invite. let Some(invite) = row else { - tx.rollback().await?; log_claim_outcome(community, None, "invalid", None, None); return Ok(ClaimOutcome::Invalid); }; @@ -252,7 +334,6 @@ pub async fn claim_relay_invite_with_identity( // not authorize fresh policy-acceptance evidence, even for an existing // member; exhausted-but-live invites remain valid for idempotent retries. if expires_at <= Utc::now() { - tx.rollback().await?; log_claim_outcome( community, Some(invite_id), @@ -264,12 +345,10 @@ pub async fn claim_relay_invite_with_identity( } let identity_binding = if let Some(identity) = identity { - match crate::identity_binding::bind_or_validate_identity_tx(&mut tx, community, identity) - .await? + match crate::identity_binding::bind_or_validate_identity_tx(tx, community, identity).await? { binding @ (BindIdentityResult::Created | BindIdentityResult::Matched) => Some(binding), BindIdentityResult::Conflict(conflict) => { - tx.rollback().await?; log_claim_outcome( community, Some(invite_id), @@ -280,7 +359,6 @@ pub async fn claim_relay_invite_with_identity( return Ok(ClaimOutcome::IdentityConflict(conflict)); } BindIdentityResult::Revoked => { - tx.rollback().await?; log_claim_outcome( community, Some(invite_id), @@ -291,7 +369,6 @@ pub async fn claim_relay_invite_with_identity( return Ok(ClaimOutcome::IdentityRevoked); } BindIdentityResult::BindingRequired => { - tx.rollback().await?; log_claim_outcome( community, Some(invite_id), @@ -313,7 +390,7 @@ pub async fn claim_relay_invite_with_identity( sqlx::query("SELECT 1 FROM relay_members WHERE community_id = $1 AND pubkey = $2") .bind(community.as_uuid()) .bind(claimer_pubkey) - .fetch_optional(&mut *tx) + .fetch_optional(&mut **tx) .await?; if existing.is_some() { @@ -326,10 +403,9 @@ pub async fn claim_relay_invite_with_identity( .bind(community.as_uuid()) .bind(claimer_pubkey) .bind(version) - .execute(&mut *tx) + .execute(&mut **tx) .await?; } - tx.commit().await?; log_claim_outcome( community, Some(invite_id), @@ -347,7 +423,6 @@ pub async fn claim_relay_invite_with_identity( // 7. Capacity check. if let Some(mu) = max_uses { if use_count >= mu { - tx.rollback().await?; log_claim_outcome( community, Some(invite_id), @@ -369,7 +444,7 @@ pub async fn claim_relay_invite_with_identity( ) .bind(community.as_uuid()) .bind(claimer_pubkey) - .execute(&mut *tx) + .execute(&mut **tx) .await? .rows_affected() > 0; @@ -384,12 +459,11 @@ pub async fn claim_relay_invite_with_identity( .bind(community.as_uuid()) .bind(claimer_pubkey) .bind(version) - .execute(&mut *tx) + .execute(&mut **tx) .await?; } if !inserted { - tx.commit().await?; log_claim_outcome( community, Some(invite_id), @@ -410,12 +484,9 @@ pub async fn claim_relay_invite_with_identity( .bind(new_use_count) .bind(community.as_uuid()) .bind(invite_id) - .execute(&mut *tx) + .execute(&mut **tx) .await?; - // 11. Commit. - tx.commit().await?; - let new_uses_remaining = max_uses.map(|mu| mu - new_use_count); log_claim_outcome( @@ -803,6 +874,12 @@ mod tests { .await .expect("age old invite"); + assert_eq!( + reap_expired_relay_invites_excluding(&pool, cutoff, &[*community.as_uuid()]) + .await + .expect("exclude protected invites"), + 0 + ); assert_eq!( reap_expired_relay_invites(&pool, cutoff) .await diff --git a/crates/buzz-db/src/relay_members.rs b/crates/buzz-db/src/relay_members.rs index affddd817..3e7ffd92e 100644 --- a/crates/buzz-db/src/relay_members.rs +++ b/crates/buzz-db/src/relay_members.rs @@ -7,7 +7,7 @@ //! lowercase hex strings. use chrono::{DateTime, Utc}; -use sqlx::{PgPool, Row as _}; +use sqlx::{PgPool, Postgres, Row as _, Transaction}; use crate::error::Result; use crate::identity_binding::{BindIdentityResult, IdentityBindingConflict, IdentityBindingInput}; @@ -93,6 +93,35 @@ pub async fn get_relay_member( .map_err(crate::error::DbError::from) } +/// Return and share-lock a relay member inside a caller-owned authorization +/// transaction so a role decision remains stable through commit. +pub async fn get_relay_member_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + pubkey: &str, +) -> Result> { + let row = sqlx::query( + "SELECT pubkey, role, added_by, created_at, updated_at \ + FROM relay_members WHERE community_id = $1 AND pubkey = $2 FOR SHARE", + ) + .bind(community.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut **transaction) + .await?; + + row.map(|r| -> std::result::Result { + Ok(RelayMember { + pubkey: r.try_get("pubkey")?, + role: r.try_get("role")?, + added_by: r.try_get("added_by")?, + created_at: r.try_get("created_at")?, + updated_at: r.try_get("updated_at")?, + }) + }) + .transpose() + .map_err(crate::error::DbError::from) +} + /// Returns all relay members of `community` ordered by `created_at` ascending. pub async fn list_relay_members(pool: &PgPool, community: CommunityId) -> Result> { let rows = sqlx::query( @@ -142,6 +171,27 @@ pub async fn add_relay_member( Ok(result.rows_affected() > 0) } +/// Transaction-owned relay member insertion. +pub async fn add_relay_member_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + pubkey: &str, + role: &str, + added_by: Option<&str>, +) -> Result { + let result = sqlx::query( + "INSERT INTO relay_members (community_id, pubkey, role, added_by) \ + VALUES ($1, $2, $3, $4) ON CONFLICT (community_id, pubkey) DO NOTHING", + ) + .bind(community.as_uuid()) + .bind(pubkey) + .bind(role) + .bind(added_by) + .execute(&mut **transaction) + .await?; + Ok(result.rows_affected() > 0) +} + /// Claims relay membership via an invite and atomically persists policy evidence. /// /// Returns `true` when membership was inserted, or `false` when the pubkey was @@ -319,6 +369,35 @@ pub async fn remove_relay_member( } } +/// Transaction-owned relay member removal with owner protection. +pub async fn remove_relay_member_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + pubkey: &str, +) -> Result { + let result = sqlx::query( + "DELETE FROM relay_members \ + WHERE community_id = $1 AND pubkey = $2 AND role <> 'owner'", + ) + .bind(community.as_uuid()) + .bind(pubkey) + .execute(&mut **transaction) + .await?; + if result.rows_affected() > 0 { + return Ok(RemoveResult::Removed); + } + let exists = sqlx::query("SELECT 1 FROM relay_members WHERE community_id = $1 AND pubkey = $2") + .bind(community.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut **transaction) + .await?; + Ok(if exists.is_some() { + RemoveResult::IsOwner + } else { + RemoveResult::NotFound + }) +} + /// Removes a relay member only if their current role matches `expected_role`. /// /// The delete and the role check are collapsed into a single @@ -375,6 +454,38 @@ pub async fn remove_relay_member_if_role( } } +/// Transaction-owned role-conditional relay member removal. +pub async fn remove_relay_member_if_role_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + pubkey: &str, + expected_role: &str, +) -> Result { + let result = sqlx::query( + "DELETE FROM relay_members WHERE community_id = $1 AND pubkey = $2 AND role = $3", + ) + .bind(community.as_uuid()) + .bind(pubkey) + .bind(expected_role) + .execute(&mut **transaction) + .await?; + if result.rows_affected() > 0 { + return Ok(RemoveResult::Removed); + } + let role = sqlx::query_scalar::<_, String>( + "SELECT role FROM relay_members WHERE community_id = $1 AND pubkey = $2 FOR SHARE", + ) + .bind(community.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut **transaction) + .await?; + Ok(match role.as_deref() { + None => RemoveResult::NotFound, + Some("owner") => RemoveResult::IsOwner, + Some(_) => RemoveResult::RoleMismatch, + }) +} + /// Updates the role of an existing relay member in `community`. Returns `true` /// if updated. pub async fn update_relay_member_role( @@ -395,6 +506,25 @@ pub async fn update_relay_member_role( Ok(result.rows_affected() > 0) } +/// Transaction-owned role update with owner protection. +pub async fn update_relay_member_role_tx( + transaction: &mut Transaction<'_, Postgres>, + community: CommunityId, + pubkey: &str, + new_role: &str, +) -> Result { + let result = sqlx::query( + "UPDATE relay_members SET role = $1, updated_at = now() \ + WHERE community_id = $2 AND pubkey = $3 AND role <> 'owner'", + ) + .bind(new_role) + .bind(community.as_uuid()) + .bind(pubkey) + .execute(&mut **transaction) + .await?; + Ok(result.rows_affected() > 0) +} + /// Ensures the configured owner pubkey holds the `"owner"` role *in /// `community`*, and demotes any other owners in that community to `"admin"`. /// This handles owner rotation: if `RELAY_OWNER_PUBKEY` changes, the old owner diff --git a/crates/buzz-db/src/user.rs b/crates/buzz-db/src/user.rs index 066fb5f5c..9902e9f00 100644 --- a/crates/buzz-db/src/user.rs +++ b/crates/buzz-db/src/user.rs @@ -4,6 +4,7 @@ use crate::error::Result; use buzz_core::CommunityId; use sqlx::PgPool; use sqlx::Row; +use sqlx::{Postgres, Transaction}; /// A user's profile fields. #[derive(Debug, Clone)] @@ -54,6 +55,63 @@ pub async fn ensure_user(pool: &PgPool, community_id: CommunityId, pubkey: &[u8] Ok(result.rows_affected() == 1) } +/// Ensure a user row exists inside a caller-owned transaction. +pub async fn ensure_user_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &[u8], +) -> Result { + let result = sqlx::query( + "INSERT INTO users (community_id, pubkey) VALUES ($1, $2) \ + ON CONFLICT (community_id, pubkey) DO NOTHING", + ) + .bind(community_id.as_uuid()) + .bind(pubkey) + .execute(&mut **tx) + .await?; + Ok(result.rows_affected() > 0) +} + +/// Apply absolute kind:0 profile state inside a caller-owned transaction. +/// A contested NIP-05 handle leaves the prior handle unchanged while updating +/// the remaining fields, matching the legacy compatibility behavior. +#[allow(clippy::too_many_arguments)] +pub async fn replace_user_profile_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &[u8], + display_name: &str, + avatar_url: &str, + about: &str, + nip05_handle: &str, +) -> Result<()> { + let contested: bool = !nip05_handle.is_empty() + && sqlx::query_scalar::<_, bool>( + "SELECT EXISTS(SELECT 1 FROM users WHERE community_id = $1 \ + AND LOWER(nip05_handle) = LOWER($2) AND pubkey <> $3)", + ) + .bind(community_id.as_uuid()) + .bind(nip05_handle) + .bind(pubkey) + .fetch_one(&mut **tx) + .await?; + sqlx::query( + "UPDATE users SET display_name = NULLIF($1, ''), avatar_url = NULLIF($2, ''), \ + about = NULLIF($3, ''), nip05_handle = CASE WHEN $4 THEN nip05_handle \ + ELSE NULLIF($5, '') END WHERE community_id = $6 AND pubkey = $7", + ) + .bind(display_name) + .bind(avatar_url) + .bind(about) + .bind(contested) + .bind(nip05_handle) + .bind(community_id.as_uuid()) + .bind(pubkey) + .execute(&mut **tx) + .await?; + Ok(()) +} + /// Get a single user record by pubkey. pub async fn get_user( pool: &PgPool, @@ -368,6 +426,26 @@ pub async fn is_agent_owner( Ok(row.unwrap_or(false)) } +/// Share-lock and validate an agent-owner relationship inside a caller-owned +/// authorization transaction. +pub async fn is_agent_owner_tx( + transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>, + community_id: CommunityId, + target_pubkey: &[u8], + actor_pubkey: &[u8], +) -> Result { + let owner = sqlx::query_scalar::<_, Vec>( + "SELECT agent_owner_pubkey FROM users \ + WHERE community_id = $1 AND pubkey = $2 AND agent_owner_pubkey IS NOT NULL \ + FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(target_pubkey) + .fetch_optional(&mut **transaction) + .await?; + Ok(owner.is_some_and(|owner| owner == actor_pubkey)) +} + /// Set the channel_add_policy for a user. /// Returns an error if the pubkey is not found (rows_affected == 0). /// Returns an error if `policy` is not one of the valid ENUM values. @@ -398,6 +476,35 @@ pub async fn set_channel_add_policy( Ok(()) } +/// Set a channel-add policy inside a caller-owned transaction. +pub async fn set_channel_add_policy_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &[u8], + policy: &str, +) -> Result<()> { + if !matches!(policy, "anyone" | "owner_only" | "nobody") { + return Err(crate::error::DbError::InvalidData(format!( + "invalid channel_add_policy: {policy}" + ))); + } + let result = sqlx::query( + "UPDATE users SET channel_add_policy = $1::channel_add_policy \ + WHERE community_id = $2 AND pubkey = $3", + ) + .bind(policy) + .bind(community_id.as_uuid()) + .bind(pubkey) + .execute(&mut **tx) + .await?; + if result.rows_affected() == 0 { + return Err(crate::error::DbError::NotFound( + "pubkey not found in users table".into(), + )); + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/buzz-relay/src/api/invites.rs b/crates/buzz-relay/src/api/invites.rs index c8cd2b112..b78c4b81a 100644 --- a/crates/buzz-relay/src/api/invites.rs +++ b/crates/buzz-relay/src/api/invites.rs @@ -23,8 +23,17 @@ use axum::{ }; use serde::Deserialize; use serde_json::Value; +use sha2::{Digest, Sha256}; +use crate::authorization_runtime::executor::{ + begin_authorized_enrollment, begin_authorized_operation, AuthorizedEnrollmentStart, + AuthorizedOperationStart, ProtectedOperationId, +}; +use crate::authorization_runtime::finalization::AuthorizationMode; +use crate::authorization_runtime::transport::authorize_enrollment_if_configured; +use crate::authorization_runtime::transport::authorize_if_configured; use crate::handlers::side_effects::{publish_nip43_member_added, publish_nip43_membership_list}; +use buzz_auth::AuthorizationCapability; use buzz_core::invite::{ hash_v2_code, validate_v2_code, DEFAULT_INVITE_TTL_SECS, MAX_INVITE_TTL_SECS, MAX_INVITE_USES, MIN_INVITE_TTL_SECS, V2_PREFIX, @@ -108,6 +117,16 @@ pub struct AcceptPolicyRequest { pub age_confirmed: bool, } +fn stable_correlation_from_proof(proof: &buzz_auth::VerifiedNostrProof) -> uuid::Uuid { + let fingerprint = proof.operation_binding().fingerprint(); + let mut bytes = [0_u8; 16]; + bytes.copy_from_slice(&fingerprint[..16]); + if bytes == [0; 16] { + bytes[15] = 1; + } + uuid::Uuid::from_bytes(bytes) +} + /// Public join policy shared by every client-side join surface. pub async fn join_policy(State(state): State>) -> Json { match &state.config.join_policy { @@ -236,7 +255,9 @@ async fn authenticate( ( buzz_core::TenantContext, nostr::PublicKey, - crate::corporate_identity::CorporateIdentityProof, + Option, + Arc, + Option>, ), (StatusCode, Json), > { @@ -254,34 +275,89 @@ async fn authenticate( })?; let url = bridge::nip98_expected_url(&state.config.relay_url, &tenant, path); - let (pubkey, event_id_bytes) = bridge::verify_bridge_auth_with_options( + let (pubkey, event_id_bytes, verified_proof) = bridge::verify_protected_bridge_auth( headers, "POST", &url, Some(body), true, // invites always require NIP-98; no X-Pubkey dev fallback true, // POST bodies must be covered by a payload tag + tenant.community(), )?; - bridge::check_nip98_replay(state, &tenant, event_id_bytes).await?; + if state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(tenant.community())) + != Some(AuthorizationMode::Enforce) + { + bridge::check_nip98_replay(state, &tenant, event_id_bytes).await?; + } - let identity_jwt = crate::corporate_identity::identity_jwt_from_headers( + let identity_assertion = crate::corporate_identity::identity_assertion_from_headers( + state, + tenant.community(), headers, - &state.config.corporate_identity, - ); + ) + .map_err(crate::corporate_identity::CorporateIdentityError::into_api_error)?; let auth_tag = headers .get("x-auth-tag") .and_then(|value| value.to_str().ok()); - let identity_proof = crate::corporate_identity::verify_corporate_identity( + let identity_lane = + crate::authorization_runtime::transport::legacy_identity_lane(state, tenant.community()); + let identity_proof = match crate::corporate_identity::verify_corporate_identity( state, tenant.community(), pubkey, - identity_jwt.as_deref(), + identity_assertion.as_ref(), auth_tag, ) .await - .map_err(|error| error.into_api_error())?; + { + Ok(proof) => Some(proof), + Err(error) + if identity_lane + == crate::authorization_runtime::transport::LegacyIdentityLane::ObserveOnly => + { + tracing::warn!(error = ?error, "observational invite identity verification unavailable"); + None + } + Err(error) => return Err(error.into_api_error()), + }; - Ok((tenant, pubkey, identity_proof)) + let verified_proof = bridge::retain_bridge_proof(verified_proof, auth_tag)? + .ok_or_else(|| api_error(StatusCode::UNAUTHORIZED, "NIP-98 evidence required"))?; + + let enrollment_assertion = if state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(tenant.community())) + == Some(AuthorizationMode::Enforce) + { + let now = state + .corporate_identity + .as_ref() + .ok_or_else(|| api_error(StatusCode::FORBIDDEN, "relay identity verification failed"))? + .authorization_now() + .map_err(|error| error.into_api_error())?; + crate::corporate_identity::verified_assertion_for_proof( + identity_proof.as_ref().ok_or_else(|| { + api_error(StatusCode::FORBIDDEN, "relay identity verification failed") + })?, + tenant.community(), + buzz_auth::AuthTransport::HttpBridge, + now, + ) + .map_err(|error| error.into_api_error())? + .map(Arc::new) + } else { + None + }; + + Ok(( + tenant, + pubkey, + identity_proof, + verified_proof, + enrollment_assertion, + )) } async fn record_atomic_identity_rejection( @@ -314,7 +390,7 @@ pub async fn mint_invite( headers: HeaderMap, body: axum::body::Bytes, ) -> Result, (StatusCode, Json)> { - let (tenant, pubkey, identity_proof) = + let (tenant, pubkey, identity_proof, verified_proof, verified_assertion) = authenticate(&state, &headers, "/api/invites", &body).await?; // Authz mirrors kind:9030 (add member): owner or admin only. @@ -344,24 +420,33 @@ pub async fn mint_invite( }; let (ttl, max_uses) = validate_mint_request(&request)?; - crate::corporate_identity::finalize_corporate_identity( + let protected_authority = authorize_if_configured( &state, - tenant.community(), - pubkey, - identity_proof, + Arc::clone(&verified_proof), + verified_assertion, + AuthorizationCapability::InviteMint, + stable_correlation_from_proof(&verified_proof), + "invite.mint", ) .await - .map_err(|error| error.into_api_error())?; - - // Mint a v2 opaque, database-backed invite. - let invite = state - .db - .mint_relay_invite(tenant.community(), &sender_hex, ttl, max_uses) - .await - .map_err(|error| match error { - buzz_db::DbError::InvalidData(message) => api_error(StatusCode::BAD_REQUEST, &message), - error => internal_error(&format!("invite mint: {error}")), - })?; + .map_err(|error| { + tracing::warn!(error = %error, "invite mint: protected authorization denied"); + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })?; + if crate::authorization_runtime::transport::legacy_identity_lane(&state, tenant.community()) + == crate::authorization_runtime::transport::LegacyIdentityLane::Legacy + { + if let Some(identity_proof) = identity_proof { + crate::corporate_identity::finalize_corporate_identity( + &state, + tenant.community(), + pubkey, + identity_proof, + ) + .await + .map_err(|error| error.into_api_error())?; + } + } // Same TLS-posture logic as nip98_expected_url: wss deployments get an // https landing page URL, ws dev/test deployments get http. @@ -371,25 +456,90 @@ pub async fn mint_invite( "http" }; - tracing::info!( - community = %tenant.community(), - minted_by = %sender_hex, - invite_id = %invite.invite_id, - expires_at = %invite.expires_at, - max_uses = ?invite.max_uses, - "relay invite minted" - ); + let build_response = |invite: buzz_db::relay_invite::MintedInvite| { + tracing::info!( + community = %tenant.community(), + minted_by = %sender_hex, + invite_id = %invite.invite_id, + expires_at = %invite.expires_at, + max_uses = ?invite.max_uses, + "relay invite minted" + ); + serde_json::json!({ + "code": invite.code, + "expires_at": invite.expires_at.timestamp() as u64, + "max_uses": invite.max_uses, + "uses_remaining": invite.uses_remaining, + "url": format!("{scheme}://{}/invite/{}", tenant.host(), invite.code), + }) + }; - // expires_at as unix seconds for the response contract. - let expires_at_unix = invite.expires_at.timestamp() as u64; + if protected_authority.is_enforcing() { + let operation_id = ProtectedOperationId::derive( + tenant.community(), + "invite.mint.v1", + &verified_proof.operation_binding().fingerprint(), + ) + .map_err(|error| internal_error(&format!("invite mint: {error}")))?; + let mut digest = Sha256::new(); + digest.update(b"buzz-invite-mint-v1"); + digest.update(ttl.to_be_bytes()); + digest.update(max_uses.unwrap_or_default().to_be_bytes()); + let permit = protected_authority + .seal_postgres_mutation(operation_id, "invite.mint.v1", digest.finalize().into()) + .map_err(|error| { + tracing::warn!(error = %error, "invite mint: protected authorization denied"); + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })? + .ok_or_else(|| api_error(StatusCode::FORBIDDEN, "protected authorization denied"))?; + let response = match begin_authorized_operation(&state, permit) + .await + .map_err(|error| internal_error(&format!("invite mint: {error}")))? + { + AuthorizedOperationStart::Replay(payload) => serde_json::from_slice(&payload) + .map_err(|error| internal_error(&format!("invite mint replay: {error}")))?, + AuthorizedOperationStart::Execute(mut operation) => { + buzz_db::relay_invite::validate_relay_invite_minter_tx( + operation.transaction(), + tenant.community(), + &sender_hex, + ) + .await + .map_err(|error| { + tracing::warn!(error = %error, "invite mint authority changed before commit"); + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })?; + let invite = buzz_db::relay_invite::mint_relay_invite_tx( + operation.transaction(), + tenant.community(), + &sender_hex, + ttl, + max_uses, + ) + .await + .map_err(|error| internal_error(&format!("invite mint: {error}")))?; + let response = build_response(invite); + let payload = serde_json::to_vec(&response) + .map_err(|error| internal_error(&format!("invite mint: {error}")))?; + operation + .commit(&payload) + .await + .map_err(|error| internal_error(&format!("invite mint: {error}")))?; + response + } + }; + return Ok(Json(response)); + } - Ok(Json(serde_json::json!({ - "code": invite.code, - "expires_at": expires_at_unix, - "max_uses": invite.max_uses, - "uses_remaining": invite.uses_remaining, - "url": format!("{scheme}://{}/invite/{}", tenant.host(), invite.code), - }))) + let invite = state + .db + .mint_relay_invite(tenant.community(), &sender_hex, ttl, max_uses) + .await + .map_err(|error| match error { + buzz_db::DbError::InvalidData(message) => api_error(StatusCode::BAD_REQUEST, &message), + error => internal_error(&format!("invite mint: {error}")), + })?; + Ok(Json(build_response(invite))) } /// Claim an invite code — `POST /api/invites/claim`, NIP-98 signed by the @@ -403,9 +553,14 @@ pub async fn claim_invite( headers: HeaderMap, body: axum::body::Bytes, ) -> Result, (StatusCode, Json)> { - let (tenant, pubkey, identity_proof) = + let (tenant, pubkey, identity_proof, verified_proof, enrollment_assertion) = authenticate(&state, &headers, "/api/invites/claim", &body).await?; + let enforcing = state + .protected_transport() + .and_then(|runtime| runtime.mode_for_domain(tenant.community())) + == Some(AuthorizationMode::Enforce); + if claim_rate_limited(&state, tenant.community(), &pubkey) { return Err(api_error( StatusCode::TOO_MANY_REQUESTS, @@ -418,7 +573,10 @@ pub async fn claim_invite( // Invite admission must be coupled to the identity being admitted. A // delegated owner proof can become stale between verification and the // invite transaction, so bootstrap claims require the joiner's direct JWT. - if crate::corporate_identity::proof_is_delegated(&identity_proof) { + if identity_proof + .as_ref() + .is_some_and(crate::corporate_identity::proof_is_delegated) + { return Err(api_error( StatusCode::FORBIDDEN, "direct relay identity required for invite claim", @@ -428,6 +586,13 @@ pub async fn claim_invite( let claimer_hex = pubkey.to_hex(); let key = invite_token::derive_invite_key(&state.relay_keypair); + if enforcing && !request.code.starts_with(V2_PREFIX) { + return Err(api_error( + StatusCode::SERVICE_UNAVAILABLE, + "invite_unavailable", + )); + } + // --- v2 database-backed path --- // // Route by exact prefix: v2. codes use the durable invite table. No @@ -448,8 +613,141 @@ pub async fn claim_invite( } let token_hash = hash_v2_code(&request.code); - let identity_binding = - crate::corporate_identity::binding_input_for_proof(&identity_proof, &pubkey); + if enforcing { + let assertion = enrollment_assertion.ok_or_else(|| { + api_error(StatusCode::FORBIDDEN, "relay identity verification failed") + })?; + let enrollment = authorize_enrollment_if_configured( + &state, + Arc::clone(&verified_proof), + assertion, + stable_correlation_from_proof(&verified_proof), + ) + .await + .map_err(|error| { + tracing::warn!(error = %error, "invite claim: protected enrollment denied"); + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })?; + let mut stable_key = Vec::with_capacity(64); + stable_key.extend_from_slice(&token_hash); + stable_key.extend_from_slice(pubkey.as_bytes()); + let operation_id = + ProtectedOperationId::derive(tenant.community(), "invite.claim.v1", &stable_key) + .map_err(|error| internal_error(&format!("invite claim: {error}")))?; + let mut digest = Sha256::new(); + digest.update(b"buzz-invite-claim-v1"); + digest.update(token_hash); + digest.update(pubkey.as_bytes()); + if let Some(policy) = &state.config.join_policy { + digest.update(policy.version.as_bytes()); + } + let request_fingerprint: [u8; 32] = digest.finalize().into(); + let permit = enrollment + .seal_postgres_enrollment(operation_id, "invite.claim.v1", request_fingerprint) + .map_err(|error| { + tracing::warn!(error = %error, "invite claim: protected enrollment stale"); + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })? + .ok_or_else(|| { + api_error(StatusCode::FORBIDDEN, "protected authorization denied") + })?; + let response = match begin_authorized_enrollment(&state, permit) + .await + .map_err(|error| internal_error(&format!("invite claim: {error}")))? + { + AuthorizedEnrollmentStart::Replay(payload) => serde_json::from_slice(&payload) + .map_err(|error| internal_error(&format!("invite claim replay: {error}")))?, + AuthorizedEnrollmentStart::Execute(mut operation) => { + let issuer = operation.issuer().to_owned(); + let subject = operation.subject().to_owned(); + let actor = *operation.actor_pubkey(); + let identity = buzz_db::identity_binding::IdentityBindingInput { + issuer: &issuer, + uid: &subject, + pubkey: &actor, + display_name: None, + source: buzz_db::identity_binding::SOURCE_JWT_NPUB, + }; + let outcome = buzz_db::relay_invite::claim_relay_invite_with_identity_tx( + operation.transaction(), + tenant.community(), + &token_hash, + &claimer_hex, + state + .config + .join_policy + .as_ref() + .map(|policy| policy.version.as_str()), + Some(&identity), + ) + .await + .map_err(|error| internal_error(&format!("invite claim: {error}")))?; + let response = match outcome { + buzz_db::relay_invite::ClaimOutcome::Joined { .. } => serde_json::json!({ + "status": "joined", + "community_id": tenant.community().to_string(), + "host": tenant.host(), + "role": "member", + }), + buzz_db::relay_invite::ClaimOutcome::AlreadyMember { .. } => { + serde_json::json!({ + "status": "already_member", + "community_id": tenant.community().to_string(), + "host": tenant.host(), + "role": "member", + }) + } + buzz_db::relay_invite::ClaimOutcome::Expired => { + return Err(api_error(StatusCode::FORBIDDEN, "invite_expired")); + } + buzz_db::relay_invite::ClaimOutcome::Exhausted => { + return Err(api_error(StatusCode::FORBIDDEN, "invite_exhausted")); + } + buzz_db::relay_invite::ClaimOutcome::Invalid => { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + } + buzz_db::relay_invite::ClaimOutcome::IdentityConflict(_) => { + return Err(api_error( + StatusCode::FORBIDDEN, + "relay identity binding conflict", + )); + } + buzz_db::relay_invite::ClaimOutcome::IdentityRevoked => { + return Err(api_error( + StatusCode::FORBIDDEN, + "relay identity binding revoked", + )); + } + buzz_db::relay_invite::ClaimOutcome::IdentityBindingRequired => { + return Err(api_error( + StatusCode::FORBIDDEN, + "relay identity binding required", + )); + } + }; + let payload = serde_json::to_vec(&response) + .map_err(|error| internal_error(&format!("invite claim: {error}")))?; + operation + .commit(&payload) + .await + .map_err(|error| internal_error(&format!("invite claim: {error}")))?; + response + } + }; + return Ok(Json(response)); + } + let legacy_identity = crate::authorization_runtime::transport::legacy_identity_lane( + &state, + tenant.community(), + ) + == crate::authorization_runtime::transport::LegacyIdentityLane::Legacy; + let identity_binding = if legacy_identity { + identity_proof.as_ref().and_then(|proof| { + crate::corporate_identity::binding_input_for_proof(proof, &pubkey) + }) + } else { + None + }; let outcome = state .db .claim_relay_invite_with_identity( @@ -470,15 +768,19 @@ pub async fn claim_invite( buzz_db::relay_invite::ClaimOutcome::Joined { identity_binding, .. } => { - crate::corporate_identity::finalize_atomic_corporate_identity_result( - &state, - tenant.community(), - pubkey, - identity_proof, - identity_binding, - ) - .await - .map_err(|error| error.into_api_error())?; + if legacy_identity { + if let Some(identity_proof) = identity_proof { + crate::corporate_identity::finalize_atomic_corporate_identity_result( + &state, + tenant.community(), + pubkey, + identity_proof, + identity_binding, + ) + .await + .map_err(|error| error.into_api_error())?; + } + } tracing::info!( community = %tenant.community(), member = %claimer_hex, @@ -503,15 +805,19 @@ pub async fn claim_invite( buzz_db::relay_invite::ClaimOutcome::AlreadyMember { identity_binding, .. } => { - crate::corporate_identity::finalize_atomic_corporate_identity_result( - &state, - tenant.community(), - pubkey, - identity_proof, - identity_binding, - ) - .await - .map_err(|error| error.into_api_error())?; + if legacy_identity { + if let Some(identity_proof) = identity_proof { + crate::corporate_identity::finalize_atomic_corporate_identity_result( + &state, + tenant.community(), + pubkey, + identity_proof, + identity_binding, + ) + .await + .map_err(|error| error.into_api_error())?; + } + } Ok(Json(serde_json::json!({ "status": "already_member", "community_id": tenant.community().to_string(), @@ -529,6 +835,9 @@ pub async fn claim_invite( Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")) } buzz_db::relay_invite::ClaimOutcome::IdentityConflict(conflict) => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -539,6 +848,9 @@ pub async fn claim_invite( .await) } buzz_db::relay_invite::ClaimOutcome::IdentityRevoked => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -549,6 +861,9 @@ pub async fn claim_invite( .await) } buzz_db::relay_invite::ClaimOutcome::IdentityBindingRequired => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -582,8 +897,16 @@ pub async fn claim_invite( .map_err(|_| api_error(StatusCode::FORBIDDEN, "join_policy_required"))?; } - let identity_binding = - crate::corporate_identity::binding_input_for_proof(&identity_proof, &pubkey); + let legacy_identity = + crate::authorization_runtime::transport::legacy_identity_lane(&state, tenant.community()) + == crate::authorization_runtime::transport::LegacyIdentityLane::Legacy; + let identity_binding = if legacy_identity { + identity_proof + .as_ref() + .and_then(|proof| crate::corporate_identity::binding_input_for_proof(proof, &pubkey)) + } else { + None + }; let claim_outcome = state .db .claim_relay_membership_with_identity( @@ -605,6 +928,9 @@ pub async fn claim_invite( identity_binding, } => (inserted, identity_binding), buzz_db::relay_members::MembershipClaimOutcome::IdentityConflict(conflict) => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; return Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -615,6 +941,9 @@ pub async fn claim_invite( .await); } buzz_db::relay_members::MembershipClaimOutcome::IdentityRevoked => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; return Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -625,6 +954,9 @@ pub async fn claim_invite( .await); } buzz_db::relay_members::MembershipClaimOutcome::IdentityBindingRequired => { + let Some(identity_proof) = identity_proof else { + return Err(api_error(StatusCode::FORBIDDEN, "invite_invalid")); + }; return Err(record_atomic_identity_rejection( &state, tenant.community(), @@ -635,15 +967,19 @@ pub async fn claim_invite( .await); } }; - crate::corporate_identity::finalize_atomic_corporate_identity_result( - &state, - tenant.community(), - pubkey, - identity_proof, - identity_binding, - ) - .await - .map_err(|error| error.into_api_error())?; + if legacy_identity { + if let Some(identity_proof) = identity_proof { + crate::corporate_identity::finalize_atomic_corporate_identity_result( + &state, + tenant.community(), + pubkey, + identity_proof, + identity_binding, + ) + .await + .map_err(|error| error.into_api_error())?; + } + } if was_inserted { tracing::info!( diff --git a/crates/buzz-relay/src/handlers/identity_archive.rs b/crates/buzz-relay/src/handlers/identity_archive.rs index 9da920483..113e9fe02 100644 --- a/crates/buzz-relay/src/handlers/identity_archive.rs +++ b/crates/buzz-relay/src/handlers/identity_archive.rs @@ -138,6 +138,61 @@ pub async fn handle_identity_archive_event( Ok(()) } +/// Validate and apply an identity archive request inside the caller's protected +/// authorization transaction. The request event is persisted by the caller in +/// that same transaction; relay-signed deltas remain unavailable background +/// effects in Enforce. +pub async fn handle_identity_archive_event_tx( + tenant: &TenantContext, + state: &Arc, + event: &Event, + transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>, +) -> Result { + let kind = event.kind.as_u16() as u32; + let actor_hex = event.pubkey.to_hex(); + if kind != KIND_IA_ARCHIVE_REQUEST && kind != KIND_IA_UNARCHIVE_REQUEST { + return Err(format!("unexpected identity archive kind: {kind}")); + } + enforce_freshness(event)?; + require_single_protected_tag(event)?; + let target_hex = extract_single_p_tag_hex(event) + .ok_or_else(|| "missing or invalid p tag".to_string())? + .to_ascii_lowercase(); + let replaced_by = extract_optional_replaced_by(event, &target_hex)?; + if kind == KIND_IA_UNARCHIVE_REQUEST && replaced_by.is_some() { + return Err("replaced-by is not valid on unarchive requests".into()); + } + let reason = extract_tag_value(event, "reason"); + let consent_path = determine_consent_path_tx( + tenant.community(), + state, + event, + &target_hex, + &actor_hex, + transaction, + ) + .await?; + let request_event_id = event.id.to_hex(); + if kind == KIND_IA_ARCHIVE_REQUEST { + buzz_db::archived_identities::archive_tx( + transaction, + tenant.community(), + &target_hex, + consent_path.as_str(), + &actor_hex, + reason.as_deref(), + replaced_by.as_deref(), + &request_event_id, + ) + .await + .map_err(|error| format!("database error: {error}")) + } else { + buzz_db::archived_identities::unarchive_tx(transaction, tenant.community(), &target_hex) + .await + .map_err(|error| format!("database error: {error}")) + } +} + fn enforce_freshness(event: &Event) -> Result<(), String> { let event_ts = event.created_at.as_secs() as i64; let now = std::time::SystemTime::now() @@ -250,6 +305,40 @@ async fn determine_consent_path( Ok(ConsentPath::Owner) } +async fn determine_consent_path_tx( + community_id: CommunityId, + state: &Arc, + event: &Event, + target_hex: &str, + actor_hex: &str, + transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>, +) -> Result { + if actor_hex == target_hex { + return Ok(ConsentPath::SelfSigned); + } + let actor_member = + buzz_db::relay_members::get_relay_member_tx(transaction, community_id, actor_hex) + .await + .map_err(|error| format!("database error: {error}"))?; + let actor_role = actor_member + .as_ref() + .map(|member| member.role.as_str()) + .unwrap_or(""); + if actor_role == "owner" || actor_role == "admin" { + return Ok(ConsentPath::Admin); + } + verify_owner_consent_tx( + community_id, + state, + event, + target_hex, + actor_hex, + transaction, + ) + .await?; + Ok(ConsentPath::Owner) +} + async fn verify_owner_consent( community_id: CommunityId, state: &Arc, @@ -297,6 +386,76 @@ async fn verify_owner_consent( Ok(()) } +async fn verify_owner_consent_tx( + community_id: CommunityId, + _state: &Arc, + event: &Event, + target_hex: &str, + actor_hex: &str, + transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>, +) -> Result<(), String> { + let request_auth = extract_single_auth_tag_json(event)?; + let request_owner = verify_auth_tag_owner(&request_auth, target_hex) + .map_err(|error| format!("invalid request auth tag: {error}"))?; + if request_owner != actor_hex { + return Err("request auth owner must equal request signer".into()); + } + enforce_request_auth_time_bounds(&request_auth, event.created_at.as_secs())?; + + let target_pubkey = PublicKey::from_hex(target_hex) + .map_err(|error| format!("invalid target pubkey: {error}"))?; + let target_author = target_pubkey.to_bytes().to_vec(); + + // The user-row share lock is also taken by the Enforce profile projection. + // It serializes archive consent against a concurrent kind:0 replacement so + // the profile read below remains the consent state through commit. + let target_exists = sqlx::query_scalar::<_, i32>( + "SELECT 1 FROM users \ + WHERE community_id = $1 AND pubkey = $2 FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(&target_author) + .fetch_optional(&mut **transaction) + .await + .map_err(|error| format!("database error: {error}"))? + .is_some(); + if !target_exists { + return Err("target has no live user profile".into()); + } + + let profile = buzz_db::event::query_events_tx( + transaction, + &EventQuery { + kinds: Some(vec![KIND_PROFILE as i32]), + authors: Some(vec![target_author]), + limit: Some(1), + global_only: true, + ..EventQuery::for_community(community_id) + }, + ) + .await + .map_err(|error| format!("database error: {error}"))? + .into_iter() + .next() + .ok_or_else(|| "target has no live kind:0 profile".to_string())?; + if !buzz_db::event::lock_live_event_tx(transaction, community_id, profile.event.id.as_bytes()) + .await + .map_err(|error| format!("database error: {error}"))? + { + return Err("live kind:0 changed during authorization".into()); + } + if profile.event.pubkey.to_hex() != target_hex { + return Err("live kind:0 author did not match target".into()); + } + let live_auth = extract_single_auth_tag_json(&profile.event)?; + let live_owner = verify_auth_tag_owner(&live_auth, target_hex) + .map_err(|error| format!("invalid live kind:0 auth tag: {error}"))?; + if live_owner != actor_hex { + return Err("live kind:0 no longer attests to request signer".into()); + } + Ok(()) +} + fn extract_single_auth_tag_json(event: &Event) -> Result { let mut found: Option> = None; for tag in event.tags.iter() {