Merge audit community-chain (per-community hash chain on frozen audit_log DDL)

* commit 'ba11d6663630ff020bf54f79a185932e26085d32':
  feat(audit): per-community hash chain on the frozen audit_log DDL

Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d
2026-06-26 17:16:49 -04:00
co-authored by Tyler Longwell
6 changed files with 430 additions and 290 deletions
+46 -32
View File
@@ -4,46 +4,60 @@ use uuid::Uuid;
use crate::action::AuditAction;
/// Materialised audit log entry as stored in the DB.
#[derive(Debug, Clone, Serialize, Deserialize)]
/// A materialised audit log entry as stored in `audit_log`.
///
/// Rows are keyed `(community_id, seq)`: `seq` is monotonic *within one
/// community*, and `prev_hash` chains to the previous entry *of the same
/// community*. The chain is independent per tenant.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AuditEntry {
/// Monotonically increasing sequence number.
/// Server-resolved community this entry belongs to. Leads the primary key.
pub community_id: Uuid,
/// Sequence number, monotonic within `community_id` (starts at 1).
pub seq: i64,
/// When the entry was recorded.
pub timestamp: DateTime<Utc>,
/// Nostr event ID that triggered this action.
pub event_id: String,
/// Nostr event kind number.
pub event_kind: u32,
/// Hex-encoded Nostr pubkey.
pub actor_pubkey: String,
/// SHA-256 of this entry's fields including `community_id` and `prev_hash`.
pub hash: Vec<u8>,
/// SHA-256 of the previous entry in *this community's* chain, or `None` for
/// the community's first entry (hashed as [`crate::hash::GENESIS_HASH`]).
pub prev_hash: Option<Vec<u8>>,
/// Action that was performed.
pub action: AuditAction,
/// Channel this action applies to, if any.
pub channel_id: Option<Uuid>,
/// Arbitrary JSON context. **Included in hash computation** (serialized with
/// sorted keys for determinism) so that metadata tampering is detectable.
pub metadata: serde_json::Value,
/// SHA-256 hex hash of the previous entry (or [`crate::hash::GENESIS_HASH`] for the first).
pub prev_hash: String,
/// SHA-256 hex hash of this entry's fields including `prev_hash`.
pub hash: String,
/// Raw bytes of the actor's Nostr pubkey, if the action has one.
pub actor_pubkey: Option<Vec<u8>>,
/// Generic identifier of the object acted upon (event id hex, channel UUID,
/// media sha256, …), if any. The relay resolves it under `community_id`;
/// it never names an object in another community.
pub object_id: Option<String>,
/// Arbitrary JSON context. **Included in the hash** (serialized with sorted
/// keys for determinism) so tampering with it is detectable.
pub detail: serde_json::Value,
/// When the entry was recorded.
pub created_at: DateTime<Utc>,
}
/// Input for creating a new audit entry. `seq`, `prev_hash`, `hash` are computed by `AuditService::log`.
#[derive(Debug, Clone, Serialize, Deserialize)]
/// Input for appending a new audit entry. `seq`, `prev_hash`, `hash`, and
/// `created_at` are assigned by [`crate::service::AuditService::log`].
///
/// `community_id` is the **server-resolved** tenant (from the request's
/// `TenantContext`), never a client-supplied value — the same provenance rule
/// the whole multi-tenant model rests on.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct NewAuditEntry {
/// Nostr event ID that triggered this action.
pub event_id: String,
/// Must not be 22242 (NIP-42 AUTH).
pub event_kind: u32,
/// Hex-encoded Nostr pubkey of the actor.
pub actor_pubkey: String,
/// Server-resolved community this entry belongs to.
pub community_id: Uuid,
/// Action that was performed.
pub action: AuditAction,
/// Channel this action applies to, if any.
pub channel_id: Option<Uuid>,
/// Arbitrary JSON context included in hash computation.
/// Raw bytes of the actor's Nostr pubkey, if the action has one.
pub actor_pubkey: Option<Vec<u8>>,
/// Generic identifier of the object acted upon, if any.
pub object_id: Option<String>,
/// Arbitrary JSON context included in the hash.
///
/// **Never bearer-token material.** This field is opaque to the audit
/// crate and persisted verbatim; callers must not write tokens, passwords,
/// or other secrets here. `AuthSuccess`/`AuthFailure` entries carry only
/// outcome metadata — the token has no slot in this type, and `detail` must
/// not become one.
#[serde(default)]
pub metadata: serde_json::Value,
pub detail: serde_json::Value,
}
+16 -20
View File
@@ -1,45 +1,41 @@
use thiserror::Error;
/// Errors that can occur during audit log operations.
///
/// These are **operator-internal** diagnostics (logged by the audit worker, or
/// returned to an operator-scoped verification call) — they are never relayed to
/// a client on the wire. Even so, no variant embeds a `community_id` or any
/// cross-community object identifier: a `seq` is per-community and meaningless
/// without its chain, and hashes are opaque. An error raised while verifying
/// community A's chain therefore cannot reveal a fact about community B.
#[derive(Debug, Error)]
pub enum AuditError {
/// A database operation failed.
#[error("database error: {0}")]
Database(#[from] sqlx::Error),
/// Attempted to log a NIP-42 AUTH event (kind 22242), which is forbidden.
#[error("auth events (kind 22242) must never appear in the audit log")]
AuthEventForbidden,
/// The `prev_hash` of an entry does not match the hash of the preceding entry.
/// The `prev_hash` of an entry does not match the hash of the preceding
/// entry in the same community's chain.
#[error(
"hash chain integrity violation at seq {seq}: expected prev_hash {expected}, got {actual}"
"hash chain integrity violation at seq {seq}: prev_hash does not match preceding entry"
)]
ChainViolation {
/// Sequence number of the offending entry.
/// Per-community sequence number of the offending entry.
seq: i64,
/// Hash that was expected based on the previous entry.
expected: String,
/// Hash that was actually found in the entry.
actual: String,
},
/// The stored hash of an entry does not match the recomputed hash.
#[error("hash mismatch at seq {seq}: stored {stored}, computed {computed}")]
#[error("hash mismatch at seq {seq}: stored hash does not match recomputed hash")]
HashMismatch {
/// Sequence number of the offending entry.
/// Per-community sequence number of the offending entry.
seq: i64,
/// Hash value stored in the database.
stored: String,
/// Hash value recomputed from the entry fields.
computed: String,
},
/// An unrecognised action string was found in the database.
#[error("unknown audit action in DB: {0:?}")]
UnknownAction(String),
#[error("unknown audit action in database")]
UnknownAction,
/// A JSON serialization error occurred (e.g. while canonicalising metadata).
/// A JSON serialization error occurred (e.g. while canonicalising `detail`).
#[error("serialization error: {0}")]
Serialization(#[from] serde_json::Error),
}
+79 -50
View File
@@ -3,41 +3,53 @@ use sha2::{Digest, Sha256};
use crate::entry::AuditEntry;
use crate::error::AuditError;
/// Sentinel `prev_hash` value used for the first entry in the chain.
pub const GENESIS_HASH: &str = "0000000000000000000000000000000000000000000000000000000000000000";
/// The 32-byte sentinel hashed in place of `prev_hash` for a community's first
/// entry. Stored as `prev_hash = NULL`; hashed as all-zero bytes.
pub const GENESIS_HASH: [u8; 32] = [0u8; 32];
/// SHA-256 over all identity, chain, and context fields.
/// Field order is fixed — changing it invalidates all existing chains.
/// SHA-256 over the entry's identity, chain, and context fields.
///
/// Metadata is serialized via `BTreeMap` to guarantee key ordering across
/// machines and Rust versions. `serde_json::Value` does not guarantee order.
/// Field order is fixed — changing it invalidates all existing chains. The
/// `community_id` is hashed first so chain identity carries the tenant: an entry
/// cannot be lifted out of one community's chain and re-verified inside another.
///
/// Returns `Err(AuditError::Serialization)` if metadata cannot be serialized.
/// Never hashes a default/empty value as a stand-in for a real payload —
/// a serialization failure is a hard error, not a silent degradation.
pub fn compute_hash(entry: &AuditEntry) -> Result<String, AuditError> {
/// `detail` is serialized via [`canonical_json`] (sorted keys) so the hash is
/// stable across machines and Rust versions. A serialization failure is a hard
/// error, never silently hashed as empty.
pub fn compute_hash(entry: &AuditEntry) -> Result<[u8; 32], AuditError> {
let mut hasher = Sha256::new();
// Tenant binding: community_id leads the hash.
hasher.update(entry.community_id.as_bytes());
hasher.update(entry.seq.to_be_bytes());
hasher.update(entry.timestamp.to_rfc3339().as_bytes());
hasher.update(entry.event_id.as_bytes());
// event_kind is u32 — 4 bytes in big-endian for the hash chain.
hasher.update(entry.event_kind.to_be_bytes());
hasher.update(entry.actor_pubkey.as_bytes());
hasher.update(entry.created_at.to_rfc3339().as_bytes());
hasher.update(entry.action.as_str().as_bytes());
match &entry.channel_id {
Some(id) => hasher.update(id.as_bytes()),
None => hasher.update([0u8; 16]),
match &entry.actor_pubkey {
Some(pk) => {
hasher.update([1u8]); // presence tag — distinguishes Some(empty) from None
hasher.update(pk);
}
None => hasher.update([0u8]),
}
hasher.update(canonical_json(&entry.metadata)?.as_bytes());
hasher.update(entry.prev_hash.as_bytes());
Ok(hex::encode(hasher.finalize()))
match &entry.object_id {
Some(id) => {
hasher.update([1u8]);
hasher.update(id.as_bytes());
}
None => hasher.update([0u8]),
}
hasher.update(canonical_json(&entry.detail)?.as_bytes());
match &entry.prev_hash {
Some(h) => hasher.update(h),
None => hasher.update(GENESIS_HASH),
}
Ok(hasher.finalize().into())
}
/// Serialize a JSON value with sorted keys for deterministic output.
/// Serialize a JSON value with sorted object keys for deterministic output.
///
/// Returns `Err` if any scalar value cannot be serialized. This should never
/// happen for well-formed `serde_json::Value`, but we propagate rather than
/// silently substitute an empty string.
/// Propagates any scalar serialization error rather than substituting a
/// placeholder — a hash must never silently stand in an empty value for a real
/// payload.
fn canonical_json(value: &serde_json::Value) -> Result<String, serde_json::Error> {
use serde_json::Value;
use std::collections::BTreeMap;
@@ -81,21 +93,21 @@ mod tests {
use super::*;
use crate::{action::AuditAction, entry::AuditEntry};
use chrono::Utc;
use uuid::Uuid;
fn sample_entry() -> AuditEntry {
AuditEntry {
community_id: Uuid::from_u128(1),
seq: 1,
timestamp: chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00Z")
hash: Vec::new(),
prev_hash: None,
action: AuditAction::EventCreated,
actor_pubkey: Some(vec![0xab; 32]),
object_id: Some("abc123".into()),
detail: serde_json::Value::Null,
created_at: chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00Z")
.unwrap()
.with_timezone(&Utc),
event_id: "abc123".to_string(),
event_kind: 1,
actor_pubkey: "pubkey_alice".to_string(),
action: AuditAction::EventCreated,
channel_id: None,
metadata: serde_json::Value::Null,
prev_hash: GENESIS_HASH.to_string(),
hash: String::new(),
}
}
@@ -103,7 +115,17 @@ mod tests {
fn deterministic() {
let entry = sample_entry();
assert_eq!(compute_hash(&entry).unwrap(), compute_hash(&entry).unwrap());
assert_eq!(compute_hash(&entry).unwrap().len(), 64);
assert_eq!(compute_hash(&entry).unwrap().len(), 32);
}
#[test]
fn community_id_is_part_of_identity() {
// The whole point: the same logical entry in two communities hashes
// differently, so a row can't be replayed across chains.
let a = sample_entry();
let mut b = a.clone();
b.community_id = Uuid::from_u128(2);
assert_ne!(compute_hash(&a).unwrap(), compute_hash(&b).unwrap());
}
#[test]
@@ -111,38 +133,45 @@ mod tests {
let base = sample_entry();
let h0 = compute_hash(&base).unwrap();
let mut e = base.clone();
e.event_id = "different_event".into();
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base.clone();
e.seq = 2;
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base.clone();
e.actor_pubkey = "pubkey_bob".into();
e.action = AuditAction::EventDeleted;
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base.clone();
e.channel_id = Some(uuid::Uuid::new_v4());
e.actor_pubkey = Some(vec![0xcd; 32]);
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base;
e.metadata = serde_json::json!({"key": "value"});
let mut e = base.clone();
e.object_id = Some("different".into());
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base.clone();
e.detail = serde_json::json!({"key": "value"});
assert_ne!(h0, compute_hash(&e).unwrap());
let mut e = base.clone();
e.prev_hash = Some(vec![0xff; 32]);
assert_ne!(h0, compute_hash(&e).unwrap());
}
#[test]
fn presence_tag_distinguishes_none_from_empty() {
// Some(empty) must not collide with None — the presence tag prevents it.
let mut none = sample_entry();
none.actor_pubkey = None;
let mut empty = sample_entry();
empty.actor_pubkey = Some(Vec::new());
assert_ne!(compute_hash(&none).unwrap(), compute_hash(&empty).unwrap());
}
#[test]
fn canonical_json_key_order_is_stable() {
// Same keys in different insertion order must produce the same hash.
let a = serde_json::json!({"z": 1, "a": 2, "m": 3});
let b = serde_json::json!({"a": 2, "m": 3, "z": 1});
assert_eq!(canonical_json(&a).unwrap(), canonical_json(&b).unwrap());
}
#[test]
fn genesis_hash_format() {
assert_eq!(GENESIS_HASH.len(), 64);
assert!(GENESIS_HASH.chars().all(|c| c == '0'));
}
}
+16 -6
View File
@@ -1,8 +1,21 @@
#![deny(unsafe_code)]
#![warn(missing_docs)]
//! Tamper-evident hash-chain audit log. Each entry chains to the previous via
//! SHA-256. Single-writer via Postgres `pg_advisory_lock`. AUTH events (kind 22242)
//! are rejected — they carry bearer tokens.
//! Tamper-evident, **per-community** hash-chain audit log.
//!
//! Each community owns an independent chain: rows are keyed `(community_id, seq)`,
//! `seq` is monotonic *within a community*, and each entry chains to the previous
//! entry *of the same community* via SHA-256. The `community_id` is folded into the
//! hash, so a row lifted out of one community's chain can never verify inside
//! another's — chain identity carries the tenant. This is the audit half of the
//! non-interference floor (`auditHeads[c]` in `MultiTenantRelay.tla`): an audit
//! observation reveals only its own community's head.
//!
//! Writes for a given community are serialized by a **per-community** Postgres
//! advisory lock, so the chain stays consistent across relay processes without one
//! global lock serializing (and timing-coupling) every tenant.
//!
//! The `audit_log` table is owned by the consolidated `0001` migration — this crate
//! is pure chain logic and ships no DDL.
/// Audit action types recorded in the log.
pub mod action;
@@ -12,8 +25,6 @@ pub mod entry;
pub mod error;
/// SHA-256 hash computation for audit entries.
pub mod hash;
/// SQL schema for the audit log table.
pub mod schema;
/// Audit log service — append and verify entries.
pub mod service;
@@ -21,5 +32,4 @@ pub use action::AuditAction;
pub use entry::{AuditEntry, NewAuditEntry};
pub use error::AuditError;
pub use hash::{compute_hash, GENESIS_HASH};
pub use schema::AUDIT_SCHEMA_SQL;
pub use service::AuditService;
-18
View File
@@ -1,18 +0,0 @@
/// DDL for the `audit_log` table. Passed to [`sqlx::raw_sql`] on startup.
pub const AUDIT_SCHEMA_SQL: &str = r#"
CREATE TABLE IF NOT EXISTS audit_log (
seq BIGINT NOT NULL PRIMARY KEY,
timestamp TIMESTAMPTZ NOT NULL DEFAULT NOW(),
event_id VARCHAR(255) NOT NULL,
event_kind INT NOT NULL,
actor_pubkey VARCHAR(255) NOT NULL,
action VARCHAR(64) NOT NULL,
channel_id BYTEA,
metadata JSONB NOT NULL,
prev_hash VARCHAR(64) NOT NULL,
hash VARCHAR(64) NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_audit_log_timestamp ON audit_log (timestamp);
CREATE INDEX IF NOT EXISTS idx_audit_log_actor ON audit_log (actor_pubkey);
CREATE INDEX IF NOT EXISTS idx_audit_log_channel ON audit_log (channel_id);
"#;
+273 -164
View File
@@ -2,24 +2,29 @@ use chrono::{DateTime, Utc};
use futures_util::FutureExt as _;
use sqlx::{Acquire, PgPool, Row};
use tracing::{debug, instrument, warn};
use uuid::Uuid;
use buzz_core::kind::KIND_AUTH;
use buzz_core::CommunityId;
use crate::{
action::AuditAction,
entry::{AuditEntry, NewAuditEntry},
error::AuditError,
hash::{compute_hash, GENESIS_HASH},
schema::AUDIT_SCHEMA_SQL,
hash::compute_hash,
};
/// Advisory lock key derived from a stable hash of "buzz_audit".
const AUDIT_LOCK_KEY: i64 = 0x5370_7275_7441_7564; // "SprutAud" as hex
/// Per-community advisory lock key. Derived in Postgres from the community UUID
/// so two communities never serialize each other's audit writes (which would be
/// both a throughput bottleneck and a cross-tenant timing oracle). The lock is
/// taken with `pg_advisory_lock(hashtextextended(...))` — see [`AuditService::log`].
const AUDIT_LOCK_NAMESPACE: &str = "buzz_audit:";
/// Append-only audit log service backed by Postgres.
/// Append-only, per-community hash-chain audit log backed by Postgres.
///
/// Serialises writes via `pg_advisory_lock` so the hash chain remains consistent
/// even when multiple relay processes share the same database.
/// Each community has an independent chain keyed `(community_id, seq)`. Writes
/// for one community are serialized by a per-community advisory lock so the chain
/// stays consistent across relay processes; different communities proceed in
/// parallel.
pub struct AuditService {
pool: PgPool,
}
@@ -30,40 +35,32 @@ impl AuditService {
Self { pool }
}
/// Idempotent — safe to call on every startup.
pub async fn ensure_schema(&self) -> Result<(), AuditError> {
sqlx::raw_sql(AUDIT_SCHEMA_SQL).execute(&self.pool).await?;
Ok(())
}
/// Append a new entry to the audit log. Single-writer via `pg_advisory_lock`.
/// Append a new entry to the calling community's chain.
///
/// Postgres advisory locks are session-scoped, so we acquire before the
/// transaction and release after commit (or on any error path).
/// Serialized per-community via `pg_advisory_lock`. Postgres advisory locks
/// are session-scoped, so we acquire before the transaction and release
/// after commit (or on any error path).
#[instrument(skip(self, entry), fields(action = %entry.action))]
pub async fn log(&self, entry: NewAuditEntry) -> Result<AuditEntry, AuditError> {
if entry.event_kind == KIND_AUTH {
warn!("rejected attempt to audit AUTH event (kind 22242)");
return Err(AuditError::AuthEventForbidden);
}
let mut conn = self.pool.acquire().await?;
// Acquire session-level advisory lock (blocks until available).
sqlx::query("SELECT pg_advisory_lock($1)")
.bind(AUDIT_LOCK_KEY)
// Per-community advisory lock: hash the namespaced community id to an
// i64 lock key inside Postgres. Communities lock independently.
let lock_key = format!("{AUDIT_LOCK_NAMESPACE}{}", entry.community_id);
sqlx::query("SELECT pg_advisory_lock(hashtextextended($1, 0))")
.bind(&lock_key)
.execute(&mut *conn)
.await?;
// Run log_inner and release the lock regardless of outcome.
// We use catch_unwind to handle panics so the lock is always released
// before the connection is returned to the pool.
// Run the chain append and release the lock regardless of outcome.
// catch_unwind so a panic still releases the lock before the connection
// returns to the pool.
let result = std::panic::AssertUnwindSafe(self.log_inner(&mut conn, entry))
.catch_unwind()
.await;
let _ = sqlx::query("SELECT pg_advisory_unlock($1)")
.bind(AUDIT_LOCK_KEY)
let _ = sqlx::query("SELECT pg_advisory_unlock(hashtextextended($1, 0))")
.bind(&lock_key)
.execute(&mut *conn)
.await;
@@ -80,56 +77,59 @@ impl AuditService {
) -> Result<AuditEntry, AuditError> {
let mut tx = conn.begin().await?;
let prev_hash: String = sqlx::query("SELECT hash FROM audit_log ORDER BY seq DESC LIMIT 1")
.fetch_optional(&mut *tx)
.await?
.map(|row| row.get::<String, _>("hash"))
.unwrap_or_else(|| GENESIS_HASH.to_string());
// Head of THIS community's chain — scoped by community_id.
let head = sqlx::query(
"SELECT seq, hash FROM audit_log
WHERE community_id = $1
ORDER BY seq DESC LIMIT 1",
)
.bind(entry.community_id)
.fetch_optional(&mut *tx)
.await?;
let seq: i64 =
sqlx::query_scalar("SELECT COALESCE(MAX(seq), 0) + 1 AS next_seq FROM audit_log")
.fetch_one(&mut *tx)
.await?;
let (prev_seq, prev_hash): (i64, Option<Vec<u8>>) = match head {
Some(row) => (
row.get::<i64, _>("seq"),
Some(row.get::<Vec<u8>, _>("hash")),
),
None => (0, None), // community's first entry
};
let seq = prev_seq + 1;
let timestamp: DateTime<Utc> = Utc::now();
let channel_id_bytes: Option<Vec<u8>> = entry.channel_id.map(|u| u.as_bytes().to_vec());
let created_at: DateTime<Utc> = Utc::now();
let mut audit_entry = AuditEntry {
community_id: entry.community_id,
seq,
timestamp,
event_id: entry.event_id,
event_kind: entry.event_kind,
actor_pubkey: entry.actor_pubkey,
action: entry.action,
channel_id: entry.channel_id,
metadata: entry.metadata,
hash: Vec::new(),
prev_hash,
hash: String::new(),
action: entry.action,
actor_pubkey: entry.actor_pubkey,
object_id: entry.object_id,
detail: entry.detail,
created_at,
};
audit_entry.hash = compute_hash(&audit_entry)?;
audit_entry.hash = compute_hash(&audit_entry)?.to_vec();
debug!(seq, hash = %audit_entry.hash, "writing audit entry");
debug!(seq, "writing audit entry");
sqlx::query(
r#"
INSERT INTO audit_log
(seq, timestamp, event_id, event_kind, actor_pubkey, action,
channel_id, metadata, prev_hash, hash)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
(community_id, seq, hash, prev_hash, action, actor_pubkey, object_id, detail, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
"#,
)
.bind(audit_entry.community_id)
.bind(audit_entry.seq)
.bind(audit_entry.timestamp)
.bind(&audit_entry.event_id)
.bind(audit_entry.event_kind as i32)
.bind(&audit_entry.actor_pubkey)
.bind(audit_entry.action.as_str())
.bind(channel_id_bytes)
.bind(&audit_entry.metadata)
.bind(&audit_entry.prev_hash)
.bind(&audit_entry.hash)
.bind(audit_entry.prev_hash.as_deref())
.bind(audit_entry.action.as_str())
.bind(audit_entry.actor_pubkey.as_deref())
.bind(audit_entry.object_id.as_deref())
.bind(&audit_entry.detail)
.bind(audit_entry.created_at)
.execute(&mut *tx)
.await?;
@@ -138,19 +138,28 @@ impl AuditService {
Ok(audit_entry)
}
/// Verify the hash chain for `[from_seq, to_seq]`.
/// Returns `Ok(false)` if range is empty, `Ok(true)` if valid.
/// Verify the hash chain for one community over `[from_seq, to_seq]`.
///
/// Reads exactly that community's chain — it can never observe another
/// community's entries or head. Returns `Ok(false)` if the range is empty,
/// `Ok(true)` if the segment is internally consistent.
#[instrument(skip(self))]
pub async fn verify_chain(&self, from_seq: i64, to_seq: i64) -> Result<bool, AuditError> {
pub async fn verify_chain(
&self,
community: CommunityId,
from_seq: i64,
to_seq: i64,
) -> Result<bool, AuditError> {
let rows = sqlx::query(
r#"
SELECT seq, timestamp, event_id, event_kind, actor_pubkey,
action, channel_id, metadata, prev_hash, hash
SELECT community_id, seq, hash, prev_hash, action, actor_pubkey,
object_id, detail, created_at
FROM audit_log
WHERE seq BETWEEN $1 AND $2
WHERE community_id = $1 AND seq BETWEEN $2 AND $3
ORDER BY seq ASC
"#,
)
.bind(community.as_uuid())
.bind(from_seq)
.bind(to_seq)
.fetch_all(&self.pool)
@@ -160,30 +169,21 @@ impl AuditService {
return Ok(false);
}
let mut expected_prev: Option<String> = None;
let mut expected_prev: Option<Vec<u8>> = None;
for row in &rows {
let entry = row_to_audit_entry(row)?;
let prev_hash = entry.prev_hash.clone();
let stored_hash = entry.hash.clone();
if let Some(ref expected) = expected_prev {
if &prev_hash != expected {
return Err(AuditError::ChainViolation {
seq: entry.seq,
expected: expected.clone(),
actual: prev_hash,
});
// The previous entry's hash must equal this entry's prev_hash.
if entry.prev_hash.as_deref() != Some(expected.as_slice()) {
return Err(AuditError::ChainViolation { seq: entry.seq });
}
}
let computed = compute_hash(&entry)?;
if computed != stored_hash {
return Err(AuditError::HashMismatch {
seq: entry.seq,
stored: stored_hash,
computed,
});
if computed.as_slice() != entry.hash.as_slice() {
return Err(AuditError::HashMismatch { seq: entry.seq });
}
expected_prev = Some(entry.hash);
@@ -192,23 +192,27 @@ impl AuditService {
Ok(true)
}
/// Returns up to `limit` entries starting at `from_seq`, ordered by sequence number.
/// Returns up to `limit` entries from one community's chain starting at
/// `from_seq`, ordered by sequence number. Scoped to `community` — never
/// returns another community's rows.
#[instrument(skip(self))]
pub async fn get_entries(
&self,
community: CommunityId,
from_seq: i64,
limit: i64,
) -> Result<Vec<AuditEntry>, AuditError> {
let rows = sqlx::query(
r#"
SELECT seq, timestamp, event_id, event_kind, actor_pubkey,
action, channel_id, metadata, prev_hash, hash
SELECT community_id, seq, hash, prev_hash, action, actor_pubkey,
object_id, detail, created_at
FROM audit_log
WHERE seq >= $1
WHERE community_id = $1 AND seq >= $2
ORDER BY seq ASC
LIMIT $2
LIMIT $3
"#,
)
.bind(community.as_uuid())
.bind(from_seq)
.bind(limit)
.fetch_all(&self.pool)
@@ -219,34 +223,22 @@ impl AuditService {
}
fn row_to_audit_entry(row: &sqlx::postgres::PgRow) -> Result<AuditEntry, AuditError> {
let seq: i64 = row.get("seq");
let action_str: String = row.get("action");
let action: AuditAction = action_str.parse().map_err(|_| {
warn!(seq, action = %action_str, "unknown action in audit log");
AuditError::UnknownAction(action_str.clone())
})?;
let channel_id_bytes: Option<Vec<u8>> = row.get("channel_id");
let channel_id = channel_id_bytes.and_then(|b| b.try_into().ok().map(uuid::Uuid::from_bytes));
let raw_kind: i32 = row.get("event_kind");
let event_kind = u32::try_from(raw_kind).map_err(|_| {
AuditError::Database(sqlx::Error::Protocol(format!(
"event_kind {raw_kind} out of u32 range at seq {seq}"
)))
warn!("unknown action in audit log");
AuditError::UnknownAction
})?;
Ok(AuditEntry {
seq,
timestamp: row.get("timestamp"),
event_id: row.get("event_id"),
event_kind,
actor_pubkey: row.get("actor_pubkey"),
action,
channel_id,
metadata: row.get("metadata"),
prev_hash: row.get("prev_hash"),
community_id: row.get::<Uuid, _>("community_id"),
seq: row.get("seq"),
hash: row.get("hash"),
prev_hash: row.get("prev_hash"),
action,
actor_pubkey: row.get("actor_pubkey"),
object_id: row.get("object_id"),
detail: row.get("detail"),
created_at: row.get("created_at"),
})
}
@@ -255,10 +247,12 @@ mod tests {
use super::*;
use crate::action::AuditAction;
use crate::entry::NewAuditEntry;
use crate::hash::GENESIS_HASH;
use std::sync::OnceLock;
use tokio::sync::Mutex;
use uuid::Uuid;
// The per-community advisory lock means different communities don't contend,
// but tests share one table; serialize them so seq assertions are stable.
static DB_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
fn db_lock() -> &'static Mutex<()> {
DB_LOCK.get_or_init(|| Mutex::new(()))
@@ -270,122 +264,237 @@ mod tests {
PgPool::connect(&url).await.ok()
}
fn sample_new_entry(kind: u32, action: AuditAction) -> NewAuditEntry {
/// A `community_id` known to exist in `communities` (FK target). Inserts a
/// throwaway community row with a unique host and returns its id.
async fn make_community(pool: &PgPool) -> Uuid {
let id = Uuid::new_v4();
let host = format!("test-{id}.example");
sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)")
.bind(id)
.bind(host)
.execute(pool)
.await
.expect("insert test community");
id
}
fn new_entry(community_id: Uuid, action: AuditAction) -> NewAuditEntry {
NewAuditEntry {
event_id: format!("evt_{}", uuid::Uuid::new_v4()),
event_kind: kind,
actor_pubkey: "deadbeefdeadbeef".into(),
community_id,
action,
channel_id: None,
metadata: serde_json::json!({"test": true}),
actor_pubkey: Some(vec![0xab; 32]),
object_id: Some(format!("obj_{}", Uuid::new_v4())),
detail: serde_json::json!({"test": true}),
}
}
async fn reset_audit_table(pool: &PgPool) {
sqlx::query("TRUNCATE TABLE audit_log")
.execute(pool)
.await
.unwrap();
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn genesis_entry() {
let _guard = db_lock().lock().await;
async fn community_chain_starts_at_seq_1_with_null_prev() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
svc.ensure_schema().await.unwrap();
reset_audit_table(&pool).await;
let c = make_community(&pool).await;
let entry = svc
.log(sample_new_entry(1, AuditAction::EventCreated))
let e = svc
.log(new_entry(c, AuditAction::EventCreated))
.await
.unwrap();
assert_eq!(entry.prev_hash, GENESIS_HASH);
assert_eq!(entry.seq, 1);
assert_eq!(entry.hash.len(), 64);
assert_eq!(e.seq, 1, "first entry in a community starts at seq 1");
assert!(e.prev_hash.is_none(), "genesis entry has NULL prev_hash");
assert_eq!(e.hash.len(), 32);
assert_eq!(e.community_id, c);
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn chain_integrity() {
let _guard = db_lock().lock().await;
async fn chain_links_within_one_community() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
svc.ensure_schema().await.unwrap();
reset_audit_table(&pool).await;
let c = make_community(&pool).await;
let e1 = svc
.log(sample_new_entry(1, AuditAction::EventCreated))
.log(new_entry(c, AuditAction::EventCreated))
.await
.unwrap();
let e2 = svc
.log(sample_new_entry(1, AuditAction::ChannelCreated))
.log(new_entry(c, AuditAction::ChannelCreated))
.await
.unwrap();
let e3 = svc
.log(sample_new_entry(1, AuditAction::MemberAdded))
.log(new_entry(c, AuditAction::MemberAdded))
.await
.unwrap();
assert_eq!(e1.prev_hash, GENESIS_HASH);
assert_eq!(e2.prev_hash, e1.hash);
assert_eq!(e3.prev_hash, e2.hash);
assert!(svc.verify_chain(e1.seq, e3.seq).await.unwrap());
assert_eq!(e1.seq, 1);
assert_eq!(e2.seq, 2);
assert_eq!(e3.seq, 3);
assert!(e1.prev_hash.is_none());
assert_eq!(e2.prev_hash.as_deref(), Some(e1.hash.as_slice()));
assert_eq!(e3.prev_hash.as_deref(), Some(e2.hash.as_slice()));
assert!(svc
.verify_chain(CommunityId::from_uuid(c), 1, 3)
.await
.unwrap());
}
/// THE isolation property: two communities keep independent chains. Each
/// starts at seq 1; interleaving writes does not link them; verifying one
/// never traverses the other.
#[tokio::test]
#[ignore = "requires Postgres"]
async fn verify_chain_detects_tampering() {
let _guard = db_lock().lock().await;
async fn chains_are_independent_per_community() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
svc.ensure_schema().await.unwrap();
reset_audit_table(&pool).await;
let a = make_community(&pool).await;
let b = make_community(&pool).await;
let e1 = svc
.log(sample_new_entry(1, AuditAction::EventCreated))
// Interleave A and B writes.
let a1 = svc
.log(new_entry(a, AuditAction::EventCreated))
.await
.unwrap();
let b1 = svc
.log(new_entry(b, AuditAction::EventCreated))
.await
.unwrap();
let a2 = svc
.log(new_entry(a, AuditAction::ChannelCreated))
.await
.unwrap();
let b2 = svc
.log(new_entry(b, AuditAction::ChannelCreated))
.await
.unwrap();
// Each community's seq is independent and starts at 1.
assert_eq!((a1.seq, a2.seq), (1, 2));
assert_eq!((b1.seq, b2.seq), (1, 2));
// A's chain links only within A; B's only within B. A2 must NOT chain to
// B1 even though B1 was written between A1 and A2.
assert_eq!(a2.prev_hash.as_deref(), Some(a1.hash.as_slice()));
assert_eq!(b2.prev_hash.as_deref(), Some(b1.hash.as_slice()));
assert_ne!(a2.prev_hash, b1.prev_hash);
// Verifying A's chain traverses only A; same for B.
assert!(svc
.verify_chain(CommunityId::from_uuid(a), 1, 2)
.await
.unwrap());
assert!(svc
.verify_chain(CommunityId::from_uuid(b), 1, 2)
.await
.unwrap());
// get_entries scoped to A returns only A's rows.
let a_rows = svc
.get_entries(CommunityId::from_uuid(a), 1, 100)
.await
.unwrap();
assert!(
a_rows.iter().all(|e| e.community_id == a),
"A read leaked another community"
);
assert_eq!(a_rows.len(), 2);
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn verify_detects_tampering_within_a_community() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
let c = make_community(&pool).await;
svc.log(new_entry(c, AuditAction::EventCreated))
.await
.unwrap();
let e2 = svc
.log(sample_new_entry(1, AuditAction::EventDeleted))
.log(new_entry(c, AuditAction::EventDeleted))
.await
.unwrap();
let e3 = svc
.log(sample_new_entry(1, AuditAction::ChannelDeleted))
svc.log(new_entry(c, AuditAction::ChannelDeleted))
.await
.unwrap();
sqlx::query("UPDATE audit_log SET actor_pubkey = 'tampered' WHERE seq = $1")
// Tamper with e2's stored actor_pubkey.
let tampered: Vec<u8> = vec![0xff; 32];
sqlx::query("UPDATE audit_log SET actor_pubkey = $1 WHERE community_id = $2 AND seq = $3")
.bind(tampered)
.bind(c)
.bind(e2.seq)
.execute(&pool)
.await
.unwrap();
let result = svc.verify_chain(e1.seq, e3.seq).await;
assert!(matches!(result, Err(AuditError::HashMismatch { seq, .. }) if seq == e2.seq));
let r = svc.verify_chain(CommunityId::from_uuid(c), 1, 3).await;
assert!(matches!(r, Err(AuditError::HashMismatch { seq }) if seq == e2.seq));
}
/// A row forged with another community's id cannot pass verification against
/// the chain it was stamped for, because community_id is hashed in. (Models
/// "a row can't be replayed across chains and still verify".)
#[tokio::test]
#[ignore = "requires Postgres"]
async fn auth_events_rejected() {
async fn cross_community_row_does_not_verify() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
let a = make_community(&pool).await;
let b = make_community(&pool).await;
let result = svc
.log(sample_new_entry(KIND_AUTH, AuditAction::AuthSuccess))
.await;
let a1 = svc
.log(new_entry(a, AuditAction::EventCreated))
.await
.unwrap();
assert!(matches!(result, Err(AuditError::AuthEventForbidden)));
// Forge: copy A's seq-1 row's hash into B's chain at seq 1.
sqlx::query(
"INSERT INTO audit_log (community_id, seq, hash, prev_hash, action, actor_pubkey, object_id, detail, created_at)
VALUES ($1, 1, $2, NULL, $3, $4, $5, $6, NOW())",
)
.bind(b)
.bind(&a1.hash) // A's hash, which was computed over community_id = A
.bind(a1.action.as_str())
.bind(a1.actor_pubkey.as_deref())
.bind(a1.object_id.as_deref())
.bind(&a1.detail)
.execute(&pool)
.await
.unwrap();
// Verifying B's chain recomputes the hash with community_id = B, which
// won't match A's stored hash → HashMismatch. The forge is rejected.
let r = svc.verify_chain(CommunityId::from_uuid(b), 1, 1).await;
assert!(matches!(r, Err(AuditError::HashMismatch { seq: 1 })));
}
#[tokio::test]
#[ignore = "requires Postgres"]
async fn verify_empty_range_is_false() {
let _g = db_lock().lock().await;
let Some(pool) = test_pool().await else {
return;
};
let svc = AuditService::new(pool.clone());
let c = make_community(&pool).await;
// No entries for this fresh community.
assert!(!svc
.verify_chain(CommunityId::from_uuid(c), 1, 100)
.await
.unwrap());
}
}