diff --git a/crates/buzz-core/src/client_binding_status.rs b/crates/buzz-core/src/client_binding_status.rs new file mode 100644 index 000000000..6428938d2 --- /dev/null +++ b/crates/buzz-core/src/client_binding_status.rs @@ -0,0 +1,1289 @@ +//! Relay-authenticated client binding status. +//! +//! Kind `24244` is a short-lived, ephemeral envelope whose JSON content names +//! the exact authorization domain and event-author key to which presentation +//! applies. The status is display-only: it is not identity proof, membership, +//! an authorization decision, or an access lease. Consumers must obtain a +//! value through +//! [`validate_client_binding_status_event`](crate::client_binding_status::validate_client_binding_status_event) +//! or [`ClientBindingStatusTracker`](crate::client_binding_status::ClientBindingStatusTracker) +//! rather than mutable profile fields or client-supplied claims. + +use std::fmt; + +use nostr::{Event, EventBuilder, EventId, Keys, Kind, PublicKey, Timestamp}; +use serde::{Deserialize, Serialize}; +use thiserror::Error; +use uuid::Uuid; + +use crate::{kind::KIND_CLIENT_BINDING_STATUS, verify_event, CommunityId}; + +/// Wire version accepted by this module. +pub const CLIENT_BINDING_STATUS_VERSION: u64 = 1; + +/// Maximum lifetime of a client binding status, in seconds. +/// +/// Producers may choose a shorter lifetime. A longer lifetime fails closed. +pub const MAX_CLIENT_BINDING_STATUS_LIFETIME_SECS: u64 = 300; + +/// Explicit client-status clock-skew allowance. +/// +/// Status is presentation-only and issued from centrally injected relay time, +/// so the portable profile permits no future issue-time skew. +pub const CLIENT_BINDING_STATUS_CLOCK_SKEW_SECS: u64 = 0; + +/// Maximum encoded payload length. +pub const MAX_CLIENT_BINDING_STATUS_PAYLOAD_BYTES: usize = 4096; + +/// Maximum encoded length of the opaque policy revision. +pub const MAX_CLIENT_BINDING_STATUS_POLICY_VERSION_BYTES: usize = 256; + +/// Maximum encoded length of the optional privacy-approved display label. +pub const MAX_CLIENT_BINDING_STATUS_LABEL_BYTES: usize = 80; + +/// Server-selected, display-only status disposition. +/// +/// V1 deliberately exposes only current verification or an opaque withdrawal. +/// Revocation, rotation, lineage, retirement, and other lifecycle causes are +/// durable server-side history and are never part of the client contract. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ClientBindingStatusDisposition { + /// The client may display current verification while the envelope is fresh. + DisplayCurrent, + /// Clear current presentation and advance only the scoped replay floor. + Withdrawn, +} + +/// A validated v1 client binding status. +/// +/// Construction proves that one correctly signed event from the expected +/// relay matched the caller's server-resolved authorization domain and exact +/// message-author key at the injected validation time. It does not grant or +/// deny any capability. +#[derive(Clone, PartialEq, Eq)] +pub struct ClientBindingStatusV1 { + event_id: EventId, + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + binding_version: Option, + policy_version: Option, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + disposition: ClientBindingStatusDisposition, + display_label: Option, +} + +impl fmt::Debug for ClientBindingStatusV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ClientBindingStatusV1") + .field("event_id", &"[redacted]") + .field("authorization_domain", &"[redacted]") + .field("event_author_pubkey", &"[redacted]") + .field("binding_version", &"[redacted]") + .field("policy_version", &"[redacted]") + .field("status_revision", &"[redacted]") + .field("issued_at", &"[redacted]") + .field("fresh_until", &"[redacted]") + .field("disposition", &self.disposition) + .field( + "display_label", + &self.display_label.as_ref().map(|_| "[redacted]"), + ) + .finish() + } +} + +impl ClientBindingStatusV1 { + /// Signed event identifier used for equal-revision idempotency. + pub const fn event_id(&self) -> EventId { + self.event_id + } + + /// Server-resolved authorization domain for which this status is valid. + pub const fn authorization_domain(&self) -> CommunityId { + self.authorization_domain + } + + /// Exact event-author key whose messages may consume this status. + pub const fn event_author_pubkey(&self) -> PublicKey { + self.event_author_pubkey + } + + /// Positive version of the identity-to-key binding. + pub const fn binding_version(&self) -> Option { + self.binding_version + } + + /// Opaque provider-neutral policy revision. + pub fn policy_version(&self) -> Option<&str> { + self.policy_version.as_deref() + } + + /// Positive, monotonically increasing revision for this scoped status. + pub const fn status_revision(&self) -> u64 { + self.status_revision + } + + /// Relay issue time as Unix seconds. + pub const fn issued_at(&self) -> u64 { + self.issued_at + } + + /// Exclusive freshness bound as Unix seconds. + pub const fn fresh_until(&self) -> u64 { + self.fresh_until + } + + /// Server-selected, display-only disposition. + pub const fn disposition(&self) -> ClientBindingStatusDisposition { + self.disposition + } + + /// Optional privacy-approved label for current presentation. + pub fn display_label(&self) -> Option<&str> { + self.display_label.as_deref() + } + + /// Returns `true` only for an active, current presentation disposition. + /// + /// Freshness and relay authentication have already been checked by + /// [`validate_client_binding_status_event`]. This method must not be used + /// for access control. + pub const fn displays_current_binding(&self) -> bool { + matches!( + self.disposition, + ClientBindingStatusDisposition::DisplayCurrent + ) + } +} + +/// Validated producer input for one v1 client binding status. +/// +/// This is a serialization/signing input only. It intentionally carries no +/// authorization context, capability set, membership state, or access lease. +pub struct ClientBindingStatusInputV1 { + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + binding_version: Option, + policy_version: Option, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + disposition: ClientBindingStatusDisposition, + display_label: Option, +} + +impl fmt::Debug for ClientBindingStatusInputV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ClientBindingStatusInputV1") + .field("authorization_domain", &"[redacted]") + .field("event_author_pubkey", &"[redacted]") + .field("binding_version", &"[redacted]") + .field("policy_version", &"[redacted]") + .field("status_revision", &"[redacted]") + .field("issued_at", &"[redacted]") + .field("fresh_until", &"[redacted]") + .field("disposition", &self.disposition) + .field( + "display_label", + &self.display_label.as_ref().map(|_| "[redacted]"), + ) + .finish() + } +} + +impl ClientBindingStatusInputV1 { + /// Construct bounded, provider-neutral current-status input. + #[allow(clippy::too_many_arguments)] + pub fn current( + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + binding_version: u64, + policy_version: impl Into, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + display_label: Option, + ) -> Result { + let value = Self { + authorization_domain, + event_author_pubkey, + binding_version: Some(binding_version), + policy_version: Some(policy_version.into()), + status_revision, + issued_at, + fresh_until, + disposition: ClientBindingStatusDisposition::DisplayCurrent, + display_label, + }; + validate_payload_fields( + value.authorization_domain, + value.binding_version, + value.policy_version.as_deref(), + value.status_revision, + value.issued_at, + value.fresh_until, + value.disposition, + value.display_label.as_deref(), + )?; + Ok(value) + } + + /// Construct an opaque withdrawal carrying no binding or lifecycle data. + pub fn withdrawn( + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + ) -> Result { + let value = Self { + authorization_domain, + event_author_pubkey, + binding_version: None, + policy_version: None, + status_revision, + issued_at, + fresh_until, + disposition: ClientBindingStatusDisposition::Withdrawn, + display_label: None, + }; + validate_payload_fields( + value.authorization_domain, + value.binding_version, + value.policy_version.as_deref(), + value.status_revision, + value.issued_at, + value.fresh_until, + value.disposition, + value.display_label.as_deref(), + )?; + Ok(value) + } + + /// Sign this status with the relay key advertised through NIP-11 `self`. + pub fn sign_with_relay_keys( + self, + relay_keys: &Keys, + ) -> Result { + let wire = WireClientBindingStatusV1 { + version: CLIENT_BINDING_STATUS_VERSION, + authorization_domain: self.authorization_domain.as_uuid().to_string(), + event_author_pubkey: self.event_author_pubkey.to_hex(), + status_revision: self.status_revision, + issued_at: self.issued_at, + fresh_until: self.fresh_until, + status: self.disposition, + binding_version: self.binding_version, + policy_version: self.policy_version, + display_label: self.display_label, + }; + let content = serde_json::to_string(&wire) + .map_err(|_| ClientBindingStatusBuildError::Serialization)?; + if content.len() > MAX_CLIENT_BINDING_STATUS_PAYLOAD_BYTES { + return Err(ClientBindingStatusBuildError::PayloadTooLarge); + } + EventBuilder::new(Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), content) + .tags([]) + .custom_created_at(Timestamp::from(wire.issued_at)) + .sign_with_keys(relay_keys) + .map_err(|_| ClientBindingStatusBuildError::Signing) + } +} + +#[derive(Deserialize)] +struct VersionHeader { + version: u64, +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct WireClientBindingStatusV1 { + version: u64, + authorization_domain: String, + event_author_pubkey: String, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + status: ClientBindingStatusDisposition, + #[serde(default, skip_serializing_if = "Option::is_none")] + binding_version: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + policy_version: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + display_label: Option, +} + +/// Validate and authenticate a relay-issued v1 client binding status event. +/// +/// `trusted_relay_pubkey`, `expected_authorization_domain`, and +/// `expected_event_author_pubkey` must come from connection or message context, +/// never from the status payload. `now` is injected Unix time; the exclusive +/// expiry rule means `now == fresh_until` is expired. Future-issued statuses +/// are rejected under [`CLIENT_BINDING_STATUS_CLOCK_SKEW_SECS`]. +pub fn validate_client_binding_status_event( + event: &Event, + trusted_relay_pubkey: &PublicKey, + expected_authorization_domain: CommunityId, + expected_event_author_pubkey: &PublicKey, + now: u64, +) -> Result { + if event.kind.as_u16() as u32 != KIND_CLIENT_BINDING_STATUS { + return Err(ClientBindingStatusError::WrongKind); + } + if event.content.len() > MAX_CLIENT_BINDING_STATUS_PAYLOAD_BYTES { + return Err(ClientBindingStatusError::PayloadTooLarge); + } + verify_event(event).map_err(|_| ClientBindingStatusError::UnauthenticatedEvent)?; + if event.pubkey != *trusted_relay_pubkey { + return Err(ClientBindingStatusError::UnexpectedRelay); + } + if !event.tags.is_empty() { + return Err(ClientBindingStatusError::UnexpectedTags); + } + + let header: VersionHeader = serde_json::from_str(&event.content) + .map_err(|_| ClientBindingStatusError::MalformedPayload)?; + if header.version != CLIENT_BINDING_STATUS_VERSION { + return Err(ClientBindingStatusError::UnsupportedVersion); + } + + let wire: WireClientBindingStatusV1 = serde_json::from_str(&event.content) + .map_err(|_| ClientBindingStatusError::MalformedPayload)?; + + let authorization_domain = Uuid::parse_str(&wire.authorization_domain) + .map_err(|_| ClientBindingStatusError::InvalidAuthorizationDomain)?; + if authorization_domain.is_nil() + || authorization_domain.to_string() != wire.authorization_domain + { + return Err(ClientBindingStatusError::InvalidAuthorizationDomain); + } + let authorization_domain = CommunityId::from_uuid(authorization_domain); + if authorization_domain != expected_authorization_domain { + return Err(ClientBindingStatusError::AuthorizationDomainMismatch); + } + + if wire.event_author_pubkey.len() != 64 + || !wire + .event_author_pubkey + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + return Err(ClientBindingStatusError::InvalidEventAuthorPubkey); + } + let event_author_pubkey = PublicKey::from_hex(&wire.event_author_pubkey) + .map_err(|_| ClientBindingStatusError::InvalidEventAuthorPubkey)?; + if event_author_pubkey != *expected_event_author_pubkey { + return Err(ClientBindingStatusError::EventAuthorMismatch); + } + + let disposition = wire.status; + let binding_version = wire.binding_version; + let policy_version = wire.policy_version; + let display_label = wire.display_label; + + validate_payload_fields( + authorization_domain, + binding_version, + policy_version.as_deref(), + wire.status_revision, + wire.issued_at, + wire.fresh_until, + disposition, + display_label.as_deref(), + )?; + if event.created_at.as_secs() != wire.issued_at { + return Err(ClientBindingStatusError::EventTimeMismatch); + } + if wire.issued_at > now.saturating_add(CLIENT_BINDING_STATUS_CLOCK_SKEW_SECS) { + return Err(ClientBindingStatusError::NotYetValid); + } + if now >= wire.fresh_until { + return Err(ClientBindingStatusError::Expired); + } + + Ok(ClientBindingStatusV1 { + event_id: event.id, + authorization_domain, + event_author_pubkey, + binding_version, + policy_version, + status_revision: wire.status_revision, + issued_at: wire.issued_at, + fresh_until: wire.fresh_until, + disposition, + display_label, + }) +} + +#[allow(clippy::too_many_arguments)] +fn validate_payload_fields( + authorization_domain: CommunityId, + binding_version: Option, + policy_version: Option<&str>, + status_revision: u64, + issued_at: u64, + fresh_until: u64, + disposition: ClientBindingStatusDisposition, + display_label: Option<&str>, +) -> Result<(), ClientBindingStatusError> { + if authorization_domain.as_uuid().is_nil() { + return Err(ClientBindingStatusError::InvalidAuthorizationDomain); + } + match disposition { + ClientBindingStatusDisposition::DisplayCurrent => { + if binding_version.is_none_or(|version| version == 0) { + return Err(ClientBindingStatusError::InvalidBindingVersion); + } + let Some(policy_version) = policy_version else { + return Err(ClientBindingStatusError::InvalidPolicyVersion); + }; + if policy_version.is_empty() + || policy_version.len() > MAX_CLIENT_BINDING_STATUS_POLICY_VERSION_BYTES + || policy_version.trim() != policy_version + || policy_version.chars().any(char::is_control) + { + return Err(ClientBindingStatusError::InvalidPolicyVersion); + } + } + ClientBindingStatusDisposition::Withdrawn => { + if binding_version.is_some() || policy_version.is_some() || display_label.is_some() { + return Err(ClientBindingStatusError::WithdrawalContainsCurrentState); + } + } + } + if status_revision == 0 { + return Err(ClientBindingStatusError::InvalidStatusRevision); + } + if issued_at == 0 { + return Err(ClientBindingStatusError::InvalidIssueTime); + } + if fresh_until <= issued_at { + return Err(ClientBindingStatusError::InvalidFreshnessBound); + } + if fresh_until - issued_at > MAX_CLIENT_BINDING_STATUS_LIFETIME_SECS { + return Err(ClientBindingStatusError::FreshnessWindowTooLong); + } + validate_display_label(disposition, display_label) +} + +fn validate_display_label( + disposition: ClientBindingStatusDisposition, + display_label: Option<&str>, +) -> Result<(), ClientBindingStatusError> { + let Some(label) = display_label else { + return Ok(()); + }; + if disposition != ClientBindingStatusDisposition::DisplayCurrent + || label.is_empty() + || label.len() > MAX_CLIENT_BINDING_STATUS_LABEL_BYTES + || label.trim() != label + || label.chars().any(char::is_control) + { + return Err(ClientBindingStatusError::InvalidDisplayLabel); + } + Ok(()) +} + +/// One accepted high-water update from [`ClientBindingStatusTracker`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ClientBindingStatusUpdate { + /// A strictly newer revision replaced presentation state. + Accepted, + /// The exact same signed event was observed again. + Duplicate, +} + +#[derive(Clone, Copy)] +struct StatusHighWater { + revision: u64, + event_id: EventId, +} + +/// Client-side, scope-keyed status revision fold. +/// +/// The tracker authenticates every event against one trusted relay/domain/ +/// author tuple. A lower revision or a different event at the same revision is +/// rejected. Expiry and disconnect clear presentation while retaining the +/// high-water mark, so a previously seen envelope cannot restore a badge. +pub struct ClientBindingStatusTracker { + trusted_relay_pubkey: PublicKey, + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + high_water: Option, + status: Option, +} + +impl fmt::Debug for ClientBindingStatusTracker { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ClientBindingStatusTracker") + .field("trusted_relay_pubkey", &"[redacted]") + .field("authorization_domain", &"[redacted]") + .field("event_author_pubkey", &"[redacted]") + .field("high_water", &self.high_water.map(|_| "[redacted]")) + .field("status", &self.status.as_ref().map(|_| "[redacted]")) + .finish() + } +} + +impl ClientBindingStatusTracker { + /// Start an empty fold for one trusted relay/domain/author scope. + pub const fn new( + trusted_relay_pubkey: PublicKey, + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + ) -> Self { + Self { + trusted_relay_pubkey, + authorization_domain, + event_author_pubkey, + high_water: None, + status: None, + } + } + + /// Authenticate and fold one signed status event at injected time `now`. + pub fn accept( + &mut self, + event: &Event, + now: u64, + ) -> Result { + let status = validate_client_binding_status_event( + event, + &self.trusted_relay_pubkey, + self.authorization_domain, + &self.event_author_pubkey, + now, + )?; + if let Some(high_water) = self.high_water { + if status.status_revision < high_water.revision { + return Err(ClientBindingStatusFoldError::LowerRevisionReplay); + } + if status.status_revision == high_water.revision { + if status.event_id != high_water.event_id { + return Err(ClientBindingStatusFoldError::ConflictingEqualRevision); + } + return Ok(ClientBindingStatusUpdate::Duplicate); + } + } + self.high_water = Some(StatusHighWater { + revision: status.status_revision, + event_id: status.event_id, + }); + self.status = status.displays_current_binding().then_some(status); + Ok(ClientBindingStatusUpdate::Accepted) + } + + /// Return the fresh accepted status, clearing expired presentation. + pub fn status(&mut self, now: u64) -> Option<&ClientBindingStatusV1> { + if self + .status + .as_ref() + .is_some_and(|status| now >= status.fresh_until) + { + self.status = None; + } + self.status.as_ref() + } + + /// Return only a fresh status allowed to display current verification. + pub fn current_presentation(&mut self, now: u64) -> Option<&ClientBindingStatusV1> { + self.status(now) + .filter(|status| status.displays_current_binding()) + } + + /// Clear presentation on relay disconnect while retaining replay defense. + pub fn on_disconnect(&mut self) { + self.status = None; + } + + /// Replace the trusted scope and clear both presentation and revision state. + /// + /// Call this on relay-identity, authorization-domain, or event-author + /// changes. Evidence from the old scope is never carried into the new one. + pub fn change_scope( + &mut self, + trusted_relay_pubkey: PublicKey, + authorization_domain: CommunityId, + event_author_pubkey: PublicKey, + ) { + self.trusted_relay_pubkey = trusted_relay_pubkey; + self.authorization_domain = authorization_domain; + self.event_author_pubkey = event_author_pubkey; + self.high_water = None; + self.status = None; + } + + /// Highest revision accepted for the current scope. + pub const fn high_water_revision(&self) -> Option { + match self.high_water { + Some(value) => Some(value.revision), + None => None, + } + } +} + +/// Fail-closed client binding status validation error. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +#[non_exhaustive] +pub enum ClientBindingStatusError { + /// The event did not use the dedicated ephemeral status kind. + #[error("client binding status event has the wrong kind")] + WrongKind, + /// The event body exceeded the public bound. + #[error("client binding status payload is too large")] + PayloadTooLarge, + /// The event ID or Schnorr signature was invalid. + #[error("client binding status event is not authenticated")] + UnauthenticatedEvent, + /// The signer did not match the relay key established by the connection. + #[error("client binding status signer is not the trusted relay")] + UnexpectedRelay, + /// Status events must not carry tags or private indexing material. + #[error("client binding status event contains unexpected tags")] + UnexpectedTags, + /// The JSON shape, field set, or enum encoding was malformed. + #[error("client binding status payload is malformed")] + MalformedPayload, + /// The payload used a version this client does not understand. + #[error("client binding status version is unsupported")] + UnsupportedVersion, + /// The authorization domain was nil or not a canonical UUID. + #[error("client binding status authorization domain is invalid")] + InvalidAuthorizationDomain, + /// The payload did not match the server-resolved authorization domain. + #[error("client binding status authorization domain does not match")] + AuthorizationDomainMismatch, + /// The event-author key was not canonical lowercase 64-character hex. + #[error("client binding status event-author key is invalid")] + InvalidEventAuthorPubkey, + /// The payload named a key other than the displayed event's author. + #[error("client binding status event-author key does not match")] + EventAuthorMismatch, + /// The binding version was zero. + #[error("client binding status binding version must be positive")] + InvalidBindingVersion, + /// The opaque policy revision was empty, unsafe, or exceeded its bound. + #[error("client binding status policy version is invalid")] + InvalidPolicyVersion, + /// A generic withdrawal attempted to carry current binding state. + #[error("client binding status withdrawal contains current binding state")] + WithdrawalContainsCurrentState, + /// The status revision was zero. + #[error("client binding status revision must be positive")] + InvalidStatusRevision, + /// The issue time was zero. + #[error("client binding status issue time must be positive")] + InvalidIssueTime, + /// Freshness did not strictly follow issue time. + #[error("client binding status freshness bound is invalid")] + InvalidFreshnessBound, + /// The freshness window exceeded the public short-lived maximum. + #[error("client binding status freshness window is too long")] + FreshnessWindowTooLong, + /// The signed Nostr timestamp did not equal the payload issue time. + #[error("client binding status event time does not match its issue time")] + EventTimeMismatch, + /// The payload was issued after the explicit skew allowance. + #[error("client binding status is not yet valid")] + NotYetValid, + /// The exclusive freshness bound was reached. + #[error("client binding status has expired")] + Expired, + /// The optional display label was unsafe, out of bounds, or attached to a + /// non-current disposition. + #[error("client binding status display label is invalid")] + InvalidDisplayLabel, +} + +/// Status-event serialization/signing failure. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum ClientBindingStatusBuildError { + /// JSON serialization failed. + #[error("client binding status serialization failed")] + Serialization, + /// The serialized payload exceeded its public bound. + #[error("client binding status payload is too large")] + PayloadTooLarge, + /// Nostr event signing failed. + #[error("client binding status signing failed")] + Signing, +} + +/// Status revision-fold failure. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum ClientBindingStatusFoldError { + /// Cryptographic or wire validation failed. + #[error(transparent)] + InvalidStatus(#[from] ClientBindingStatusError), + /// A lower status revision attempted to restore older presentation. + #[error("client binding status revision is below the accepted high-water mark")] + LowerRevisionReplay, + /// Another signed event reused an accepted revision. + #[error("client binding status revision conflicts with another event")] + ConflictingEqualRevision, +} + +#[cfg(test)] +mod tests { + use super::*; + use nostr::JsonUtil; + use serde_json::{json, Value}; + + const DOMAIN: &str = "00000000-0000-4000-8000-000000000123"; + const ISSUED_AT: u64 = 1_800_000_000; + const FRESH_UNTIL: u64 = ISSUED_AT + 120; + + fn domain() -> CommunityId { + CommunityId::from_uuid(Uuid::parse_str(DOMAIN).expect("synthetic domain is valid")) + } + + fn input( + author: PublicKey, + revision: u64, + disposition: ClientBindingStatusDisposition, + ) -> ClientBindingStatusInputV1 { + match disposition { + ClientBindingStatusDisposition::DisplayCurrent => ClientBindingStatusInputV1::current( + domain(), + author, + 7, + "synthetic-policy-v1", + revision, + ISSUED_AT, + FRESH_UNTIL, + Some("Synthetic Example".to_string()), + ), + ClientBindingStatusDisposition::Withdrawn => ClientBindingStatusInputV1::withdrawn( + domain(), + author, + revision, + ISSUED_AT, + FRESH_UNTIL, + ), + } + .expect("synthetic status input is valid") + } + + fn signed_status( + relay: &Keys, + author: PublicKey, + revision: u64, + disposition: ClientBindingStatusDisposition, + ) -> Event { + input(author, revision, disposition) + .sign_with_relay_keys(relay) + .expect("synthetic event signs") + } + + fn validate( + event: &Event, + relay: &Keys, + author: &Keys, + now: u64, + ) -> Result { + validate_client_binding_status_event( + event, + &relay.public_key(), + domain(), + &author.public_key(), + now, + ) + } + + #[test] + fn validates_relay_authenticated_current_status() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + + let status = validate(&event, &relay, &author, ISSUED_AT) + .expect("synthetic current status validates"); + + assert_eq!(status.event_id(), event.id); + assert_eq!(status.authorization_domain(), domain()); + assert_eq!(status.event_author_pubkey(), author.public_key()); + assert_eq!(status.binding_version(), Some(7)); + assert_eq!(status.policy_version(), Some("synthetic-policy-v1")); + assert_eq!(status.status_revision(), 11); + assert_eq!(status.issued_at(), ISSUED_AT); + assert_eq!(status.fresh_until(), FRESH_UNTIL); + assert_eq!( + status.disposition(), + ClientBindingStatusDisposition::DisplayCurrent + ); + assert_eq!(status.display_label(), Some("Synthetic Example")); + assert!(status.displays_current_binding()); + assert!(event.tags.is_empty()); + } + + #[test] + fn wire_is_current_or_opaque_withdrawal() { + let relay = Keys::generate(); + let author = Keys::generate(); + let cases = [ + ( + ClientBindingStatusDisposition::DisplayCurrent, + "display_current", + ), + (ClientBindingStatusDisposition::Withdrawn, "withdrawn"), + ]; + for (disposition, expected) in cases { + let event = signed_status(&relay, author.public_key(), 11, disposition); + let content: Value = serde_json::from_str(&event.content).expect("content parses"); + assert_eq!(content["status"], expected); + if disposition == ClientBindingStatusDisposition::Withdrawn { + for forbidden in [ + "binding_version", + "policy_version", + "display_label", + "reason", + ] { + assert!( + content.get(forbidden).is_none(), + "unexpected field {forbidden}" + ); + } + } + } + } + + #[test] + fn withdrawal_removes_current_presentation_without_lifecycle_data() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::Withdrawn, + ); + let status = + validate(&event, &relay, &author, ISSUED_AT).expect("synthetic withdrawal validates"); + assert!(!status.displays_current_binding()); + assert_eq!(status.binding_version(), None); + assert_eq!(status.policy_version(), None); + assert!(status.display_label().is_none()); + } + + #[test] + fn exact_expiry_and_future_issue_fail_closed() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + + assert_eq!( + validate(&event, &relay, &author, FRESH_UNTIL), + Err(ClientBindingStatusError::Expired) + ); + assert_eq!( + validate(&event, &relay, &author, ISSUED_AT - 1), + Err(ClientBindingStatusError::NotYetValid) + ); + assert_eq!(CLIENT_BINDING_STATUS_CLOCK_SKEW_SECS, 0); + } + + #[test] + fn rejects_wrong_signer_tampering_kind_scope_and_tags() { + let relay = Keys::generate(); + let wrong_relay = Keys::generate(); + let author = Keys::generate(); + let other_author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + + assert_eq!( + validate(&event, &wrong_relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::UnexpectedRelay) + ); + assert_eq!( + validate(&event, &relay, &other_author, ISSUED_AT), + Err(ClientBindingStatusError::EventAuthorMismatch) + ); + assert_eq!( + validate_client_binding_status_event( + &event, + &relay.public_key(), + CommunityId::from_uuid(Uuid::from_u128(999)), + &author.public_key(), + ISSUED_AT, + ), + Err(ClientBindingStatusError::AuthorizationDomainMismatch) + ); + + let wrong_kind = EventBuilder::new(Kind::TextNote, event.content.clone()) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&wrong_kind, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::WrongKind) + ); + + let tagged = EventBuilder::new( + Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), + event.content.clone(), + ) + .tags([nostr::Tag::parse(["p", &author.public_key().to_hex()]) + .expect("synthetic tag is valid")]) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&tagged, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::UnexpectedTags) + ); + + let mut json: Value = serde_json::from_str(&event.as_json()).expect("event parses"); + json["content"] = Value::String("{}".to_string()); + let tampered = Event::from_json(json.to_string()).expect("tampered event parses"); + assert_eq!( + validate(&tampered, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::UnauthenticatedEvent) + ); + } + + #[test] + fn unknown_version_fields_and_enums_fail_closed() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + let mut payload: Value = serde_json::from_str(&event.content).expect("content parses"); + + payload["version"] = json!(2); + let unknown_version = EventBuilder::new( + Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), + payload.to_string(), + ) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&unknown_version, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::UnsupportedVersion) + ); + + payload["version"] = json!(1); + payload["synthetic_extension"] = json!(true); + let unknown_field = EventBuilder::new( + Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), + payload.to_string(), + ) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&unknown_field, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::MalformedPayload) + ); + } + + #[test] + fn historical_status_and_lifecycle_fields_cannot_reappear() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + let mut payload: Value = serde_json::from_str(&event.content).expect("content parses"); + let current_keys = payload + .as_object() + .expect("current payload is an object") + .keys() + .map(String::as_str) + .collect::>(); + assert_eq!( + current_keys, + std::collections::BTreeSet::from([ + "authorization_domain", + "binding_version", + "display_label", + "event_author_pubkey", + "fresh_until", + "issued_at", + "policy_version", + "status", + "status_revision", + "version", + ]) + ); + let withdrawn = signed_status( + &relay, + author.public_key(), + 12, + ClientBindingStatusDisposition::Withdrawn, + ); + let withdrawn_payload: Value = + serde_json::from_str(&withdrawn.content).expect("withdrawn content parses"); + let withdrawn_keys = withdrawn_payload + .as_object() + .expect("withdrawn payload is an object") + .keys() + .map(String::as_str) + .collect::>(); + assert_eq!( + withdrawn_keys, + std::collections::BTreeSet::from([ + "authorization_domain", + "event_author_pubkey", + "fresh_until", + "issued_at", + "status", + "status_revision", + "version", + ]) + ); + payload["status"] = json!("historical_only"); + payload["reason"] = json!("rotated"); + let historical = EventBuilder::new( + Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), + payload.to_string(), + ) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&historical, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::MalformedPayload) + ); + + for forbidden in [ + "history", + "lineage", + "predecessor", + "replacement", + "tombstone", + "principal", + "issuer", + "subject", + "retirement_reason", + "corporate_history", + "historical_label", + "employment_history", + ] { + let withdrawal = signed_status( + &relay, + author.public_key(), + 12, + ClientBindingStatusDisposition::Withdrawn, + ); + let mut payload: Value = + serde_json::from_str(&withdrawal.content).expect("content parses"); + payload[forbidden] = json!("forbidden"); + let injected = EventBuilder::new( + Kind::Custom(KIND_CLIENT_BINDING_STATUS as u16), + payload.to_string(), + ) + .custom_created_at(Timestamp::from(ISSUED_AT)) + .sign_with_keys(&relay) + .expect("synthetic event signs"); + assert_eq!( + validate(&injected, &relay, &author, ISSUED_AT), + Err(ClientBindingStatusError::MalformedPayload), + "accepted forbidden field {forbidden}" + ); + } + } + + #[test] + fn label_is_privacy_bounded_and_current_only() { + let author = Keys::generate(); + for label in [ + "", + " synthetic.example", + "synthetic.example\n", + &"x".repeat(MAX_CLIENT_BINDING_STATUS_LABEL_BYTES + 1), + ] { + assert!(matches!( + ClientBindingStatusInputV1::current( + domain(), + author.public_key(), + 7, + "synthetic-policy-v1", + 11, + ISSUED_AT, + FRESH_UNTIL, + Some(label.to_string()), + ), + Err(ClientBindingStatusError::InvalidDisplayLabel) + )); + } + let withdrawal = ClientBindingStatusInputV1::withdrawn( + domain(), + author.public_key(), + 11, + ISSUED_AT, + FRESH_UNTIL, + ) + .expect("withdrawal needs no label"); + assert!(!format!("{withdrawal:?}").contains("Synthetic Example")); + } + + #[test] + fn revision_fold_rejects_lower_and_conflicting_equal_replays() { + let relay = Keys::generate(); + let author = Keys::generate(); + let current = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + let withdrawn = signed_status( + &relay, + author.public_key(), + 12, + ClientBindingStatusDisposition::Withdrawn, + ); + let conflicting_equal = ClientBindingStatusInputV1::current( + domain(), + author.public_key(), + 8, + "synthetic-policy-v2", + 12, + ISSUED_AT, + FRESH_UNTIL, + None, + ) + .expect("conflicting current input") + .sign_with_relay_keys(&relay) + .expect("conflicting current signs"); + let mut tracker = + ClientBindingStatusTracker::new(relay.public_key(), domain(), author.public_key()); + + assert_eq!( + tracker.accept(¤t, ISSUED_AT), + Ok(ClientBindingStatusUpdate::Accepted) + ); + assert!(tracker.current_presentation(ISSUED_AT).is_some()); + assert_eq!( + tracker.accept(¤t, ISSUED_AT), + Ok(ClientBindingStatusUpdate::Duplicate) + ); + assert_eq!( + tracker.accept(&withdrawn, ISSUED_AT), + Ok(ClientBindingStatusUpdate::Accepted) + ); + assert!(tracker.current_presentation(ISSUED_AT).is_none()); + assert_eq!( + tracker.accept(¤t, ISSUED_AT), + Err(ClientBindingStatusFoldError::LowerRevisionReplay) + ); + assert_eq!( + tracker.accept(&conflicting_equal, ISSUED_AT), + Err(ClientBindingStatusFoldError::ConflictingEqualRevision) + ); + } + + #[test] + fn expiry_disconnect_and_scope_change_clear_presentation() { + let relay = Keys::generate(); + let other_relay = Keys::generate(); + let author = Keys::generate(); + let other_author = Keys::generate(); + let current = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + let mut tracker = + ClientBindingStatusTracker::new(relay.public_key(), domain(), author.public_key()); + tracker + .accept(¤t, ISSUED_AT) + .expect("current status accepted"); + + assert!(tracker.current_presentation(FRESH_UNTIL).is_none()); + assert_eq!(tracker.high_water_revision(), Some(11)); + assert_eq!( + tracker.accept(¤t, ISSUED_AT), + Ok(ClientBindingStatusUpdate::Duplicate) + ); + assert!(tracker.current_presentation(ISSUED_AT).is_none()); + + let newer = signed_status( + &relay, + author.public_key(), + 12, + ClientBindingStatusDisposition::DisplayCurrent, + ); + tracker + .accept(&newer, ISSUED_AT) + .expect("newer status accepted"); + tracker.on_disconnect(); + assert!(tracker.current_presentation(ISSUED_AT).is_none()); + assert_eq!(tracker.high_water_revision(), Some(12)); + + tracker.change_scope( + other_relay.public_key(), + CommunityId::from_uuid(Uuid::from_u128(999)), + other_author.public_key(), + ); + assert!(tracker.current_presentation(ISSUED_AT).is_none()); + assert_eq!(tracker.high_water_revision(), None); + } + + #[test] + fn debug_output_and_wire_omit_private_identity_material() { + let relay = Keys::generate(); + let author = Keys::generate(); + let event = signed_status( + &relay, + author.public_key(), + 11, + ClientBindingStatusDisposition::DisplayCurrent, + ); + let status = validate(&event, &relay, &author, ISSUED_AT) + .expect("synthetic current status validates"); + + let debug = format!("{status:?}"); + assert!(!debug.contains(DOMAIN)); + assert!(!debug.contains(&author.public_key().to_hex())); + assert!(!debug.contains("synthetic-policy-v1")); + assert!(!debug.contains("Synthetic Example")); + + let payload: Value = serde_json::from_str(&event.content).expect("content parses"); + for forbidden in [ + "iss", + "sub", + "issuer", + "audience", + "email", + "display_name", + "binding_id", + "bearer", + ] { + assert!( + payload.get(forbidden).is_none(), + "unexpected field {forbidden}" + ); + } + } +} diff --git a/crates/buzz-core/src/kind.rs b/crates/buzz-core/src/kind.rs index 76943c2ab..495fe8465 100644 --- a/crates/buzz-core/src/kind.rs +++ b/crates/buzz-core/src/kind.rs @@ -85,6 +85,12 @@ pub const KIND_AUTH: u32 = 22242; pub const KIND_BLOSSOM_AUTH: u32 = 24242; /// Buzz custom one-time identity binding proof (ephemeral, not stored). pub const KIND_NOSTR_IDENTITY_BINDING: u32 = 24243; +/// Buzz relay-authenticated client binding status (ephemeral, not stored). +/// +/// This provisional allocation carries short-lived, display-only status. It +/// is intentionally absent from relay ingest and storage allowlists until the +/// binding lifecycle and client-presentation joins are complete. +pub const KIND_CLIENT_BINDING_STATUS: u32 = 24244; /// NIP-98: HTTP auth event (used in nip98.rs, not stored). pub const KIND_HTTP_AUTH: u32 = 27235; @@ -823,6 +829,7 @@ pub const fn is_relay_only_kind(kind: u32) -> bool { matches!( kind, KIND_NIP43_MEMBERSHIP_LIST + | KIND_CLIENT_BINDING_STATUS | KIND_CHANNEL_SUMMARY | KIND_PRESENCE_SNAPSHOT | KIND_DM_VISIBILITY @@ -904,6 +911,12 @@ mod tests { assert!(!is_relay_only_kind(KIND_NIP43_LEAVE_REQUEST)); } + #[test] + fn client_binding_status_is_relay_only() { + assert!(is_relay_only_kind(KIND_CLIENT_BINDING_STATUS)); + assert!(is_ephemeral(KIND_CLIENT_BINDING_STATUS)); + } + #[test] fn parameterized_replaceable_range() { assert!(!is_parameterized_replaceable(29999)); diff --git a/crates/buzz-core/src/lib.rs b/crates/buzz-core/src/lib.rs index 66b7708f1..6be3e97c4 100644 --- a/crates/buzz-core/src/lib.rs +++ b/crates/buzz-core/src/lib.rs @@ -9,6 +9,8 @@ pub mod agent_turn_metric; /// Channel and membership enums shared across crates. pub mod channel; +/// Relay-authenticated, display-only client binding status contract. +pub mod client_binding_status; /// NIP-AE Agent Engrams — slug grammar, conversation key, d-tag derivation, /// body parse/serialize, envelope build/validate, head selection. pub mod engram; diff --git a/crates/buzz-db/src/client_status.rs b/crates/buzz-db/src/client_status.rs new file mode 100644 index 000000000..46e70b80e --- /dev/null +++ b/crates/buzz-db/src/client_status.rs @@ -0,0 +1,710 @@ +//! Durable, current-only client verification-status revisions. +//! +//! Allocation is transaction-owned: the exact active binding, membership, +//! invalidation generation/floors, database-clock freshness, revision row, +//! authority epoch, and idempotency receipt are committed together. + +use buzz_core::CommunityId; +use sqlx::{Postgres, Row, Transaction}; +use thiserror::Error; +use uuid::Uuid; + +use crate::authorization_invalidation::AuthorizationSelector; +use crate::Db; + +const CURRENT_KIND: &str = "client.status.current.v1"; +const WITHDRAW_KIND: &str = "client.status.withdraw.v1"; + +/// Exact private requirement for one current-status issuance. +pub struct CurrentStatusAllocation<'a> { + /// Server-resolved authorization domain. + pub community_id: CommunityId, + /// Exact event-author key. + pub event_author_pubkey: &'a [u8; 32], + /// Stable active binding ID. + pub binding_id: Uuid, + /// Exact positive binding version. + pub binding_version: u64, + /// Opaque current provider policy version. + pub policy_version: &'a str, + /// Invalidation generation captured before provider evaluation. + pub evaluation_generation: u64, + /// Database-clock freshness boundary. + pub fresh_until: u64, + /// Stable issuance operation ID. + pub operation_id: Uuid, + /// Exact event-input fingerprint. + pub request_fingerprint: [u8; 32], +} + +/// Exact private requirement for an opaque withdrawal. +pub struct WithdrawalStatusAllocation<'a> { + /// Server-resolved authorization domain. + pub community_id: CommunityId, + /// Exact event-author key. + pub event_author_pubkey: &'a [u8; 32], + /// Revision of the actual current issuance being withdrawn. + pub supersedes_revision: u64, + /// Fingerprint of the durable current issuance receipt. + pub issuance_fingerprint: [u8; 32], + /// Stable withdrawal operation ID. + pub operation_id: Uuid, + /// Exact withdrawal-input fingerprint. + pub request_fingerprint: [u8; 32], +} + +/// Result of a transaction-owned revision allocation. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct AllocatedStatusRevision { + /// Strictly positive revision. + pub revision: u64, + /// Domain-wide durable status floor after this allocation. + pub floor: u64, +} + +/// Allocation failure classified by whether PostgreSQL commit was attempted. +#[derive(Debug, Error)] +pub enum ClientStatusAllocationError { + /// Current private authority did not satisfy the exact requirement. + #[error("client status authority is not current")] + NotCurrent, + /// A stable operation ID was reused for different input. + #[error("client status operation conflicts with a committed request")] + ConflictingRetry, + /// Input could not be represented safely. + #[error("client status allocation input is invalid")] + InvalidInput, + /// PostgreSQL failed before commit was attempted; the transaction rolls back. + #[error("client status allocation failed before commit")] + Database(#[source] sqlx::Error), + /// PostgreSQL commit acknowledgement was ambiguous. The receipt decides. + #[error("client status commit acknowledgement is ambiguous")] + CommitUnknown(#[source] sqlx::Error), +} + +impl From for ClientStatusAllocationError { + fn from(error: sqlx::Error) -> Self { + Self::Database(error) + } +} + +impl Db { + /// Read an exact committed status allocation after ambiguous commit acknowledgement. + pub async fn committed_status_revision( + &self, + community_id: CommunityId, + operation_id: Uuid, + request_fingerprint: [u8; 32], + ) -> Result, ClientStatusAllocationError> { + let row = sqlx::query( + "SELECT request_fingerprint, result_payload FROM authorization_operation_receipts \ + WHERE community_id=$1 AND operation_id=$2", + ) + .bind(community_id.as_uuid()) + .bind(operation_id) + .fetch_optional(&self.pool) + .await?; + let Some(row) = row else { return Ok(None) }; + let fingerprint: Vec = row.try_get("request_fingerprint")?; + let payload: Vec = row.try_get("result_payload")?; + if fingerprint.as_slice() != request_fingerprint { + return Err(ClientStatusAllocationError::ConflictingRetry); + } + Ok(Some(decode_receipt_payload(&payload)?)) + } + + /// Allocate or replay one strictly monotonic current-status revision. + pub async fn allocate_current_status_revision( + &self, + request: CurrentStatusAllocation<'_>, + ) -> Result { + validate_common( + request.community_id, + request.event_author_pubkey, + request.operation_id, + request.binding_version, + )?; + if request.policy_version.is_empty() + || request.evaluation_generation > i64::MAX as u64 + || request.fresh_until > i64::MAX as u64 + { + return Err(ClientStatusAllocationError::InvalidInput); + } + let mut tx = self.begin_transaction().await.map_err(db_error)?; + lock_scope(&mut tx, request.community_id, request.event_author_pubkey).await?; + let (issuer, subject) = validate_current_authority(&mut tx, &request).await?; + validate_invalidation( + &mut tx, + request.community_id, + request.evaluation_generation, + request.binding_id, + request.binding_version, + request.event_author_pubkey, + request.policy_version, + &issuer, + &subject, + ) + .await?; + if let Some(revision) = replay_revision( + &mut tx, + request.community_id, + request.operation_id, + CURRENT_KIND, + request.request_fingerprint, + ) + .await? + { + tx.commit() + .await + .map_err(ClientStatusAllocationError::CommitUnknown)?; + return Ok(revision); + } + let allocated = next_revision(&mut tx, request.community_id).await?; + sqlx::query( + "INSERT INTO client_status_revisions \ + (community_id, event_author_pubkey, revision, disposition, binding_id, binding_version) \ + VALUES ($1, $2, $3, 'current', $4, $5) \ + ON CONFLICT (community_id, event_author_pubkey) DO UPDATE SET \ + revision=EXCLUDED.revision, disposition='current', binding_id=EXCLUDED.binding_id, \ + binding_version=EXCLUDED.binding_version, supersedes_revision=NULL, \ + updated_at=clock_timestamp()", + ) + .bind(request.community_id.as_uuid()) + .bind(request.event_author_pubkey.as_slice()) + .bind(allocated as i64) + .bind(request.binding_id) + .bind(request.binding_version as i64) + .execute(&mut *tx) + .await?; + insert_receipt( + &mut tx, + request.community_id, + request.operation_id, + CURRENT_KIND, + request.request_fingerprint, + allocated, + ) + .await?; + tx.commit() + .await + .map_err(ClientStatusAllocationError::CommitUnknown)?; + Ok(AllocatedStatusRevision { + revision: allocated, + floor: allocated, + }) + } + + /// Allocate or replay a withdrawal strictly after its exact current receipt. + pub async fn allocate_withdrawn_status_revision( + &self, + request: WithdrawalStatusAllocation<'_>, + ) -> Result { + validate_common( + request.community_id, + request.event_author_pubkey, + request.operation_id, + request.supersedes_revision, + )?; + let mut tx = self.begin_transaction().await.map_err(db_error)?; + lock_scope(&mut tx, request.community_id, request.event_author_pubkey).await?; + if let Some(revision) = replay_revision( + &mut tx, + request.community_id, + request.operation_id, + WITHDRAW_KIND, + request.request_fingerprint, + ) + .await? + { + tx.commit() + .await + .map_err(ClientStatusAllocationError::CommitUnknown)?; + return Ok(revision); + } + let issuance_exists: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM authorization_operation_receipts \ + WHERE community_id=$1 AND operation_kind=$2 AND request_fingerprint=$3 \ + AND octet_length(result_payload) IN (8,16) \ + AND substring(result_payload FROM 1 FOR 8)=$4)", + ) + .bind(request.community_id.as_uuid()) + .bind(CURRENT_KIND) + .bind(request.issuance_fingerprint.as_slice()) + .bind(request.supersedes_revision.to_be_bytes().as_slice()) + .fetch_one(&mut *tx) + .await?; + if !issuance_exists { + return Err(ClientStatusAllocationError::NotCurrent); + } + let row: Option<(i64, String, Option)> = sqlx::query_as( + "SELECT revision, disposition, supersedes_revision FROM client_status_revisions \ + WHERE community_id=$1 AND event_author_pubkey=$2 FOR UPDATE", + ) + .bind(request.community_id.as_uuid()) + .bind(request.event_author_pubkey.as_slice()) + .fetch_optional(&mut *tx) + .await?; + let Some((revision, disposition, prior_supersedes)) = row else { + return Err(ClientStatusAllocationError::NotCurrent); + }; + let receipt_revision = request.supersedes_revision as i64; + let superseded_revision = if disposition == "current" { + if receipt_revision > revision { + return Err(ClientStatusAllocationError::NotCurrent); + } + revision + } else if disposition == "withdrawn" + && prior_supersedes.is_some_and(|superseded| receipt_revision <= superseded) + && revision > receipt_revision + { + prior_supersedes.expect("withdrawn status has a superseded revision") + } else { + return Err(ClientStatusAllocationError::NotCurrent); + }; + let floor: i64 = sqlx::query_scalar( + "SELECT status_revision FROM authorization_authority_epochs \ + WHERE community_id=$1 FOR UPDATE", + ) + .bind(request.community_id.as_uuid()) + .fetch_one(&mut *tx) + .await?; + let allocated = if disposition == "current" || revision < floor { + next_revision(&mut tx, request.community_id).await? + } else { + revision as u64 + }; + if allocated <= superseded_revision as u64 { + return Err(ClientStatusAllocationError::NotCurrent); + } + sqlx::query( + "UPDATE client_status_revisions SET revision=$3, disposition='withdrawn', \ + binding_id=NULL, binding_version=NULL, supersedes_revision=$4, \ + updated_at=clock_timestamp() \ + WHERE community_id=$1 AND event_author_pubkey=$2", + ) + .bind(request.community_id.as_uuid()) + .bind(request.event_author_pubkey.as_slice()) + .bind(allocated as i64) + .bind(superseded_revision) + .execute(&mut *tx) + .await?; + insert_receipt( + &mut tx, + request.community_id, + request.operation_id, + WITHDRAW_KIND, + request.request_fingerprint, + allocated, + ) + .await?; + tx.commit() + .await + .map_err(ClientStatusAllocationError::CommitUnknown)?; + Ok(AllocatedStatusRevision { + revision: allocated, + floor: allocated, + }) + } +} + +fn validate_common( + community_id: CommunityId, + pubkey: &[u8; 32], + operation_id: Uuid, + positive: u64, +) -> Result<(), ClientStatusAllocationError> { + if community_id.as_uuid().is_nil() + || operation_id.is_nil() + || positive == 0 + || positive > i64::MAX as u64 + || pubkey.iter().all(|byte| *byte == 0) + { + return Err(ClientStatusAllocationError::InvalidInput); + } + Ok(()) +} + +fn db_error(error: crate::DbError) -> ClientStatusAllocationError { + match error { + crate::DbError::Sqlx(error) => ClientStatusAllocationError::Database(error), + _ => ClientStatusAllocationError::InvalidInput, + } +} + +async fn lock_scope( + tx: &mut Transaction<'static, Postgres>, + community_id: CommunityId, + pubkey: &[u8; 32], +) -> Result<(), ClientStatusAllocationError> { + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1, 0))") + .bind(format!( + "client-status:{}:{}", + community_id, + hex::encode(pubkey) + )) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn validate_current_authority( + tx: &mut Transaction<'static, Postgres>, + request: &CurrentStatusAllocation<'_>, +) -> Result<(String, String), ClientStatusAllocationError> { + let row: Option<(String, String)> = sqlx::query_as( + "SELECT binding.issuer, binding.uid FROM identity_bindings binding \ + JOIN identity_principals principal ON principal.community_id=binding.community_id \ + AND principal.issuer=binding.issuer AND principal.uid=binding.uid \ + JOIN relay_members member ON member.community_id=binding.community_id \ + AND member.pubkey=encode(binding.pubkey, 'hex') \ + WHERE binding.community_id=$1 AND binding.binding_id=$2 AND binding.pubkey=$3 \ + AND binding.binding_version=$4 AND binding.binding_state='active' \ + AND principal.disabled_at IS NULL \ + AND NOT EXISTS (SELECT 1 FROM identity_revoked_keys revoked \ + WHERE revoked.community_id=binding.community_id AND revoked.pubkey=binding.pubkey) \ + FOR SHARE OF binding, principal, member", + ) + .bind(request.community_id.as_uuid()) + .bind(request.binding_id) + .bind(request.event_author_pubkey.as_slice()) + .bind(request.binding_version as i64) + .fetch_optional(&mut **tx) + .await?; + let Some(principal) = row else { + return Err(ClientStatusAllocationError::NotCurrent); + }; + let fresh: bool = + sqlx::query_scalar("SELECT clock_timestamp() < to_timestamp($1::double precision)") + .bind(request.fresh_until as f64) + .fetch_one(&mut **tx) + .await?; + if !fresh { + return Err(ClientStatusAllocationError::NotCurrent); + } + Ok(principal) +} + +#[allow(clippy::too_many_arguments)] +async fn validate_invalidation( + tx: &mut Transaction<'static, Postgres>, + community_id: CommunityId, + evaluation_generation: u64, + binding_id: Uuid, + binding_version: u64, + pubkey: &[u8; 32], + policy_version: &str, + issuer: &str, + subject: &str, +) -> Result<(), ClientStatusAllocationError> { + let generation: Option = sqlx::query_scalar( + "SELECT generation FROM authorization_invalidation_domains \ + WHERE community_id=$1 FOR SHARE", + ) + .bind(community_id.as_uuid()) + .fetch_optional(&mut **tx) + .await?; + if generation != Some(evaluation_generation as i64) { + return Err(ClientStatusAllocationError::NotCurrent); + } + let selectors = [ + AuthorizationSelector::domain(), + AuthorizationSelector::principal(issuer, subject) + .map_err(|_| ClientStatusAllocationError::InvalidInput)?, + AuthorizationSelector::nostr_key(*pubkey), + AuthorizationSelector::binding(binding_id, binding_version) + .map_err(|_| ClientStatusAllocationError::InvalidInput)?, + AuthorizationSelector::policy_version(policy_version) + .map_err(|_| ClientStatusAllocationError::InvalidInput)?, + ]; + for selector in selectors { + let row = sqlx::query( + "SELECT generation, sticky_deny, binding_version_floor \ + FROM authorization_invalidation_floors WHERE community_id=$1 \ + AND selector_kind=$2 AND selector_fingerprint=$3 FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(selector.kind().as_str()) + .bind(selector.fingerprint().as_slice()) + .fetch_optional(&mut **tx) + .await?; + let Some(row) = row else { continue }; + let floor_generation: i64 = row.try_get("generation")?; + let sticky: bool = row.try_get("sticky_deny")?; + let version_floor: Option = row.try_get("binding_version_floor")?; + if sticky + || floor_generation > evaluation_generation as i64 + || version_floor.is_some_and(|floor| binding_version <= floor as u64) + { + return Err(ClientStatusAllocationError::NotCurrent); + } + } + Ok(()) +} + +async fn replay_revision( + tx: &mut Transaction<'static, Postgres>, + community_id: CommunityId, + operation_id: Uuid, + operation_kind: &str, + request_fingerprint: [u8; 32], +) -> Result, ClientStatusAllocationError> { + let row = sqlx::query( + "SELECT operation_kind, request_fingerprint, result_payload \ + FROM authorization_operation_receipts \ + WHERE community_id=$1 AND operation_id=$2 FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(operation_id) + .fetch_optional(&mut **tx) + .await?; + let Some(row) = row else { return Ok(None) }; + let kind: String = row.try_get("operation_kind")?; + let fingerprint: Vec = row.try_get("request_fingerprint")?; + let payload: Vec = row.try_get("result_payload")?; + if kind != operation_kind || fingerprint.as_slice() != request_fingerprint { + return Err(ClientStatusAllocationError::ConflictingRetry); + } + Ok(Some(decode_receipt_payload(&payload)?)) +} + +fn decode_receipt_payload( + payload: &[u8], +) -> Result { + let (revision_bytes, floor_bytes) = match payload.len() { + // Compatibility with receipts written before allocation-time floors + // were retained. The allocation revision was also its floor. + 8 => (&payload[..8], &payload[..8]), + 16 => (&payload[..8], &payload[8..16]), + _ => return Err(ClientStatusAllocationError::ConflictingRetry), + }; + let revision = u64::from_be_bytes( + revision_bytes + .try_into() + .map_err(|_| ClientStatusAllocationError::ConflictingRetry)?, + ); + let floor = u64::from_be_bytes( + floor_bytes + .try_into() + .map_err(|_| ClientStatusAllocationError::ConflictingRetry)?, + ); + if revision == 0 || floor == 0 || revision < floor { + return Err(ClientStatusAllocationError::ConflictingRetry); + } + Ok(AllocatedStatusRevision { revision, floor }) +} + +async fn next_revision( + tx: &mut Transaction<'static, Postgres>, + community_id: CommunityId, +) -> Result { + let revision: i64 = sqlx::query_scalar( + "UPDATE authorization_authority_epochs SET \ + authority_epoch=authority_epoch+1, status_revision=status_revision+1, \ + updated_at=clock_timestamp() WHERE community_id=$1 RETURNING status_revision", + ) + .bind(community_id.as_uuid()) + .fetch_one(&mut **tx) + .await?; + u64::try_from(revision).map_err(|_| ClientStatusAllocationError::InvalidInput) +} + +async fn insert_receipt( + tx: &mut Transaction<'static, Postgres>, + community_id: CommunityId, + operation_id: Uuid, + operation_kind: &str, + request_fingerprint: [u8; 32], + revision: u64, +) -> Result<(), ClientStatusAllocationError> { + let mut result_payload = Vec::with_capacity(16); + result_payload.extend_from_slice(&revision.to_be_bytes()); + result_payload.extend_from_slice(&revision.to_be_bytes()); + sqlx::query( + "INSERT INTO authorization_operation_receipts \ + (community_id, operation_id, operation_kind, request_fingerprint, \ + result_version, result_payload, lease_expires_at) \ + VALUES ($1,$2,$3,$4,1,$5,clock_timestamp()+interval '100 years')", + ) + .bind(community_id.as_uuid()) + .bind(operation_id) + .bind(operation_kind) + .bind(request_fingerprint.as_slice()) + .bind(result_payload) + .execute(&mut **tx) + .await?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + #[ignore = "requires migrated Postgres"] + async fn current_replay_withdrawal_and_revocation_are_transaction_owned() { + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_owned()); + let pool = sqlx::PgPool::connect(&database_url) + .await + .expect("test database"); + crate::migration::run_migrations(&pool) + .await + .expect("migrations"); + let db = Db::from_pool(pool); + let community = CommunityId::from_uuid(Uuid::new_v4()); + let binding_id = Uuid::new_v4(); + let author = [0x41; 32]; + sqlx::query("INSERT INTO communities (id, host) VALUES ($1,$2)") + .bind(community.as_uuid()) + .bind(format!("status-{}.example", community.as_uuid())) + .execute(&db.pool) + .await + .expect("community"); + sqlx::query("INSERT INTO identity_principals (community_id,issuer,uid) VALUES ($1,$2,$3)") + .bind(community.as_uuid()) + .bind("https://idp.example") + .bind("subject") + .execute(&db.pool) + .await + .expect("principal"); + sqlx::query( + "INSERT INTO identity_bindings \ + (community_id,issuer,uid,pubkey,source,binding_id,binding_version, \ + binding_state,binding_provenance) \ + VALUES ($1,$2,$3,$4,'jwt_npub',$5,1,'active','attested_key')", + ) + .bind(community.as_uuid()) + .bind("https://idp.example") + .bind("subject") + .bind(author.as_slice()) + .bind(binding_id) + .execute(&db.pool) + .await + .expect("binding"); + sqlx::query("INSERT INTO relay_members (community_id,pubkey,role) VALUES ($1,$2,'member')") + .bind(community.as_uuid()) + .bind(hex::encode(author)) + .execute(&db.pool) + .await + .expect("member"); + sqlx::query( + "INSERT INTO authorization_invalidation_domains (community_id) VALUES ($1) \ + ON CONFLICT (community_id) DO NOTHING", + ) + .bind(community.as_uuid()) + .execute(&db.pool) + .await + .expect("invalidation domain"); + let generation: i64 = sqlx::query_scalar( + "SELECT generation FROM authorization_invalidation_domains WHERE community_id=$1", + ) + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("generation"); + let fresh_until = chrono::Utc::now().timestamp() as u64 + 300; + let operation_id = Uuid::new_v4(); + let allocate = |operation_id, fingerprint| CurrentStatusAllocation { + community_id: community, + event_author_pubkey: &author, + binding_id, + binding_version: 1, + policy_version: "policy-v1", + evaluation_generation: generation as u64, + fresh_until, + operation_id, + request_fingerprint: fingerprint, + }; + let first = db + .allocate_current_status_revision(allocate(operation_id, [1; 32])) + .await + .expect("first current"); + let replay = db + .allocate_current_status_revision(allocate(operation_id, [1; 32])) + .await + .expect("exact replay"); + assert_eq!(first, replay); + let second = db + .allocate_current_status_revision(allocate(Uuid::new_v4(), [2; 32])) + .await + .expect("new issuance"); + assert!(second.revision > first.revision); + let withdrawal_operation = Uuid::new_v4(); + let withdrawn = db + .allocate_withdrawn_status_revision(WithdrawalStatusAllocation { + community_id: community, + event_author_pubkey: &author, + supersedes_revision: second.revision, + issuance_fingerprint: [2; 32], + operation_id: withdrawal_operation, + request_fingerprint: [3; 32], + }) + .await + .expect("withdrawal"); + assert!(withdrawn.revision > second.revision); + let fanout = db + .allocate_withdrawn_status_revision(WithdrawalStatusAllocation { + community_id: community, + event_author_pubkey: &author, + supersedes_revision: first.revision, + issuance_fingerprint: [1; 32], + operation_id: Uuid::new_v4(), + request_fingerprint: [4; 32], + }) + .await + .expect("older displayed current receives the same withdrawal"); + assert_eq!(fanout, withdrawn); + + sqlx::query( + "UPDATE authorization_authority_epochs \ + SET authority_epoch=authority_epoch+1, status_revision=status_revision+1 \ + WHERE community_id=$1", + ) + .bind(community.as_uuid()) + .execute(&db.pool) + .await + .expect("advance unrelated durable status floor"); + let delayed_replay = db + .allocate_withdrawn_status_revision(WithdrawalStatusAllocation { + community_id: community, + event_author_pubkey: &author, + supersedes_revision: second.revision, + issuance_fingerprint: [2; 32], + operation_id: withdrawal_operation, + request_fingerprint: [3; 32], + }) + .await + .expect("exact delayed fan-out replay retains allocation-time floor"); + assert_eq!(delayed_replay, withdrawn); + + let reissued = db + .allocate_current_status_revision(allocate(Uuid::new_v4(), [6; 32])) + .await + .expect("a fresh current status can replace a withdrawn projection"); + assert!(reissued.revision > withdrawn.revision); + let row: (String, Option) = sqlx::query_as( + "SELECT disposition, supersedes_revision FROM client_status_revisions \ + WHERE community_id=$1 AND event_author_pubkey=$2", + ) + .bind(community.as_uuid()) + .bind(author.as_slice()) + .fetch_one(&db.pool) + .await + .expect("reissued projection"); + assert_eq!(row, ("current".to_owned(), None)); + + assert!(matches!( + db.allocate_withdrawn_status_revision(WithdrawalStatusAllocation { + community_id: community, + event_author_pubkey: &author, + supersedes_revision: second.revision, + issuance_fingerprint: [9; 32], + operation_id: Uuid::new_v4(), + request_fingerprint: [5; 32], + }) + .await, + Err(ClientStatusAllocationError::NotCurrent) + )); + } +} diff --git a/crates/buzz-db/src/public_projection.rs b/crates/buzz-db/src/public_projection.rs new file mode 100644 index 000000000..516699e67 --- /dev/null +++ b/crates/buzz-db/src/public_projection.rs @@ -0,0 +1,2592 @@ +//! Durable reconciliation for the optional relay-authored identity projection. +//! +//! O3 lifecycle rows remain authoritative. This module stores only public +//! event coordinates and opaque binding generations; it is neither an +//! operator API nor a durable audit surface. + +use std::{fmt, time::Duration}; + +use buzz_core::{CommunityId, StoredEvent}; +use chrono::{DateTime, Utc}; +use nostr::Event; +use sqlx::{Postgres, Row, Transaction}; +use uuid::Uuid; + +use crate::{ + event::{self, EventQuery}, + identity_binding::{key_lock_coordinate, lock_identity_coordinates_tx}, + Db, DbError, Result, +}; + +const ASSERTION_KIND: i32 = 30382; +const CLAIM_LEASE: Duration = Duration::from_secs(30); + +/// Opaque source generation for one relay-authored public projection. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct ProjectionBindingOrigin { + binding_id: Uuid, + binding_version: u64, +} + +impl ProjectionBindingOrigin { + /// Stable binding identifier used only for server-side fencing. + pub const fn binding_id(self) -> Uuid { + self.binding_id + } + + /// Positive binding generation used only for server-side fencing. + pub const fn binding_version(self) -> u64 { + self.binding_version + } +} + +impl fmt::Debug for ProjectionBindingOrigin { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProjectionBindingOrigin") + .field("binding_id", &"[redacted]") + .field("binding_version", &"[redacted]") + .finish() + } +} + +/// Server-only disposition recorded for the current public projection head. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ProjectionDisposition { + /// A current binding owns a label-bearing assertion. + Active, + /// The assertion is the canonical label-free inactive replacement. + Inactive, +} + +impl ProjectionDisposition { + const fn as_str(self) -> &'static str { + match self { + Self::Active => "active", + Self::Inactive => "inactive", + } + } +} + +/// Current server-only ownership metadata for an assertion coordinate. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct ProjectionHead { + event_id: [u8; 32], + disposition: ProjectionDisposition, + origin: Option, +} + +impl ProjectionHead { + /// Exact signed event installed with this ownership record. + pub const fn event_id(self) -> [u8; 32] { + self.event_id + } + + /// Whether the current projection is active or inactive. + pub const fn disposition(self) -> ProjectionDisposition { + self.disposition + } + + /// Exact binding generation that created the current projection, if known. + pub const fn origin(self) -> Option { + self.origin + } +} + +impl fmt::Debug for ProjectionHead { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProjectionHead") + .field("event_id", &"[redacted]") + .field("disposition", &self.disposition) + .field("origin", &"[redacted]") + .finish() + } +} + +fn checked_version(value: i64) -> Result { + u64::try_from(value) + .map_err(|_| DbError::InvalidData("public projection version is invalid".to_owned())) +} + +fn checked_i64(value: u64) -> Result { + i64::try_from(value) + .map_err(|_| DbError::InvalidData("public projection version is invalid".to_owned())) +} + +fn validate_pubkey(value: &[u8]) -> Result<()> { + if value.len() != 32 { + return Err(DbError::InvalidData( + "public projection key must be 32 bytes".to_owned(), + )); + } + Ok(()) +} + +async fn authoritative_binding_for_exact_principal_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + issuer: &str, + subject: &str, + pubkey: &[u8], +) -> Result> { + let row = sqlx::query( + r#" + SELECT binding.binding_id, binding.binding_version + FROM identity_bindings binding + WHERE binding.community_id=$1 + AND binding.issuer=$2 + AND binding.uid=$3 + AND binding.pubkey=$4 + AND binding.binding_state='active' + AND binding.revoked_at IS NULL + AND binding.rotation_completed_at IS NULL + AND NOT EXISTS ( + SELECT 1 FROM identity_migration_denials denial + WHERE denial.community_id=binding.community_id + AND denial.issuer=binding.issuer AND denial.subject=binding.uid) + AND NOT EXISTS ( + SELECT 1 FROM identity_migration_denied_keys denial + WHERE denial.community_id=binding.community_id + AND denial.pubkey=binding.pubkey) + AND NOT EXISTS ( + SELECT 1 FROM identity_principals principal + WHERE principal.community_id=binding.community_id + AND principal.issuer=binding.issuer AND principal.uid=binding.uid + AND principal.disabled_at IS NOT NULL) + AND NOT EXISTS ( + SELECT 1 FROM identity_revoked_keys revoked + WHERE revoked.community_id=binding.community_id + AND revoked.pubkey=binding.pubkey) + AND NOT EXISTS ( + SELECT 1 FROM identity_pending_replacements pending + WHERE pending.community_id=binding.community_id + AND pending.issuer=binding.issuer AND pending.subject=binding.uid + AND pending.cleared_at IS NULL) + AND NOT EXISTS ( + SELECT 1 FROM identity_retired_pairs retired + WHERE retired.community_id=binding.community_id + AND retired.issuer=binding.issuer AND retired.subject=binding.uid + AND retired.pubkey=binding.pubkey) + FOR SHARE OF binding + "#, + ) + .bind(community_id.as_uuid()) + .bind(issuer) + .bind(subject) + .bind(pubkey) + .fetch_optional(&mut **tx) + .await?; + row.map(|row| { + Ok(ProjectionBindingOrigin { + binding_id: row.try_get("binding_id")?, + binding_version: checked_version(row.try_get("binding_version")?)?, + }) + }) + .transpose() +} + +async fn authoritative_binding_for_key_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + pubkey: &[u8], +) -> Result> { + let row = sqlx::query( + r#" + SELECT binding.binding_id, binding.binding_version + FROM identity_bindings binding + WHERE binding.community_id=$1 + AND binding.pubkey=$2 + AND binding.binding_state='active' + AND binding.revoked_at IS NULL + AND binding.rotation_completed_at IS NULL + AND NOT EXISTS ( + SELECT 1 FROM identity_migration_denials denial + WHERE denial.community_id=binding.community_id + AND denial.issuer=binding.issuer AND denial.subject=binding.uid) + AND NOT EXISTS ( + SELECT 1 FROM identity_migration_denied_keys denial + WHERE denial.community_id=binding.community_id + AND denial.pubkey=binding.pubkey) + AND NOT EXISTS ( + SELECT 1 FROM identity_principals principal + WHERE principal.community_id=binding.community_id + AND principal.issuer=binding.issuer AND principal.uid=binding.uid + AND principal.disabled_at IS NOT NULL) + AND NOT EXISTS ( + SELECT 1 FROM identity_revoked_keys revoked + WHERE revoked.community_id=binding.community_id + AND revoked.pubkey=binding.pubkey) + AND NOT EXISTS ( + SELECT 1 FROM identity_pending_replacements pending + WHERE pending.community_id=binding.community_id + AND pending.issuer=binding.issuer AND pending.subject=binding.uid + AND pending.cleared_at IS NULL) + AND NOT EXISTS ( + SELECT 1 FROM identity_retired_pairs retired + WHERE retired.community_id=binding.community_id + AND retired.issuer=binding.issuer AND retired.subject=binding.uid + AND retired.pubkey=binding.pubkey) + FOR SHARE OF binding + "#, + ) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut **tx) + .await?; + row.map(|row| { + Ok(ProjectionBindingOrigin { + binding_id: row.try_get("binding_id")?, + binding_version: checked_version(row.try_get("binding_version")?)?, + }) + }) + .transpose() +} + +async fn current_projection_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + relay_pubkey: &[u8], + subject_pubkey: &[u8], +) -> Result> { + let subject = hex::encode(subject_pubkey); + Ok(event::query_events_tx( + tx, + &EventQuery { + kinds: Some(vec![ASSERTION_KIND]), + pubkey: Some(relay_pubkey.to_vec()), + d_tag: Some(subject), + global_only: true, + limit: Some(1), + ..EventQuery::for_community(community_id) + }, + ) + .await? + .into_iter() + .next()) +} + +async fn projection_by_id_including_deleted_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + event_id: &[u8], +) -> Result> { + let row = sqlx::query( + "SELECT id, pubkey, created_at, kind, tags, content, sig, received_at, channel_id \ + FROM events WHERE community_id=$1 AND id=$2 \ + ORDER BY created_at DESC LIMIT 1 FOR SHARE", + ) + .bind(community_id.as_uuid()) + .bind(event_id) + .fetch_optional(&mut **tx) + .await?; + row.map(event::row_to_stored_event) + .transpose() + .map(Option::flatten) +} + +async fn projection_head_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + relay_pubkey: &[u8], + subject_pubkey: &[u8], +) -> Result> { + let row = sqlx::query( + "SELECT event_id, disposition, source_binding_id, source_binding_version \ + FROM identity_public_projection_heads \ + WHERE community_id=$1 AND relay_pubkey=$2 AND subject_pubkey=$3 FOR UPDATE", + ) + .bind(community_id.as_uuid()) + .bind(relay_pubkey) + .bind(subject_pubkey) + .fetch_optional(&mut **tx) + .await?; + row.map(|row| { + let event_id: Vec = row.try_get("event_id")?; + let event_id = event_id.try_into().map_err(|_| { + DbError::InvalidData("public projection head event id is invalid".to_owned()) + })?; + let disposition: String = row.try_get("disposition")?; + let disposition = match disposition.as_str() { + "active" => ProjectionDisposition::Active, + "inactive" => ProjectionDisposition::Inactive, + _ => { + return Err(DbError::InvalidData( + "public projection head disposition is invalid".to_owned(), + )) + } + }; + let binding_id: Option = row.try_get("source_binding_id")?; + let binding_version: Option = row.try_get("source_binding_version")?; + let origin = match (binding_id, binding_version) { + (Some(binding_id), Some(binding_version)) => Some(ProjectionBindingOrigin { + binding_id, + binding_version: checked_version(binding_version)?, + }), + (None, None) => None, + _ => { + return Err(DbError::InvalidData( + "public projection head origin is incomplete".to_owned(), + )) + } + }; + Ok(ProjectionHead { + event_id, + disposition, + origin, + }) + }) + .transpose() +} + +async fn upsert_projection_head_tx( + tx: &mut Transaction<'_, Postgres>, + community_id: CommunityId, + relay_pubkey: &[u8], + subject_pubkey: &[u8], + event: &Event, + disposition: ProjectionDisposition, + origin: Option, +) -> Result<()> { + let created_at = DateTime::::from_timestamp(event.created_at.as_secs() as i64, 0) + .ok_or(DbError::InvalidTimestamp(event.created_at.as_secs() as i64))?; + sqlx::query( + r#" + INSERT INTO identity_public_projection_heads + (community_id, relay_pubkey, subject_pubkey, event_id, + event_created_at, disposition, source_binding_id, + source_binding_version, updated_at) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,NOW()) + ON CONFLICT (community_id, relay_pubkey, subject_pubkey) DO UPDATE + SET event_id=EXCLUDED.event_id, + event_created_at=EXCLUDED.event_created_at, + disposition=EXCLUDED.disposition, + source_binding_id=EXCLUDED.source_binding_id, + source_binding_version=EXCLUDED.source_binding_version, + updated_at=NOW() + "#, + ) + .bind(community_id.as_uuid()) + .bind(relay_pubkey) + .bind(subject_pubkey) + .bind(event.id.as_bytes().as_slice()) + .bind(created_at) + .bind(disposition.as_str()) + .bind(origin.map(ProjectionBindingOrigin::binding_id)) + .bind( + origin + .map(ProjectionBindingOrigin::binding_version) + .map(checked_i64) + .transpose()?, + ) + .execute(&mut **tx) + .await?; + Ok(()) +} + +fn validate_event_coordinate( + event: &Event, + relay_pubkey: &[u8], + subject_pubkey: &[u8], +) -> Result { + let subject = hex::encode(subject_pubkey); + let exact_tag_count = |name: &str, value: &str| { + event + .tags + .iter() + .filter(|tag| { + let parts = tag.as_slice(); + parts.len() == 2 && parts[0] == name && parts[1] == value + }) + .count() + }; + if event.kind.as_u16() as i32 != ASSERTION_KIND + || event.pubkey.as_bytes() != relay_pubkey + || exact_tag_count("d", &subject) != 1 + || exact_tag_count("p", &subject) != 1 + || exact_tag_count("verified", "relay") != 1 + || !event.verify_id() + || !event.verify_signature() + { + return Err(DbError::InvalidData( + "public projection event coordinate is invalid".to_owned(), + )); + } + Ok(subject) +} + +fn validate_event_disposition(event: &Event, disposition: ProjectionDisposition) -> Result<()> { + let exact_tag_count = |name: &str, value: &str| { + event + .tags + .iter() + .filter(|tag| { + let parts = tag.as_slice(); + parts.len() == 2 && parts[0] == name && parts[1] == value + }) + .count() + }; + let expiration = event + .tags + .iter() + .filter(|tag| { + let parts = tag.as_slice(); + parts.len() == 2 && parts[0] == "expiration" + }) + .map(|tag| tag.as_slice()[1].parse::()) + .collect::, _>>() + .map_err(|_| DbError::InvalidData("public projection expiration is invalid".to_owned()))?; + let display_names = event + .tags + .iter() + .filter(|tag| { + let parts = tag.as_slice(); + parts.len() == 2 && parts[0] == "display_name" + }) + .map(|tag| tag.as_slice()[1].as_str()) + .collect::>(); + let valid = event.content.is_empty() + && match disposition { + ProjectionDisposition::Active => { + event.tags.len() == 6 + && exact_tag_count("active", "true") == 1 + && exact_tag_count("active", "false") == 0 + && expiration.len() == 1 + && expiration[0] > 0 + && display_names.len() == 1 + && !display_names[0].is_empty() + } + ProjectionDisposition::Inactive => { + event.tags.len() == 5 + && exact_tag_count("active", "false") == 1 + && exact_tag_count("active", "true") == 0 + && expiration.as_slice() == [0] + && display_names.is_empty() + } + }; + if !valid { + return Err(DbError::InvalidData( + "public projection disposition is invalid".to_owned(), + )); + } + Ok(()) +} + +fn canonical_later_projection_is_proven( + current: &StoredEvent, + head: ProjectionHead, + expected: &StoredEvent, + source: ProjectionBindingOrigin, + relay_pubkey: &[u8], + subject_pubkey: &[u8], +) -> Result { + validate_event_coordinate(¤t.event, relay_pubkey, subject_pubkey)?; + validate_event_disposition(¤t.event, head.disposition)?; + validate_event_coordinate(&expected.event, relay_pubkey, subject_pubkey)?; + validate_event_disposition(&expected.event, ProjectionDisposition::Inactive)?; + + let head_matches_current = head.event_id.as_slice() == current.event.id.as_bytes().as_slice(); + let current_is_later = current.event.created_at > expected.event.created_at + || (current.event.created_at == expected.event.created_at + && current.event.id.as_bytes().as_slice() < expected.event.id.as_bytes().as_slice()); + let later_owner_is_proven = head.origin.is_some_and(|origin| { + origin.binding_id != source.binding_id || origin.binding_version > source.binding_version + }); + + Ok(head_matches_current && current_is_later && later_owner_is_proven) +} + +/// Transaction-owned active publication permit. +pub struct ActivePublicProjectionPermit { + tx: Transaction<'static, Postgres>, + community_id: CommunityId, + relay_pubkey: Vec, + subject_pubkey: Vec, + origin: ProjectionBindingOrigin, + current: Option, +} + +impl fmt::Debug for ActivePublicProjectionPermit { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ActivePublicProjectionPermit") + .field("community_id", &"[redacted]") + .field("relay_pubkey", &"[redacted]") + .field("subject_pubkey", &"[redacted]") + .field("origin", &"[redacted]") + .finish_non_exhaustive() + } +} + +impl ActivePublicProjectionPermit { + /// Current stored projection at the locked coordinate. + pub fn current_projection(&self) -> Option<&StoredEvent> { + self.current.as_ref() + } + + /// Exact active binding generation retained through commit. + pub const fn origin(&self) -> ProjectionBindingOrigin { + self.origin + } + + /// Atomically replace or accept the candidate and record private ownership. + pub async fn commit( + mut self, + event: &Event, + disposition: ProjectionDisposition, + ) -> Result { + let subject = validate_event_coordinate(event, &self.relay_pubkey, &self.subject_pubkey)?; + validate_event_disposition(event, disposition)?; + let (candidate, inserted) = event::replace_parameterized_event_tx( + &mut self.tx, + self.community_id, + event, + &subject, + None, + ) + .await?; + let stored = if inserted { + candidate + } else { + let current = current_projection_tx( + &mut self.tx, + self.community_id, + &self.relay_pubkey, + &self.subject_pubkey, + ) + .await? + .ok_or_else(|| { + DbError::InvalidData("public projection replacement disappeared".to_owned()) + })?; + if current.event.id != event.id { + return Err(DbError::InvalidData( + "public projection candidate lost ordering".to_owned(), + )); + } + current + }; + upsert_projection_head_tx( + &mut self.tx, + self.community_id, + &self.relay_pubkey, + &self.subject_pubkey, + &stored.event, + disposition, + Some(self.origin), + ) + .await?; + self.tx.commit().await?; + Ok(stored) + } +} + +/// Begin active publication while retaining exact binding authority to commit. +pub async fn begin_active_public_projection( + db: &Db, + community_id: CommunityId, + relay_pubkey: &[u8], + issuer: &str, + subject: &str, + subject_pubkey: &[u8], +) -> Result> { + validate_pubkey(relay_pubkey)?; + validate_pubkey(subject_pubkey)?; + if issuer.is_empty() || subject.is_empty() { + return Err(DbError::InvalidData( + "public projection principal is invalid".to_owned(), + )); + } + let mut tx = db.pool.begin().await?; + sqlx::query("SET LOCAL lock_timeout = '3s'") + .execute(&mut *tx) + .await?; + lock_identity_coordinates_tx( + &mut tx, + vec![key_lock_coordinate(community_id, subject_pubkey)], + ) + .await?; + let Some(origin) = authoritative_binding_for_exact_principal_tx( + &mut tx, + community_id, + issuer, + subject, + subject_pubkey, + ) + .await? + else { + tx.rollback().await?; + return Ok(None); + }; + let current = + current_projection_tx(&mut tx, community_id, relay_pubkey, subject_pubkey).await?; + Ok(Some(ActivePublicProjectionPermit { + tx, + community_id, + relay_pubkey: relay_pubkey.to_vec(), + subject_pubkey: subject_pubkey.to_vec(), + origin, + current, + })) +} + +/// Materialize committed O3 revoke/rotate operations as retryable O4 work. +pub async fn materialize_public_projection_retirements( + db: &Db, + domains: &[CommunityId], + relay_pubkey: &[u8], +) -> Result { + validate_pubkey(relay_pubkey)?; + if domains.is_empty() { + return Ok(0); + } + let domain_ids = domains + .iter() + .map(|domain| *domain.as_uuid()) + .collect::>(); + let result = sqlx::query( + r#" + INSERT INTO identity_public_projection_retirements + (community_id, operation_id, relay_pubkey, old_pubkey, + source_binding_id, source_binding_version, operation_kind) + SELECT operation.community_id, operation.operation_id, $2, + operation.pubkey, retired.binding_id, + retired.binding_version - 1, operation.operation_kind + FROM identity_lifecycle_operations operation + LEFT JOIN LATERAL ( + SELECT history.binding_id, history.binding_version + FROM identity_binding_history history + WHERE history.community_id=operation.community_id + AND history.operation_id=operation.operation_id + AND history.binding_id IS NOT DISTINCT FROM operation.binding_id + AND history.pubkey=operation.pubkey + AND history.binding_state IN ('revoked', 'rotated') + AND history.binding_version > 1 + ORDER BY history.recorded_at DESC, history.history_id + LIMIT 1 + ) retired ON TRUE + WHERE operation.community_id = ANY($1) + AND operation.operation_kind IN ('revoke_key', 'rotate') + AND operation.pubkey IS NOT NULL + ON CONFLICT (community_id, operation_id, relay_pubkey) DO NOTHING + "#, + ) + .bind(domain_ids) + .bind(relay_pubkey) + .execute(&db.pool) + .await?; + Ok(result.rows_affected()) +} + +/// Retryable retirement work claimed by one relay replica. +#[derive(Clone, PartialEq, Eq)] +pub struct ProjectionRetirementClaim { + community_id: CommunityId, + operation_id: Uuid, + relay_pubkey: Vec, + old_pubkey: Vec, + source_origin: Option, + claim_token: Uuid, +} + +impl ProjectionRetirementClaim { + /// Server-resolved authorization domain. + pub const fn community_id(&self) -> CommunityId { + self.community_id + } + + /// Public subject key whose old assertion may need retirement. + pub fn old_pubkey(&self) -> &[u8] { + &self.old_pubkey + } + + /// Exact retired source generation, when the lifecycle transition had one. + pub const fn source_origin(&self) -> Option { + self.source_origin + } +} + +impl fmt::Debug for ProjectionRetirementClaim { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProjectionRetirementClaim") + .field("community_id", &"[redacted]") + .field("operation_id", &"[redacted]") + .field("relay_pubkey", &"[redacted]") + .field("old_pubkey", &"[redacted]") + .field("source_origin", &"[redacted]") + .finish() + } +} + +async fn claim_next( + db: &Db, + domains: &[CommunityId], + relay_pubkey: &[u8], + phase: &str, +) -> Result> { + validate_pubkey(relay_pubkey)?; + if domains.is_empty() { + return Ok(None); + } + let domain_ids = domains + .iter() + .map(|domain| *domain.as_uuid()) + .collect::>(); + let lease_seconds = i64::try_from(CLAIM_LEASE.as_secs()) + .map_err(|_| DbError::InvalidData("projection claim lease is invalid".to_owned()))?; + let row = sqlx::query( + r#" + WITH candidate AS ( + SELECT community_id, operation_id, relay_pubkey + FROM identity_public_projection_retirements + WHERE community_id=ANY($1) AND relay_pubkey=$2 AND phase=$3 + AND next_attempt_at <= NOW() + AND (claim_token IS NULL OR lease_until <= NOW()) + ORDER BY next_attempt_at, community_id, operation_id + FOR UPDATE SKIP LOCKED + LIMIT 1 + ) + UPDATE identity_public_projection_retirements work + SET claim_token=gen_random_uuid(), + lease_until=NOW() + ($4::DOUBLE PRECISION * INTERVAL '1 second'), + attempts=attempts+1, + updated_at=NOW() + FROM candidate + WHERE work.community_id=candidate.community_id + AND work.operation_id=candidate.operation_id + AND work.relay_pubkey=candidate.relay_pubkey + RETURNING work.community_id, work.operation_id, work.relay_pubkey, + work.old_pubkey, work.source_binding_id, + work.source_binding_version, work.claim_token + "#, + ) + .bind(domain_ids) + .bind(relay_pubkey) + .bind(phase) + .bind(lease_seconds) + .fetch_optional(&db.pool) + .await?; + row.map(|row| { + let source_binding_id: Option = row.try_get("source_binding_id")?; + let source_binding_version: Option = row.try_get("source_binding_version")?; + let source_origin = match (source_binding_id, source_binding_version) { + (Some(binding_id), Some(binding_version)) => Some(ProjectionBindingOrigin { + binding_id, + binding_version: checked_version(binding_version)?, + }), + (None, None) => None, + _ => { + return Err(DbError::InvalidData( + "projection retirement source is incomplete".to_owned(), + )) + } + }; + Ok(ProjectionRetirementClaim { + community_id: CommunityId::from_uuid(row.try_get("community_id")?), + operation_id: row.try_get("operation_id")?, + relay_pubkey: row.try_get("relay_pubkey")?, + old_pubkey: row.try_get("old_pubkey")?, + source_origin, + claim_token: row.try_get("claim_token")?, + }) + }) + .transpose() +} + +/// Claim one ready public-projection retirement operation. +pub async fn claim_public_projection_retirement( + db: &Db, + domains: &[CommunityId], + relay_pubkey: &[u8], +) -> Result> { + claim_next(db, domains, relay_pubkey, "projection").await +} + +/// Transaction-owned view of one claimed retirement. +pub struct ProjectionRetirementPermit { + tx: Transaction<'static, Postgres>, + claim: ProjectionRetirementClaim, + current: Option, + head: Option, + active_origin: Option, +} + +impl fmt::Debug for ProjectionRetirementPermit { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProjectionRetirementPermit") + .field("claim", &self.claim) + .field("current", &self.current.as_ref().map(|_| "[event]")) + .field("head", &self.head) + .field("active_origin", &"[redacted]") + .finish_non_exhaustive() + } +} + +impl ProjectionRetirementPermit { + /// Public subject key selected by the committed lifecycle operation. + pub fn old_pubkey(&self) -> &[u8] { + &self.claim.old_pubkey + } + + /// Current event at the locked assertion coordinate. + pub fn current_projection(&self) -> Option<&StoredEvent> { + self.current.as_ref() + } + + /// Current private projection ownership metadata. + pub const fn head(&self) -> Option { + self.head + } + + /// Current authoritative binding for the key, if it has been reused. + pub const fn active_origin(&self) -> Option { + self.active_origin + } + + /// Retired binding generation carried by the durable lifecycle work. + pub const fn source_origin(&self) -> Option { + self.claim.source_origin + } + + fn head_for_current(&self, disposition: ProjectionDisposition) -> Result { + let current = self.current.as_ref().ok_or_else(|| { + DbError::InvalidData("public projection retirement has no current event".to_owned()) + })?; + validate_event_coordinate( + ¤t.event, + &self.claim.relay_pubkey, + &self.claim.old_pubkey, + )?; + validate_event_disposition(¤t.event, disposition)?; + let head = self.head.ok_or_else(|| { + DbError::InvalidData("public projection ownership is unavailable".to_owned()) + })?; + if head.event_id != *current.event.id.as_bytes() || head.disposition != disposition { + return Err(DbError::InvalidData( + "public projection ownership does not match the current event".to_owned(), + )); + } + Ok(head) + } + + fn current_can_be_retired_by_source(&self) -> Result { + let (current, head) = match (self.current.as_ref(), self.head) { + (None, None) => return Ok(true), + (None, Some(_)) => { + return Err(DbError::InvalidData( + "public projection ownership exists without an event".to_owned(), + )) + } + (Some(_), None) => return Ok(true), + (Some(current), Some(head)) => (current, head), + }; + validate_event_coordinate( + ¤t.event, + &self.claim.relay_pubkey, + &self.claim.old_pubkey, + )?; + validate_event_disposition(¤t.event, head.disposition)?; + if head.event_id != *current.event.id.as_bytes() { + return Err(DbError::InvalidData( + "public projection ownership does not match the current event".to_owned(), + )); + } + Ok(match (self.claim.source_origin, head.origin) { + (_, None) => true, + (None, Some(_)) => false, + (Some(source), Some(origin)) if source.binding_id == origin.binding_id => { + origin.binding_version <= source.binding_version + } + (Some(_), Some(_)) => false, + }) + } + + async fn finish_terminal(mut self, phase: &str, outcome: &str) -> Result<()> { + let changed = sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET phase=$5, outcome=$6, claim_token=NULL, lease_until=NULL, \ + completed_at=NOW(), updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='projection' AND lease_until > NOW()", + ) + .bind(self.claim.community_id.as_uuid()) + .bind(self.claim.operation_id) + .bind(&self.claim.relay_pubkey) + .bind(self.claim.claim_token) + .bind(phase) + .bind(outcome) + .execute(&mut *self.tx) + .await?; + if changed.rows_affected() != 1 { + return Err(DbError::InvalidData( + "public projection retirement claim expired".to_owned(), + )); + } + self.tx.commit().await?; + Ok(()) + } + + /// Complete a job whose assertion coordinate is empty. + pub async fn finish_no_projection(self) -> Result<()> { + if self.current.is_some() || self.head.is_some() { + return Err(DbError::InvalidData( + "public projection coordinate is not empty".to_owned(), + )); + } + self.finish_terminal("completed", "no_projection").await + } + + /// Preserve a newer legitimate binding and finish the stale job. + pub async fn finish_superseded(self, current: &Event) -> Result<()> { + let origin = self.active_origin.ok_or_else(|| { + DbError::InvalidData("projection supersession lacks active binding".to_owned()) + })?; + validate_event_coordinate(current, &self.claim.relay_pubkey, &self.claim.old_pubkey)?; + validate_event_disposition(current, ProjectionDisposition::Active)?; + let current_id = self + .current + .as_ref() + .map(|stored| stored.event.id) + .ok_or_else(|| { + DbError::InvalidData("public projection retirement has no current event".to_owned()) + })?; + let head = self.head_for_current(ProjectionDisposition::Active)?; + if current.id != current_id + || head.origin != Some(origin) + || self.claim.source_origin == Some(origin) + { + return Err(DbError::InvalidData( + "active binding does not own the current public projection".to_owned(), + )); + } + self.finish_terminal("superseded", "newer_binding").await + } + + /// Preserve a projection owned by a later binding generation. + pub async fn finish_newer_projection(self) -> Result<()> { + let head = self.head_for_current(ProjectionDisposition::Active)?; + let origin = head.origin.ok_or_else(|| { + DbError::InvalidData("newer public projection lacks an owner".to_owned()) + })?; + let is_newer = match self.claim.source_origin { + None => true, + Some(source) if source.binding_id == origin.binding_id => { + origin.binding_version > source.binding_version + } + Some(source) => source != origin, + }; + if !is_newer { + return Err(DbError::InvalidData( + "public projection is not owned by a later generation".to_owned(), + )); + } + self.finish_terminal("superseded", "newer_projection").await + } + + /// Finish a stale/replayed retirement when the exact coordinate is already + /// the canonical inactive projection. + /// + /// This is a metadata-only convergence path: it preserves the event and + /// ownership head byte-for-byte and never relabels an old assertion to the + /// replayed job's source generation. + pub async fn finish_existing_inactive(self) -> Result<()> { + let current = self.current.as_ref().ok_or_else(|| { + DbError::InvalidData("public projection retirement has no current event".to_owned()) + })?; + validate_event_coordinate( + ¤t.event, + &self.claim.relay_pubkey, + &self.claim.old_pubkey, + )?; + validate_event_disposition(¤t.event, ProjectionDisposition::Inactive)?; + if let Some(head) = self.head { + if head.event_id != *current.event.id.as_bytes() + || head.disposition != ProjectionDisposition::Inactive + { + return Err(DbError::InvalidData( + "public projection ownership does not match the current event".to_owned(), + )); + } + } + self.finish_terminal("completed", "already_inactive").await + } + + /// Atomically install/accept the canonical inactive event and queue delivery. + pub async fn finish_inactive(mut self, inactive: &Event) -> Result { + if !self.current_can_be_retired_by_source()? { + return Err(DbError::InvalidData( + "public projection belongs to a later binding generation".to_owned(), + )); + } + let subject = + validate_event_coordinate(inactive, &self.claim.relay_pubkey, &self.claim.old_pubkey)?; + validate_event_disposition(inactive, ProjectionDisposition::Inactive)?; + let (candidate, inserted) = event::replace_parameterized_event_tx( + &mut self.tx, + self.claim.community_id, + inactive, + &subject, + None, + ) + .await?; + let stored = if inserted { + candidate + } else { + let current = current_projection_tx( + &mut self.tx, + self.claim.community_id, + &self.claim.relay_pubkey, + &self.claim.old_pubkey, + ) + .await? + .ok_or_else(|| { + DbError::InvalidData("inactive projection replacement disappeared".to_owned()) + })?; + if current.event.id != inactive.id { + return Err(DbError::InvalidData( + "inactive projection candidate lost ordering".to_owned(), + )); + } + current + }; + upsert_projection_head_tx( + &mut self.tx, + self.claim.community_id, + &self.claim.relay_pubkey, + &self.claim.old_pubkey, + &stored.event, + ProjectionDisposition::Inactive, + self.claim.source_origin, + ) + .await?; + let changed = sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET phase='delivery', outcome=$5, event_id=$6, claim_token=NULL, \ + lease_until=NULL, next_attempt_at=NOW(), updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='projection' AND lease_until > NOW()", + ) + .bind(self.claim.community_id.as_uuid()) + .bind(self.claim.operation_id) + .bind(&self.claim.relay_pubkey) + .bind(self.claim.claim_token) + .bind(if inserted { + "replaced_inactive" + } else { + "already_inactive" + }) + .bind(stored.event.id.as_bytes().as_slice()) + .execute(&mut *self.tx) + .await?; + if changed.rows_affected() != 1 { + return Err(DbError::InvalidData( + "public projection retirement claim expired".to_owned(), + )); + } + self.tx.commit().await?; + Ok(stored) + } + + /// Release retryable work with bounded backoff and no authority change. + pub async fn defer(mut self) -> Result<()> { + sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET claim_token=NULL, lease_until=NULL, \ + next_attempt_at=NOW() + (LEAST(60, GREATEST(1, attempts))::DOUBLE PRECISION * INTERVAL '1 second'), \ + updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='projection'", + ) + .bind(self.claim.community_id.as_uuid()) + .bind(self.claim.operation_id) + .bind(&self.claim.relay_pubkey) + .bind(self.claim.claim_token) + .execute(&mut *self.tx) + .await?; + self.tx.commit().await?; + Ok(()) + } +} + +/// Revalidate and lock one claimed retirement through its event commit boundary. +pub async fn begin_public_projection_retirement( + db: &Db, + claim: ProjectionRetirementClaim, +) -> Result { + let mut tx = db.pool.begin().await?; + sqlx::query("SET LOCAL lock_timeout = '3s'") + .execute(&mut *tx) + .await?; + lock_identity_coordinates_tx( + &mut tx, + vec![key_lock_coordinate(claim.community_id, &claim.old_pubkey)], + ) + .await?; + let claimed = sqlx::query( + "SELECT 1 FROM identity_public_projection_retirements \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='projection' AND lease_until > NOW() FOR UPDATE", + ) + .bind(claim.community_id.as_uuid()) + .bind(claim.operation_id) + .bind(&claim.relay_pubkey) + .bind(claim.claim_token) + .fetch_optional(&mut *tx) + .await? + .is_some(); + if !claimed { + return Err(DbError::InvalidData( + "public projection retirement claim expired".to_owned(), + )); + } + let active_origin = + authoritative_binding_for_key_tx(&mut tx, claim.community_id, &claim.old_pubkey).await?; + let head = projection_head_tx( + &mut tx, + claim.community_id, + &claim.relay_pubkey, + &claim.old_pubkey, + ) + .await?; + let current = current_projection_tx( + &mut tx, + claim.community_id, + &claim.relay_pubkey, + &claim.old_pubkey, + ) + .await?; + Ok(ProjectionRetirementPermit { + tx, + claim, + current, + head, + active_origin, + }) +} + +/// Delivery work retained until Redis and local fan-out have both been attempted. +pub struct ProjectionDeliveryPermit { + tx: Transaction<'static, Postgres>, + claim: ProjectionRetirementClaim, + stored: StoredEvent, +} + +impl fmt::Debug for ProjectionDeliveryPermit { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProjectionDeliveryPermit") + .field("claim", &self.claim) + .field("stored", &"[event]") + .finish_non_exhaustive() + } +} + +impl ProjectionDeliveryPermit { + /// Exact current inactive event retained under the identity-key lock. + pub const fn stored(&self) -> &StoredEvent { + &self.stored + } + + /// Server-resolved authorization domain. + pub const fn community_id(&self) -> CommunityId { + self.claim.community_id + } + + /// Mark delivery complete while the exact projection head remains locked. + pub async fn complete(mut self) -> Result<()> { + let changed = sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET phase='completed', claim_token=NULL, lease_until=NULL, \ + completed_at=NOW(), updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='delivery' AND lease_until > NOW()", + ) + .bind(self.claim.community_id.as_uuid()) + .bind(self.claim.operation_id) + .bind(&self.claim.relay_pubkey) + .bind(self.claim.claim_token) + .execute(&mut *self.tx) + .await?; + if changed.rows_affected() != 1 { + return Err(DbError::InvalidData( + "public projection delivery claim expired".to_owned(), + )); + } + self.tx.commit().await?; + Ok(()) + } + + /// Release delivery for bounded retry without changing the event head. + pub async fn defer(mut self) -> Result<()> { + sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET claim_token=NULL, lease_until=NULL, \ + next_attempt_at=NOW() + (LEAST(60, GREATEST(1, attempts))::DOUBLE PRECISION * INTERVAL '1 second'), \ + updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='delivery'", + ) + .bind(self.claim.community_id.as_uuid()) + .bind(self.claim.operation_id) + .bind(&self.claim.relay_pubkey) + .bind(self.claim.claim_token) + .execute(&mut *self.tx) + .await?; + self.tx.commit().await?; + Ok(()) + } +} + +/// Claim and lock one pending inactive-event delivery. +pub async fn begin_public_projection_delivery( + db: &Db, + domains: &[CommunityId], + relay_pubkey: &[u8], +) -> Result> { + let Some(claim) = claim_next(db, domains, relay_pubkey, "delivery").await? else { + return Ok(None); + }; + let mut tx = db.pool.begin().await?; + sqlx::query("SET LOCAL lock_timeout = '3s'") + .execute(&mut *tx) + .await?; + lock_identity_coordinates_tx( + &mut tx, + vec![key_lock_coordinate(claim.community_id, &claim.old_pubkey)], + ) + .await?; + let row = sqlx::query( + "SELECT event_id FROM identity_public_projection_retirements \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='delivery' AND lease_until > NOW() FOR UPDATE", + ) + .bind(claim.community_id.as_uuid()) + .bind(claim.operation_id) + .bind(&claim.relay_pubkey) + .bind(claim.claim_token) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| DbError::InvalidData("public projection delivery claim expired".to_owned()))?; + let expected_id: Vec = row.try_get("event_id")?; + let current = current_projection_tx( + &mut tx, + claim.community_id, + &claim.relay_pubkey, + &claim.old_pubkey, + ) + .await?; + let head = projection_head_tx( + &mut tx, + claim.community_id, + &claim.relay_pubkey, + &claim.old_pubkey, + ) + .await?; + let head_matches_inactive = head.is_some_and(|head| { + head.event_id.as_slice() == expected_id.as_slice() + && head.disposition == ProjectionDisposition::Inactive + }); + let Some(stored) = current + .as_ref() + .filter(|stored| { + stored.event.id.as_bytes().as_slice() == expected_id.as_slice() && head_matches_inactive + }) + .cloned() + else { + let expected = + projection_by_id_including_deleted_tx(&mut tx, claim.community_id, &expected_id) + .await?; + let later_projection_is_proven = match ( + current.as_ref(), + head, + expected.as_ref(), + claim.source_origin, + ) { + (Some(current), Some(head), Some(expected), Some(source)) => { + canonical_later_projection_is_proven( + current, + head, + expected, + source, + &claim.relay_pubkey, + &claim.old_pubkey, + )? + } + _ => false, + }; + if !later_projection_is_proven { + return Err(DbError::InvalidData( + "public projection delivery lost its canonical ownership proof".to_owned(), + )); + } + let changed = sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET phase='superseded', outcome='newer_projection', claim_token=NULL, \ + lease_until=NULL, completed_at=NOW(), updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4 AND phase='delivery' AND lease_until > NOW()", + ) + .bind(claim.community_id.as_uuid()) + .bind(claim.operation_id) + .bind(&claim.relay_pubkey) + .bind(claim.claim_token) + .execute(&mut *tx) + .await?; + if changed.rows_affected() != 1 { + return Err(DbError::InvalidData( + "public projection delivery claim expired".to_owned(), + )); + } + tx.commit().await?; + return Ok(None); + }; + validate_event_coordinate(&stored.event, &claim.relay_pubkey, &claim.old_pubkey)?; + validate_event_disposition(&stored.event, ProjectionDisposition::Inactive)?; + Ok(Some(ProjectionDeliveryPermit { tx, claim, stored })) +} + +/// Count unfinished work for a relay author and exact domain set. +pub async fn unfinished_public_projection_retirements( + db: &Db, + domains: &[CommunityId], + relay_pubkey: &[u8], +) -> Result { + validate_pubkey(relay_pubkey)?; + if domains.is_empty() { + return Ok(0); + } + let domain_ids = domains + .iter() + .map(|domain| *domain.as_uuid()) + .collect::>(); + let count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM identity_public_projection_retirements \ + WHERE community_id=ANY($1) AND relay_pubkey=$2 \ + AND phase IN ('projection', 'delivery')", + ) + .bind(domain_ids) + .bind(relay_pubkey) + .fetch_one(&db.pool) + .await?; + u64::try_from(count) + .map_err(|_| DbError::InvalidData("projection retirement count is invalid".to_owned())) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::identity_binding::{ + resolve_identity_binding, BindingProvenance, EnrollmentMode, ResolveBindingInput, + ResolveBindingResult, + }; + use crate::identity_lifecycle::{ + revoke_identity_key, rotate_identity_binding, IdentityPrincipal, LifecycleContext, + VerifiedReplacementKey, + }; + use nostr::{EventBuilder, Keys, Kind, Tag, Timestamp}; + use sqlx::PgPool; + + const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; + const ISSUER: &str = "https://idp.example"; + + async fn setup() -> (Db, CommunityId) { + 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()); + let pool = PgPool::connect(&database_url).await.expect("test database"); + crate::migration::run_migrations(&pool) + .await + .expect("run migrations"); + let id = Uuid::new_v4(); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1,$2)") + .bind(id) + .bind(format!("public-projection-{}.example", id.simple())) + .execute(&pool) + .await + .expect("insert community"); + (Db::from_pool(pool), CommunityId::from_uuid(id)) + } + + fn assertion(keys: &Keys, subject: nostr::PublicKey, active: bool, at: u64) -> Event { + let subject = subject.to_hex(); + let active_value = if active { "true" } else { "false" }; + let expiration = if active { "500" } else { "0" }; + let mut tags = vec![ + Tag::parse(["d", subject.as_str()]).expect("d tag"), + Tag::parse(["p", subject.as_str()]).expect("p tag"), + Tag::parse(["verified", "relay"]).expect("verified tag"), + Tag::parse(["active", active_value]).expect("active tag"), + Tag::parse(["expiration", expiration]).expect("expiration tag"), + ]; + if active { + tags.push(Tag::parse(["display_name", "Approved Label"]).expect("display tag")); + } + EventBuilder::new(Kind::Custom(ASSERTION_KIND as u16), "") + .tags(tags) + .custom_created_at(Timestamp::from(at)) + .sign_with_keys(keys) + .expect("sign assertion") + } + + fn stored(event: Event) -> StoredEvent { + StoredEvent::with_received_at(event, Utc::now(), None, true) + } + + #[test] + fn delivery_supersession_requires_a_canonical_later_projection_and_owner() { + let relay = Keys::generate(); + let subject = Keys::generate().public_key(); + let source = ProjectionBindingOrigin { + binding_id: Uuid::new_v4(), + binding_version: 7, + }; + let later = ProjectionBindingOrigin { + binding_id: source.binding_id, + binding_version: 8, + }; + let expected = stored(assertion(&relay, subject, false, 100)); + let current = stored(assertion(&relay, subject, true, 101)); + let head = ProjectionHead { + event_id: *current.event.id.as_bytes(), + disposition: ProjectionDisposition::Active, + origin: Some(later), + }; + + assert!(canonical_later_projection_is_proven( + ¤t, + head, + &expected, + source, + relay.public_key().as_bytes(), + subject.as_bytes(), + ) + .expect("canonical proof")); + + let stale_head = ProjectionHead { + event_id: *expected.event.id.as_bytes(), + ..head + }; + assert!(!canonical_later_projection_is_proven( + ¤t, + stale_head, + &expected, + source, + relay.public_key().as_bytes(), + subject.as_bytes(), + ) + .expect("stale head is a denied proof")); + + let unreplaced_head = ProjectionHead { + origin: Some(source), + ..head + }; + assert!(!canonical_later_projection_is_proven( + ¤t, + unreplaced_head, + &expected, + source, + relay.public_key().as_bytes(), + subject.as_bytes(), + ) + .expect("same owner generation is a denied proof")); + + let noncanonical_expected = stored(assertion(&relay, subject, true, 100)); + assert!(canonical_later_projection_is_proven( + ¤t, + head, + &noncanonical_expected, + source, + relay.public_key().as_bytes(), + subject.as_bytes(), + ) + .is_err()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn delivery_converges_only_after_a_later_projection_owns_the_coordinate() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let subject = Keys::generate().public_key(); + let source = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "delivery-race-subject", + pubkey: subject.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll binding") + { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected binding result: {other:?}"), + }; + let active = assertion(&relay, subject, true, 100); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "delivery-race-subject", + subject.as_bytes(), + ) + .await + .expect("begin active projection") + .expect("active binding") + .commit(&active, ProjectionDisposition::Active) + .await + .expect("commit active projection"); + + let operation_id = Uuid::new_v4(); + revoke_identity_key( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "delivery race", + }, + subject.as_bytes(), + ) + .await + .expect("commit revocation"); + materialize_public_projection_retirements(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("materialize retirement"); + let claim = + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim retirement") + .expect("retirement exists"); + let expected = assertion(&relay, subject, false, 101); + begin_public_projection_retirement(&db, claim) + .await + .expect("begin retirement") + .finish_inactive(&expected) + .await + .expect("commit inactive projection"); + + let later = assertion(&relay, subject, true, 102); + let later_origin = ProjectionBindingOrigin { + binding_id: source.binding_id, + binding_version: source.binding_version + 1, + }; + let mut tx = db.pool.begin().await.expect("begin replacement"); + let subject_hex = subject.to_hex(); + let (_, inserted) = + event::replace_parameterized_event_tx(&mut tx, community, &later, &subject_hex, None) + .await + .expect("replace inactive projection"); + assert!(inserted, "later projection must win canonical ordering"); + upsert_projection_head_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + subject.as_bytes(), + &later, + ProjectionDisposition::Active, + Some(later_origin), + ) + .await + .expect("install later projection ownership"); + tx.commit().await.expect("commit later projection"); + + assert!( + begin_public_projection_delivery(&db, &[community], relay.public_key().as_bytes(),) + .await + .expect("converge stale delivery") + .is_none(), + "a fully proven later projection must terminalize stale delivery" + ); + assert_eq!( + unfinished_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("count unfinished work"), + 0 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn committed_revoke_materializes_retries_and_completes_exactly_once() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let subject_keys = Keys::generate(); + let subject = subject_keys.public_key(); + let subject_bytes = subject.to_bytes(); + let resolved = resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "subject-one", + pubkey: subject_bytes.as_slice(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll binding"); + let expected_origin = match resolved { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected binding result: {other:?}"), + }; + assert_eq!(expected_origin.binding_version(), 1); + + let active = assertion(&relay, subject, true, 100); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "subject-one", + subject.as_bytes(), + ) + .await + .expect("begin active projection") + .expect("active binding") + .commit(&active, ProjectionDisposition::Active) + .await + .expect("commit active projection"); + + let operation_id = Uuid::new_v4(); + revoke_identity_key( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "test revocation", + }, + subject.as_bytes(), + ) + .await + .expect("commit revocation"); + + assert_eq!( + materialize_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("materialize work"), + 1 + ); + assert_eq!( + materialize_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("duplicate discovery is idempotent"), + 0 + ); + let domains = [community]; + let relay_pubkey = relay.public_key().to_bytes(); + let (first_replica, second_replica) = tokio::join!( + claim_public_projection_retirement(&db, &domains, relay_pubkey.as_slice()), + claim_public_projection_retirement(&db, &domains, relay_pubkey.as_slice()), + ); + let mut claims = [ + first_replica.expect("first replica claim"), + second_replica.expect("second replica claim"), + ] + .into_iter() + .flatten() + .collect::>(); + assert_eq!(claims.len(), 1, "only one replica may own the claim"); + let stale_claim = claims.pop().expect("one claim"); + sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET lease_until=NOW() - INTERVAL '1 second' \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND claim_token=$4", + ) + .bind(community.as_uuid()) + .bind(operation_id) + .bind(relay.public_key().as_bytes()) + .bind(stale_claim.claim_token) + .execute(&db.pool) + .await + .expect("simulate crashed claim owner"); + assert!( + begin_public_projection_retirement(&db, stale_claim) + .await + .is_err(), + "an expired owner must not mutate the projection" + ); + let claim = + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("reclaim crashed work") + .expect("reclaimable work exists"); + assert_eq!(claim.source_origin(), Some(expected_origin)); + let permit = begin_public_projection_retirement(&db, claim) + .await + .expect("begin retirement"); + assert_eq!(permit.active_origin(), None); + assert_eq!( + permit.head().and_then(ProjectionHead::origin), + Some(expected_origin) + ); + let inactive = assertion(&relay, subject, false, 101); + permit + .finish_inactive(&inactive) + .await + .expect("commit inactive projection"); + + let delivery = + begin_public_projection_delivery(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim delivery") + .expect("delivery exists"); + assert_eq!(delivery.stored().event.id, inactive.id); + delivery.defer().await.expect("defer failed delivery"); + sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET next_attempt_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3 \ + AND phase='delivery'", + ) + .bind(community.as_uuid()) + .bind(operation_id) + .bind(relay.public_key().as_bytes()) + .execute(&db.pool) + .await + .expect("make deferred delivery ready"); + let delivery = + begin_public_projection_delivery(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("retry delivery") + .expect("deferred delivery exists"); + assert_eq!(delivery.stored().event.id, inactive.id); + delivery.complete().await.expect("complete delivery"); + assert_eq!( + unfinished_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("count unfinished"), + 0 + ); + assert!(claim_public_projection_retirement( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("replay claim") + .is_none()); + assert!( + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "subject-one", + subject.as_bytes(), + ) + .await + .expect("revalidate revoked principal") + .is_none(), + "committed revocation must prevent assertion reactivation" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn committed_rotation_retires_only_the_old_projection_generation() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let old_keys = Keys::generate(); + let new_keys = Keys::generate(); + let old_key = old_keys.public_key(); + let new_key = new_keys.public_key(); + let old_origin = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "rotated-subject", + pubkey: old_key.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll rotation source") + { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected rotation enrollment: {other:?}"), + }; + let old_active = assertion(&relay, old_key, true, 200); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "rotated-subject", + old_key.as_bytes(), + ) + .await + .expect("begin old projection") + .expect("old binding is active") + .commit(&old_active, ProjectionDisposition::Active) + .await + .expect("commit old projection"); + + let operation_id = Uuid::new_v4(); + rotate_identity_binding( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "test rotation", + }, + IdentityPrincipal { + issuer: ISSUER, + subject: "rotated-subject", + }, + old_key.as_bytes(), + VerifiedReplacementKey::after_verified_proof( + new_key.as_bytes(), + None, + BindingProvenance::AttestedKey, + Some("test-policy-v1"), + ) + .expect("verified replacement"), + ) + .await + .expect("commit rotation"); + + let new_active = assertion(&relay, new_key, true, 201); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "rotated-subject", + new_key.as_bytes(), + ) + .await + .expect("begin replacement projection") + .expect("replacement binding is active") + .commit(&new_active, ProjectionDisposition::Active) + .await + .expect("commit replacement projection"); + + assert_eq!( + materialize_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("materialize rotation"), + 1 + ); + let claim = + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim rotation") + .expect("rotation work exists"); + assert_eq!(claim.source_origin(), Some(old_origin)); + let permit = begin_public_projection_retirement(&db, claim) + .await + .expect("begin old projection retirement"); + assert_eq!(permit.active_origin(), None); + let old_inactive = assertion(&relay, old_key, false, 202); + permit + .finish_inactive(&old_inactive) + .await + .expect("retire old projection"); + + let mut tx = db.pool.begin().await.expect("inspect projections"); + let current_old = current_projection_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + old_key.as_bytes(), + ) + .await + .expect("read old projection") + .expect("old projection exists"); + let current_new = current_projection_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + new_key.as_bytes(), + ) + .await + .expect("read replacement projection") + .expect("replacement projection exists"); + tx.rollback().await.expect("rollback inspection"); + assert_eq!(current_old.event.id, old_inactive.id); + assert_eq!(current_new.event.id, new_active.id); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn active_key_reuse_without_republication_cannot_claim_the_old_assertion() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let reused_keys = Keys::generate(); + let replacement_keys = Keys::generate(); + let reused_key = reused_keys.public_key(); + let replacement_key = replacement_keys.public_key(); + let original_origin = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "original-principal", + pubkey: reused_key.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll original binding") + { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected original enrollment: {other:?}"), + }; + let original = assertion(&relay, reused_key, true, 250); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "original-principal", + reused_key.as_bytes(), + ) + .await + .expect("begin original projection") + .expect("original binding is active") + .commit(&original, ProjectionDisposition::Active) + .await + .expect("publish original projection"); + + let operation_id = Uuid::new_v4(); + rotate_identity_binding( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "free key for unpublished reuse", + }, + IdentityPrincipal { + issuer: ISSUER, + subject: "original-principal", + }, + reused_key.as_bytes(), + VerifiedReplacementKey::after_verified_proof( + replacement_key.as_bytes(), + None, + BindingProvenance::AttestedKey, + Some("test-policy-v1"), + ) + .expect("verified replacement"), + ) + .await + .expect("rotate original binding"); + let reused_origin = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: "https://replacement-idp.example", + subject: "replacement-principal", + pubkey: reused_key.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("reuse key before publishing replacement") + { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected key reuse: {other:?}"), + }; + assert_ne!(reused_origin, original_origin); + + materialize_public_projection_retirements(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("materialize original retirement"); + let permit = begin_public_projection_retirement( + &db, + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim original retirement") + .expect("original retirement exists"), + ) + .await + .expect("begin original retirement"); + assert_eq!(permit.source_origin(), Some(original_origin)); + assert_eq!(permit.active_origin(), Some(reused_origin)); + assert_eq!( + permit.head().and_then(ProjectionHead::origin), + Some(original_origin), + "B has not published and must not own A's assertion" + ); + let current = permit + .current_projection() + .expect("original assertion remains current") + .event + .clone(); + assert!( + permit.finish_superseded(¤t).await.is_err(), + "active binding existence alone cannot transfer projection ownership" + ); + + let mut tx = db.pool.begin().await.expect("inspect unchanged projection"); + let stored = current_projection_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + reused_key.as_bytes(), + ) + .await + .expect("read current projection") + .expect("current projection remains"); + let head = projection_head_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + reused_key.as_bytes(), + ) + .await + .expect("read projection head") + .expect("projection head remains"); + tx.rollback().await.expect("rollback inspection"); + assert_eq!(stored.event.id, original.id); + assert_eq!(head.event_id(), *original.id.as_bytes()); + assert_eq!(head.origin(), Some(original_origin)); + assert_eq!( + unfinished_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("count retryable work"), + 1 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn binding_version_advance_without_republication_is_not_a_newer_projection() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let subject_keys = Keys::generate(); + let subject = subject_keys.public_key(); + let first_origin = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "strengthened-principal", + pubkey: subject.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::Tofu, + key_attested: false, + }, + ) + .await + .expect("enroll tofu binding") + { + ResolveBindingResult::Enrolled(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected tofu enrollment: {other:?}"), + }; + let original = assertion(&relay, subject, true, 275); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "strengthened-principal", + subject.as_bytes(), + ) + .await + .expect("begin original projection") + .expect("tofu binding is active") + .commit(&original, ProjectionDisposition::Active) + .await + .expect("publish version-one projection"); + let strengthened_origin = match resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "strengthened-principal", + pubkey: subject.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("strengthen binding without republishing") + { + ResolveBindingResult::Existing(evidence) => ProjectionBindingOrigin { + binding_id: evidence.binding_id, + binding_version: evidence.binding_version, + }, + other => panic!("unexpected strengthening result: {other:?}"), + }; + assert_eq!(strengthened_origin.binding_id(), first_origin.binding_id()); + assert!(strengthened_origin.binding_version() > first_origin.binding_version()); + + let operation_id = Uuid::new_v4(); + revoke_identity_key( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "revoke strengthened binding", + }, + subject.as_bytes(), + ) + .await + .expect("commit strengthened revocation"); + materialize_public_projection_retirements(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("materialize strengthened retirement"); + let permit = begin_public_projection_retirement( + &db, + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim strengthened retirement") + .expect("strengthened retirement exists"), + ) + .await + .expect("begin strengthened retirement"); + assert_eq!(permit.source_origin(), Some(strengthened_origin)); + assert_eq!(permit.active_origin(), None); + assert_eq!( + permit.head().and_then(ProjectionHead::origin), + Some(first_origin), + "the projection still belongs to version one" + ); + assert!( + permit.finish_newer_projection().await.is_err(), + "an older head cannot terminalize a later-generation retirement" + ); + assert_eq!( + unfinished_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("count retryable strengthened work"), + 1 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn stale_rotation_work_cannot_retire_a_reused_key_projection() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let reused_keys = Keys::generate(); + let replacement_keys = Keys::generate(); + let reused_key = reused_keys.public_key(); + let replacement_key = replacement_keys.public_key(); + resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "first-principal", + pubkey: reused_key.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll first principal"); + let first = assertion(&relay, reused_key, true, 300); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "first-principal", + reused_key.as_bytes(), + ) + .await + .expect("begin first projection") + .expect("first binding is active") + .commit(&first, ProjectionDisposition::Active) + .await + .expect("commit first projection"); + + let operation_id = Uuid::new_v4(); + rotate_identity_binding( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "free key for reuse test", + }, + IdentityPrincipal { + issuer: ISSUER, + subject: "first-principal", + }, + reused_key.as_bytes(), + VerifiedReplacementKey::after_verified_proof( + replacement_key.as_bytes(), + None, + BindingProvenance::AttestedKey, + Some("test-policy-v1"), + ) + .expect("verified replacement"), + ) + .await + .expect("rotate first principal"); + resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: "https://second-idp.example", + subject: "second-principal", + pubkey: reused_key.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("reuse key for independent principal"); + let second = assertion(&relay, reused_key, true, 301); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + "https://second-idp.example", + "second-principal", + reused_key.as_bytes(), + ) + .await + .expect("begin reused projection") + .expect("reused binding is active") + .commit(&second, ProjectionDisposition::Active) + .await + .expect("commit reused projection"); + + materialize_public_projection_retirements(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("materialize stale rotation"); + let claim = + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim stale rotation") + .expect("stale rotation exists"); + let permit = begin_public_projection_retirement(&db, claim) + .await + .expect("begin stale rotation"); + assert_ne!(permit.active_origin(), permit.source_origin()); + let current = permit + .current_projection() + .expect("reused projection exists") + .event + .clone(); + permit + .finish_superseded(¤t) + .await + .expect("preserve reused projection"); + + let mut tx = db.pool.begin().await.expect("inspect reused projection"); + let stored = current_projection_tx( + &mut tx, + community, + relay.public_key().as_bytes(), + reused_key.as_bytes(), + ) + .await + .expect("read reused projection") + .expect("reused projection remains"); + tx.rollback().await.expect("rollback inspection"); + assert_eq!(stored.event.id, second.id); + assert_eq!( + unfinished_public_projection_retirements( + &db, + &[community], + relay.public_key().as_bytes(), + ) + .await + .expect("count stale work"), + 0 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn populated_upgrade_retires_an_assertion_without_private_head_metadata() { + let (db, community) = setup().await; + let relay = Keys::generate(); + let subject_keys = Keys::generate(); + let subject = subject_keys.public_key(); + resolve_identity_binding( + &db.pool, + community, + &ResolveBindingInput { + issuer: ISSUER, + subject: "legacy-projection-subject", + pubkey: subject.as_bytes(), + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("enroll legacy projection subject"); + let active = assertion(&relay, subject, true, 400); + begin_active_public_projection( + &db, + community, + relay.public_key().as_bytes(), + ISSUER, + "legacy-projection-subject", + subject.as_bytes(), + ) + .await + .expect("begin legacy projection") + .expect("legacy binding is active") + .commit(&active, ProjectionDisposition::Active) + .await + .expect("commit legacy projection"); + sqlx::query( + "DELETE FROM identity_public_projection_heads \ + WHERE community_id=$1 AND relay_pubkey=$2 AND subject_pubkey=$3", + ) + .bind(community.as_uuid()) + .bind(relay.public_key().as_bytes()) + .bind(subject.as_bytes()) + .execute(&db.pool) + .await + .expect("simulate pre-migration projection"); + + let operation_id = Uuid::new_v4(); + revoke_identity_key( + &db.pool, + community, + LifecycleContext { + operation_id, + actor: None, + reason: "retire populated upgrade projection", + }, + subject.as_bytes(), + ) + .await + .expect("commit populated-upgrade revocation"); + materialize_public_projection_retirements(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("materialize populated-upgrade work"); + let permit = begin_public_projection_retirement( + &db, + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim populated-upgrade work") + .expect("populated-upgrade work exists"), + ) + .await + .expect("begin populated-upgrade retirement"); + assert_eq!(permit.head(), None); + let inactive = assertion(&relay, subject, false, 401); + permit + .finish_inactive(&inactive) + .await + .expect("retire populated-upgrade projection"); + let delivery = + begin_public_projection_delivery(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim populated-upgrade delivery") + .expect("populated-upgrade delivery exists"); + assert_eq!(delivery.stored().event.id, inactive.id); + delivery.complete().await.expect("complete delivery"); + + // A crash-restored/pre-head replica can rediscover already-inactive + // work without private ownership metadata. Convergence must not + // relabel that public event to the replayed job's source generation. + sqlx::query( + "DELETE FROM identity_public_projection_heads \ + WHERE community_id=$1 AND relay_pubkey=$2 AND subject_pubkey=$3", + ) + .bind(community.as_uuid()) + .bind(relay.public_key().as_bytes()) + .bind(subject.as_bytes()) + .execute(&db.pool) + .await + .expect("remove restored private head metadata"); + sqlx::query( + "UPDATE identity_public_projection_retirements \ + SET source_binding_id=NULL, source_binding_version=NULL, \ + phase='projection', outcome=NULL, event_id=NULL, \ + completed_at=NULL, next_attempt_at=NOW(), updated_at=NOW() \ + WHERE community_id=$1 AND operation_id=$2 AND relay_pubkey=$3", + ) + .bind(community.as_uuid()) + .bind(operation_id) + .bind(relay.public_key().as_bytes()) + .execute(&db.pool) + .await + .expect("simulate source-less restored retry"); + let replay = begin_public_projection_retirement( + &db, + claim_public_projection_retirement(&db, &[community], relay.public_key().as_bytes()) + .await + .expect("claim restored retry") + .expect("restored retry exists"), + ) + .await + .expect("begin restored retry"); + assert_eq!(replay.head(), None); + assert_eq!( + replay + .current_projection() + .expect("inactive projection remains") + .event + .id, + inactive.id + ); + replay + .finish_existing_inactive() + .await + .expect("source-less inactive retry converges"); + let head_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM identity_public_projection_heads \ + WHERE community_id=$1 AND relay_pubkey=$2 AND subject_pubkey=$3", + ) + .bind(community.as_uuid()) + .bind(relay.public_key().as_bytes()) + .bind(subject.as_bytes()) + .fetch_one(&db.pool) + .await + .expect("count private heads after convergence"); + assert_eq!(head_count, 0, "replay must not manufacture ownership"); + } + + #[test] + fn public_projection_disposition_is_canonical_and_label_free_when_inactive() { + let relay = Keys::generate(); + let subject = Keys::generate().public_key(); + let active = assertion(&relay, subject, true, 500); + let inactive = assertion(&relay, subject, false, 501); + assert!(validate_event_disposition(&active, ProjectionDisposition::Active).is_ok()); + assert!(validate_event_disposition(&inactive, ProjectionDisposition::Inactive).is_ok()); + assert!(validate_event_disposition(&active, ProjectionDisposition::Inactive).is_err()); + assert!(validate_event_disposition(&inactive, ProjectionDisposition::Active).is_err()); + + let subject_hex = subject.to_hex(); + let stale_label = EventBuilder::new(Kind::Custom(ASSERTION_KIND as u16), "") + .tags([ + Tag::parse(["d", subject_hex.as_str()]).expect("d tag"), + Tag::parse(["p", subject_hex.as_str()]).expect("p tag"), + Tag::parse(["verified", "relay"]).expect("verified tag"), + Tag::parse(["active", "false"]).expect("active tag"), + Tag::parse(["expiration", "0"]).expect("expiration tag"), + Tag::parse(["display_name", "Stale label"]).expect("display tag"), + ]) + .custom_created_at(Timestamp::from(502)) + .sign_with_keys(&relay) + .expect("sign malformed projection"); + assert!(validate_event_disposition(&stale_label, ProjectionDisposition::Inactive).is_err()); + + for (active, at) in [(true, 503), (false, 504)] { + let active_value = if active { "true" } else { "false" }; + let expiration = if active { "500" } else { "0" }; + let mut tags = vec![ + Tag::parse(["d", subject_hex.as_str()]).expect("d tag"), + Tag::parse(["p", subject_hex.as_str()]).expect("p tag"), + Tag::parse(["verified", "relay"]).expect("verified tag"), + Tag::parse(["active", active_value]).expect("active tag"), + Tag::parse(["expiration", expiration]).expect("expiration tag"), + ]; + if active { + tags.push( + Tag::parse(["display_name", "Private stale content"]).expect("display tag"), + ); + } + let nonempty = EventBuilder::new( + Kind::Custom(ASSERTION_KIND as u16), + "private content must never be projected", + ) + .tags(tags) + .custom_created_at(Timestamp::from(at)) + .sign_with_keys(&relay) + .expect("sign non-canonical projection"); + let disposition = if active { + ProjectionDisposition::Active + } else { + ProjectionDisposition::Inactive + }; + assert!( + validate_event_disposition(&nonempty, disposition).is_err(), + "signed projection content must be empty" + ); + } + } + + #[test] + fn debug_receipts_redact_binding_and_public_key_coordinates() { + let claim = ProjectionRetirementClaim { + community_id: CommunityId::from_uuid(Uuid::new_v4()), + operation_id: Uuid::new_v4(), + relay_pubkey: vec![3; 32], + old_pubkey: vec![4; 32], + source_origin: Some(ProjectionBindingOrigin { + binding_id: Uuid::new_v4(), + binding_version: 7, + }), + claim_token: Uuid::new_v4(), + }; + let debug = format!("{claim:?}"); + assert!(!debug.contains(&hex::encode(&claim.relay_pubkey))); + assert!(!debug.contains(&hex::encode(&claim.old_pubkey))); + assert!(!debug.contains(&claim.operation_id.to_string())); + assert!(!debug.contains("7")); + } +} diff --git a/crates/buzz-db/src/push.rs b/crates/buzz-db/src/push.rs index 04b6a7ae3..559dbaddc 100644 --- a/crates/buzz-db/src/push.rs +++ b/crates/buzz-db/src/push.rs @@ -821,10 +821,21 @@ pub async fn claim_due_match_batch( limit: i64, lease_until: DateTime, ) -> Result> { - claim_due_match_batch_with_loader( + claim_due_match_batch_excluding(pool, limit, lease_until, &[]).await +} + +/// Claim a matcher batch without touching exact protected Enforce domains. +pub async fn claim_due_match_batch_excluding( + pool: &PgPool, + limit: i64, + lease_until: DateTime, + excluded_communities: &[Uuid], +) -> Result> { + claim_due_match_batch_with_loader_excluding( pool, limit, lease_until, + excluded_communities, |pool, community, ids| async move { let refs: Vec<&[u8]> = ids.iter().map(Vec::as_slice).collect(); crate::event::get_events_by_ids(&pool, community, &refs).await @@ -833,10 +844,11 @@ pub async fn claim_due_match_batch( .await } -async fn claim_due_match_batch_with_loader( +async fn claim_due_match_batch_with_loader_excluding( pool: &PgPool, limit: i64, lease_until: DateTime, + excluded_communities: &[Uuid], load: F, ) -> Result> where @@ -850,6 +862,7 @@ where SELECT community_id FROM push_match_queue WHERE attempts < $3 + AND NOT (community_id = ANY($5::uuid[])) AND next_attempt_at <= now() AND (state = 'pending' OR (state = 'matching' AND lease_until < now())) ORDER BY next_attempt_at, created_at @@ -860,6 +873,7 @@ where FROM push_match_queue q JOIN target t ON q.community_id = t.community_id WHERE q.attempts < $3 + AND NOT (q.community_id = ANY($5::uuid[])) AND q.next_attempt_at <= now() AND (q.state = 'pending' OR (q.state = 'matching' AND q.lease_until < now())) ORDER BY q.next_attempt_at, q.created_at @@ -877,6 +891,7 @@ where .bind(lease_until) .bind(MAX_MATCH_ATTEMPTS) .bind(limit) + .bind(excluded_communities) .fetch_all(pool) .await?; if rows.is_empty() { @@ -931,11 +946,21 @@ where /// served by the due partial index, so putting it in every claim made claims /// slower exactly when a backlog needed them fastest. pub async fn reap_exhausted_matches(pool: &PgPool) -> Result { + reap_exhausted_matches_excluding(pool, &[]).await +} + +/// Reap exhausted matcher jobs outside exact protected Enforce domains. +pub async fn reap_exhausted_matches_excluding( + pool: &PgPool, + excluded_communities: &[Uuid], +) -> Result { Ok(sqlx::query( "DELETE FROM push_match_queue WHERE attempts >= $1 \ + AND NOT (community_id = ANY($2::uuid[])) \ AND (state='pending' OR (state='matching' AND lease_until < now()))", ) .bind(MAX_MATCH_ATTEMPTS) + .bind(excluded_communities) .execute(pool) .await? .rows_affected()) @@ -1891,6 +1916,18 @@ mod tests { .await .expect("read matcher queue"); assert_eq!(queued, vec![9]); + assert!( + claim_due_match_batch_excluding( + &pool, + 16, + Utc::now() + chrono::Duration::minutes(1), + &[*community.as_uuid()], + ) + .await + .expect("excluded protected matcher claim") + .is_none(), + "an excluded domain must remain unclaimed" + ); sqlx::query("UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=$2") .bind(community.as_uuid()) @@ -1927,10 +1964,11 @@ mod tests { .await .expect("insert event"); - let error = claim_due_match_batch_with_loader( + let error = claim_due_match_batch_with_loader_excluding( &pool, 16, Utc::now() - chrono::Duration::seconds(1), + &[], |_pool, _community, _event_ids| async { Err(crate::DbError::InvalidData("injected load failure".into())) }, diff --git a/desktop/src-tauri/src/commands/profile.rs b/desktop/src-tauri/src/commands/profile.rs index c870c8f28..cdc87ee60 100644 --- a/desktop/src-tauri/src/commands/profile.rs +++ b/desktop/src-tauri/src/commands/profile.rs @@ -89,33 +89,58 @@ fn verified_identities( let parts = tag.as_slice(); (parts.len() == 2).then(|| parts[1].as_str()) }; + let canonical_tag_set = |active: bool| { + let allowed: &[&str] = if active { + &["d", "p", "verified", "active", "expiration", "display_name"] + } else { + &["d", "p", "verified", "active", "expiration"] + }; + event.tags.len() == allowed.len() + && event.tags.iter().all(|tag| { + let parts = tag.as_slice(); + parts.len() == 2 + && allowed.contains(&parts[0].as_str()) + && event + .tags + .iter() + .filter(|candidate| { + candidate.as_slice().first() == parts.first() + }) + .count() + == 1 + }) + }; // Select the signed replaceable-event head before validating its // payload. Otherwise a newer malformed assertion could be skipped and // silently resurrect the older active label returned alongside it. - let identity = match (tag_value("d"), tag_value("verified"), tag_value("p")) { - (Some(assertion_d), Some("relay"), Some(asserted_subject)) - if assertion_d == subject && asserted_subject == subject => - { - match tag_value("active") { - Some("false") => None, - Some("true") => match ( - tag_value("expiration") - .and_then(|value| value.parse::().ok()) - .filter(|expiration| *expiration > now), - tag_value("display_name") - .map(str::trim) - .filter(|value| !value.is_empty()), - ) { - (Some(expires_at), Some(display_name)) => Some(VerifiedIdentity { - display_name: display_name.to_string(), - expires_at, - }), + let identity = if event.content.is_empty() { + match (tag_value("d"), tag_value("verified"), tag_value("p")) { + (Some(assertion_d), Some("relay"), Some(asserted_subject)) + if assertion_d == subject && asserted_subject == subject => + { + match tag_value("active") { + Some("false") if canonical_tag_set(false) => None, + Some("true") if canonical_tag_set(true) => match ( + tag_value("expiration") + .and_then(|value| value.parse::().ok()) + .filter(|expiration| *expiration > now), + tag_value("display_name") + .map(str::trim) + .filter(|value| !value.is_empty()), + ) { + (Some(expires_at), Some(display_name)) => Some(VerifiedIdentity { + display_name: display_name.to_string(), + expires_at, + }), + _ => None, + }, _ => None, - }, - _ => None, + } } + _ => None, } - _ => None, + } else { + None }; let created_at = event.created_at.as_secs(); let event_id = event.id.to_hex(); @@ -650,6 +675,86 @@ mod tests { ); } + #[test] + fn newer_nonempty_projection_removes_verified_identity() { + let relay = nostr::Keys::generate(); + let subject = nostr::Keys::generate().public_key().to_hex(); + let created_at = nostr::Timestamp::now().as_secs(); + let expires_at = created_at + 60; + let canonical_tags = || { + [ + nostr::Tag::parse(["d", subject.as_str()]).unwrap(), + nostr::Tag::parse(["p", subject.as_str()]).unwrap(), + nostr::Tag::parse(["verified", "relay"]).unwrap(), + nostr::Tag::parse(["active", "true"]).unwrap(), + nostr::Tag::parse(["expiration", &expires_at.to_string()]).unwrap(), + nostr::Tag::parse(["display_name", "Example User"]).unwrap(), + ] + }; + let active = + nostr::EventBuilder::new(nostr::Kind::Custom(KIND_USER_TRUSTED_ASSERTION as u16), "") + .tags(canonical_tags()) + .custom_created_at(nostr::Timestamp::from(created_at)) + .sign_with_keys(&relay) + .unwrap(); + let nonempty = nostr::EventBuilder::new( + nostr::Kind::Custom(KIND_USER_TRUSTED_ASSERTION as u16), + "private content must never be projected", + ) + .tags(canonical_tags()) + .custom_created_at(nostr::Timestamp::from(created_at + 1)) + .sign_with_keys(&relay) + .unwrap(); + + assert!( + verified_identities(&[active, nonempty], Some(&relay.public_key().to_hex())).is_empty(), + "a malformed newer head must withdraw rather than reveal or resurrect a label" + ); + } + + #[test] + fn newer_projection_with_unknown_or_duplicate_tags_removes_verified_identity() { + let relay = nostr::Keys::generate(); + let subject = nostr::Keys::generate().public_key().to_hex(); + let created_at = nostr::Timestamp::now().as_secs(); + let expires_at = created_at + 60; + let canonical = nostr::EventBuilder::new( + nostr::Kind::Custom(KIND_USER_TRUSTED_ASSERTION as u16), + "", + ) + .tags([ + nostr::Tag::parse(["d", subject.as_str()]).unwrap(), + nostr::Tag::parse(["p", subject.as_str()]).unwrap(), + nostr::Tag::parse(["verified", "relay"]).unwrap(), + nostr::Tag::parse(["active", "true"]).unwrap(), + nostr::Tag::parse(["expiration", &expires_at.to_string()]).unwrap(), + nostr::Tag::parse(["display_name", "Example User"]).unwrap(), + ]) + .custom_created_at(nostr::Timestamp::from(created_at)) + .sign_with_keys(&relay) + .unwrap(); + + for extra in [ + nostr::Tag::parse(["issuer", "private.invalid"]).unwrap(), + nostr::Tag::parse(["display_name", "Replacement"]).unwrap(), + ] { + let mut tags = canonical.tags.clone().to_vec(); + tags.push(extra); + let malformed = nostr::EventBuilder::new( + nostr::Kind::Custom(KIND_USER_TRUSTED_ASSERTION as u16), + "", + ) + .tags(tags) + .custom_created_at(nostr::Timestamp::from(created_at + 1)) + .sign_with_keys(&relay) + .unwrap(); + assert!( + verified_identities(&[canonical.clone(), malformed], Some(&relay.public_key().to_hex())) + .is_empty() + ); + } + } + #[test] fn newer_malformed_assertion_does_not_resurrect_older_identity() { let relay = nostr::Keys::generate(); diff --git a/migrations/0038_client_status_fanout_withdrawals.sql b/migrations/0038_client_status_fanout_withdrawals.sql new file mode 100644 index 000000000..eb0371cca --- /dev/null +++ b/migrations/0038_client_status_fanout_withdrawals.sql @@ -0,0 +1,21 @@ +-- Retain the exact current revision superseded by a withdrawal so every +-- authenticated connection that displayed that author can receive a +-- strictly newer opaque withdrawal. This remains server-side reconciliation +-- state and is never serialized into the client projection. + +ALTER TABLE client_status_revisions + ADD COLUMN supersedes_revision BIGINT; + +UPDATE client_status_revisions +SET supersedes_revision = revision - 1 +WHERE disposition = 'withdrawn'; + +ALTER TABLE client_status_revisions + ADD CONSTRAINT client_status_revisions_withdrawal CHECK ( + (disposition = 'current' AND supersedes_revision IS NULL) + OR + (disposition = 'withdrawn' + AND supersedes_revision IS NOT NULL + AND supersedes_revision > 0 + AND revision > supersedes_revision) + ); diff --git a/migrations/0044_identity_public_projection_retirement.sql b/migrations/0044_identity_public_projection_retirement.sql new file mode 100644 index 000000000..e6182bd20 --- /dev/null +++ b/migrations/0044_identity_public_projection_retirement.sql @@ -0,0 +1,80 @@ +-- Durable, provider-neutral reconciliation for the optional public identity +-- projection. O3 lifecycle rows remain the authority; these tables contain +-- only public event coordinates and opaque binding generations. + +CREATE TABLE identity_public_projection_heads ( + community_id UUID NOT NULL REFERENCES communities(id), + relay_pubkey BYTEA NOT NULL, + subject_pubkey BYTEA NOT NULL, + event_id BYTEA NOT NULL, + event_created_at TIMESTAMPTZ NOT NULL, + disposition TEXT NOT NULL, + source_binding_id UUID, + source_binding_version BIGINT, + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (community_id, relay_pubkey, subject_pubkey), + FOREIGN KEY (community_id, source_binding_id) + REFERENCES identity_bindings (community_id, binding_id), + CHECK (length(relay_pubkey) = 32), + CHECK (length(subject_pubkey) = 32), + CHECK (length(event_id) = 32), + CHECK (disposition IN ('active', 'inactive')), + CHECK (source_binding_version IS NULL OR source_binding_version > 0), + CHECK ( + (source_binding_id IS NULL AND source_binding_version IS NULL) + OR + (source_binding_id IS NOT NULL AND source_binding_version IS NOT NULL) + ) +); + +CREATE TABLE identity_public_projection_retirements ( + community_id UUID NOT NULL REFERENCES communities(id), + operation_id UUID NOT NULL, + relay_pubkey BYTEA NOT NULL, + old_pubkey BYTEA NOT NULL, + source_binding_id UUID, + source_binding_version BIGINT, + operation_kind TEXT NOT NULL, + phase TEXT NOT NULL DEFAULT 'projection', + outcome TEXT, + event_id BYTEA, + attempts BIGINT NOT NULL DEFAULT 0, + next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + claim_token UUID, + lease_until TIMESTAMPTZ, + completed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (community_id, operation_id, relay_pubkey), + FOREIGN KEY (community_id, operation_id) + REFERENCES identity_lifecycle_operations (community_id, operation_id), + FOREIGN KEY (community_id, source_binding_id) + REFERENCES identity_bindings (community_id, binding_id), + CHECK (length(relay_pubkey) = 32), + CHECK (length(old_pubkey) = 32), + CHECK (source_binding_version IS NULL OR source_binding_version > 0), + CHECK ( + (source_binding_id IS NULL AND source_binding_version IS NULL) + OR + (source_binding_id IS NOT NULL AND source_binding_version IS NOT NULL) + ), + CHECK (operation_kind IN ('revoke_key', 'rotate')), + CHECK (phase IN ('projection', 'delivery', 'completed', 'superseded')), + CHECK (outcome IS NULL OR outcome IN ( + 'no_projection', 'already_inactive', 'replaced_inactive', + 'newer_binding', 'newer_projection' + )), + CHECK (event_id IS NULL OR length(event_id) = 32), + CHECK (attempts >= 0), + CHECK ((claim_token IS NULL) = (lease_until IS NULL)), + CHECK ( + (phase IN ('completed', 'superseded') AND completed_at IS NOT NULL) + OR + (phase IN ('projection', 'delivery') AND completed_at IS NULL) + ) +); + +CREATE INDEX idx_identity_public_projection_retirements_ready + ON identity_public_projection_retirements + (phase, next_attempt_at, community_id, operation_id) + WHERE phase IN ('projection', 'delivery');