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>
This commit is contained in:
Cea Stapleton Cordasco
2026-08-04 11:43:41 -05:00
parent 19524a34ff
commit df4efac2f5
10 changed files with 3208 additions and 247 deletions
+37 -2
View File
@@ -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<BindingExpiry>,
source: BindingSource,
resolution_reason: AuthorizationReason,
) -> Result<Self, AuthContextError> {
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
}
}
+12
View File
@@ -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",
}
}
}
File diff suppressed because it is too large Load Diff
+46 -1
View File
@@ -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<bool> {
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<bool> {
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,
File diff suppressed because it is too large Load Diff
+100 -23
View File
@@ -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<i32>,
) -> Result<MintedInvite> {
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<i32>,
) -> Result<MintedInvite> {
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<String> = 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<uuid::Uuid>,
@@ -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<Utc>) -> Result<u64> {
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<Utc>,
excluded_communities: &[Uuid],
) -> Result<u64> {
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<ClaimOutcome> {
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<ClaimOutcome> {
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
+131 -1
View File
@@ -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<Option<RelayMember>> {
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<RelayMember, sqlx::Error> {
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<Vec<RelayMember>> {
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<bool> {
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<RemoveResult> {
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<RemoveResult> {
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<bool> {
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
+107
View File
@@ -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<bool> {
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<bool> {
let owner = sqlx::query_scalar::<_, Vec<u8>>(
"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::*;
+412 -76
View File
@@ -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<Arc<AppState>>) -> Json<Value> {
match &state.config.join_policy {
@@ -236,7 +255,9 @@ async fn authenticate(
(
buzz_core::TenantContext,
nostr::PublicKey,
crate::corporate_identity::CorporateIdentityProof,
Option<crate::corporate_identity::CorporateIdentityProof>,
Arc<buzz_auth::VerifiedNostrProof>,
Option<Arc<buzz_auth::VerifiedFederatedAssertion>>,
),
(StatusCode, Json<Value>),
> {
@@ -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<Json<Value>, (StatusCode, Json<Value>)> {
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<Json<Value>, (StatusCode, Json<Value>)> {
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!(
@@ -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<AppState>,
event: &Event,
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
) -> Result<bool, String> {
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<AppState>,
event: &Event,
target_hex: &str,
actor_hex: &str,
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
) -> Result<ConsentPath, String> {
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<AppState>,
@@ -297,6 +386,76 @@ async fn verify_owner_consent(
Ok(())
}
async fn verify_owner_consent_tx(
community_id: CommunityId,
_state: &Arc<AppState>,
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<String, String> {
let mut found: Option<Vec<String>> = None;
for tag in event.tags.iter() {