From 6196da22affaf309d9d75b1ad6985b85e85e8724 Mon Sep 17 00:00:00 2001 From: Honey <8e307ae0076a4dab6b94b036ea3edc7e08f823a625269c1e6919e881a048b4d2@buzz.block.builderlab.xyz> Date: Sat, 26 Sep 2026 14:26:02 -0700 Subject: [PATCH 1/6] fix(db): take channel-head writes through the metered writer pool insert_channel_head_checked opened its transaction with pool.begin(), bypassing acquire_writer. Every other event write in this module acquires through it, so this path went uncounted in writer-pool acquire metrics. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: Honey <8e307ae0076a4dab6b94b036ea3edc7e08f823a625269c1e6919e881a048b4d2@buzz.block.builderlab.xyz> --- crates/buzz-db/src/store/event.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index 23082102e9f..7031a60ca4c 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -1732,7 +1732,12 @@ pub async fn insert_channel_head_checked( let received_at = Utc::now(); let incoming_id = event.id.as_bytes(); - let mut tx = pool.begin().await?; + let connection = crate::observability::acquire_writer( + pool, + crate::observability::WriterOperation::EventWrite, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; // Serialize check+insert per (community, kind, channel). let lock_key = event_replacement_lock_key( From 90b4dc145a05d1786e3b7ec7b5748021d94eaa76 Mon Sep 17 00:00:00 2001 From: Honey <8e307ae0076a4dab6b94b036ea3edc7e08f823a625269c1e6919e881a048b4d2@buzz.block.builderlab.xyz> Date: Sat, 26 Sep 2026 14:28:48 -0700 Subject: [PATCH 2/6] feat(relay): implement NIP-AR channel artifacts Adds kind 45010 artifacts: editable records with a stable identity (d), one home channel (h), and full-snapshot revisions chained by prev. - One transaction per write: lock the artifact's head row, require prev to name the current revision, record the revision in a ledger that survives retention, and advance the head. Competing edits cannot both win; resubmitting an accepted event succeeds without reapplying it. - Writes pass the same gates as kind 9 in h; a move also passes them in the source channel. The move stores the destination snapshot and a relay-signed kind 45011 removal marker for the source in the same transaction, so both channels recover by replay. - Redaction is the ordinary kind-9005 removal. The ledger keeps the ID so prev checks still work; redacting the current revision retires the artifact. - HTTP /query and /count accept explicit `artifact: current|history` filters with exact multi-character tag matching. Generic REQ and search return only current revisions. Kind 5 against artifacts is rejected. NIP-11 advertises the limits. - NIP-AR updated to match: kind-9 permission checks, 9005 redaction, retirement, and no creator-only delete rule. Co-authored-by: Smartie Co-authored-by: Fizz <400e8babadcee6a7f420103f10a2849d84c4a9c71d5bd04f3948c814216648a3@buzz.block.builderlab.xyz> Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: Honey <8e307ae0076a4dab6b94b036ea3edc7e08f823a625269c1e6919e881a048b4d2@buzz.block.builderlab.xyz> --- Cargo.lock | 1 + crates/buzz-core/src/artifact.rs | 376 ++++++++++++++++++ crates/buzz-core/src/kind.rs | 9 +- crates/buzz-core/src/lib.rs | 3 + crates/buzz-db/src/lib.rs | 10 +- crates/buzz-db/src/runtime/migration.rs | 13 +- .../tests/thread_window_postgres_tests.rs | 2 +- crates/buzz-db/src/store/artifact.rs | 194 +++++++++ .../src/store/artifact_postgres_tests.rs | 344 ++++++++++++++++ crates/buzz-db/src/store/artifact_query.rs | 72 ++++ crates/buzz-db/src/store/deletion.rs | 4 + crates/buzz-db/src/store/event.rs | 19 + crates/buzz-db/src/store/mod.rs | 8 + crates/buzz-relay/src/api/artifact.rs | 77 ++++ .../src/api/artifact_postgres_tests.rs | 289 ++++++++++++++ crates/buzz-relay/src/api/bridge.rs | 13 + crates/buzz-relay/src/api/mod.rs | 2 + crates/buzz-relay/src/handlers/artifact.rs | 148 +++++++ crates/buzz-relay/src/handlers/ingest.rs | 75 +++- crates/buzz-relay/src/handlers/mod.rs | 3 + .../buzz-relay/src/handlers/side_effects.rs | 10 + crates/buzz-relay/src/nip11.rs | 18 +- crates/buzz-relay/src/protocol.rs | 45 +++ crates/buzz-search/Cargo.toml | 1 + crates/buzz-search/src/query.rs | 3 + .../tests/postgres_fts_integration.rs | 163 ++++---- docs/nips/NIP-AR.md | 15 +- migrations/0052_channel_artifacts.sql | 41 ++ schema/schema.sql | 42 ++ scripts/reconcile-schema-after-pgschema.sql | 8 + 30 files changed, 1905 insertions(+), 103 deletions(-) create mode 100644 crates/buzz-core/src/artifact.rs create mode 100644 crates/buzz-db/src/store/artifact.rs create mode 100644 crates/buzz-db/src/store/artifact_postgres_tests.rs create mode 100644 crates/buzz-db/src/store/artifact_query.rs create mode 100644 crates/buzz-relay/src/api/artifact.rs create mode 100644 crates/buzz-relay/src/api/artifact_postgres_tests.rs create mode 100644 crates/buzz-relay/src/handlers/artifact.rs create mode 100644 migrations/0052_channel_artifacts.sql diff --git a/Cargo.lock b/Cargo.lock index 95e80f1c129..e57532c25cc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1525,6 +1525,7 @@ version = "0.1.0" dependencies = [ "buzz-core", "buzz-datastore-tracing", + "buzz-db", "metrics", "sqlx", "thiserror 2.0.18", diff --git a/crates/buzz-core/src/artifact.rs b/crates/buzz-core/src/artifact.rs new file mode 100644 index 00000000000..b4b8a27638a --- /dev/null +++ b/crates/buzz-core/src/artifact.rs @@ -0,0 +1,376 @@ +//! NIP-AR envelope validation; client-defined payloads and annotations are opaque. +use nostr::Event; +use uuid::Uuid; +/// Maximum artifact tags, including the envelope. +pub const MAX_TAGS: usize = 256; +/// Maximum UTF-8 bytes in a tag name. +pub const MAX_TAG_NAME_BYTES: usize = 128; +/// Maximum UTF-8 bytes in each tag value. +pub const MAX_TAG_VALUE_BYTES: usize = 4096; +/// Maximum UTF-8 bytes across all tags. +pub const MAX_TAG_BYTES: usize = 65536; +/// Maximum predicates per artifact query. +pub const MAX_PREDICATES: usize = 32; +/// Maximum total requested predicate values. +pub const MAX_QUERY_VALUES: usize = 256; +/// Maximum artifact page size. +pub const MAX_PAGE_SIZE: usize = 1000; +/// Lifecycle operation named by a revision's `op` tag. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ArtifactOp { + /// First revision of a new identity. + Create, + /// Content change within the same home. + Update, + /// Content snapshot published into a new home. + Move, + /// Soft delete; only `Restore` may follow. + Delete, + /// Complete snapshot that revives a deleted artifact. + Restore, +} +impl ArtifactOp { + fn parse(value: &str) -> Option { + Some(match value { + "create" => Self::Create, + "update" => Self::Update, + "move" => Self::Move, + "delete" => Self::Delete, + "restore" => Self::Restore, + _ => return None, + }) + } +} +/// The authoritative envelope, independent of any client's content schema. +#[derive(Debug)] +pub struct ArtifactEnvelope { + /// Community-local stable identity. + pub id: Uuid, + /// Home channel. + pub home: Uuid, + /// Immutable type name. + pub artifact_type: String, + /// Lifecycle operation. + pub op: ArtifactOp, + /// Expected previous revision, absent on create. + pub prev: Option>, + /// Optional conversation anchor. + pub root: Option>, +} +/// Parse a canonical, non-nil UUID. +pub fn canonical_uuid(value: &str) -> Result { + let id = Uuid::parse_str(value).map_err(|_| "invalid UUID")?; + if id.is_nil() || id.to_string() != value { + return Err("UUID must be canonical lowercase and non-nil"); + } + Ok(id) +} +/// Parse a canonical Nostr event identifier. +pub fn event_id(value: &str) -> Result, &'static str> { + if value.len() != 64 + || !value + .bytes() + .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b)) + { + return Err("event ID must be 64 lowercase hex characters"); + } + hex::decode(value).map_err(|_| "invalid event ID") +} +/// Validate the complete envelope, leaving auth-tag verification to the relay. +pub fn validate(event: &Event) -> Result { + if event.kind.as_u16() != 45010 { + return Err("not an artifact revision"); + } + if event.tags.len() > MAX_TAGS { + return Err("too many tags"); + } + let names = ["ar", "d", "h", "type", "title", "op", "root", "prev"]; + let mut fields = std::collections::HashMap::new(); + let mut bytes = 0; + for tag in event.tags.iter() { + let parts = tag.as_slice(); + let Some(name) = parts.first() else { + return Err("empty tag"); + }; + if name.len() > MAX_TAG_NAME_BYTES { + return Err("tag name too long"); + } + for value in parts { + bytes += value.len(); + if value.len() > MAX_TAG_VALUE_BYTES { + return Err("tag value too long"); + } + } + if names.contains(&name.as_str()) + && (parts.len() != 2 || fields.insert(name.as_str(), parts[1].as_str()).is_some()) + { + return Err("envelope tags must occur once with exactly two elements"); + } + } + if bytes > MAX_TAG_BYTES { + return Err("total tag bytes exceeded"); + } + let field = |name| { + fields + .get(name) + .copied() + .ok_or("missing required envelope tag") + }; + if field("ar")? != "1" { + return Err("unsupported artifact envelope version"); + } + let id = canonical_uuid(field("d")?)?; + let home = canonical_uuid(field("h")?)?; + let artifact_type = field("type")?; + if artifact_type.len() > 128 + || !artifact_type.contains('.') + || !artifact_type.split('.').all(|c| { + !c.is_empty() + && c.as_bytes()[0].is_ascii_lowercase() + && c.bytes() + .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'_' || b == b'-') + }) + { + return Err("invalid namespaced artifact type"); + } + let op = ArtifactOp::parse(field("op")?).ok_or("invalid artifact operation")?; + let prev = fields.get("prev").map(|s| event_id(s)).transpose()?; + if (op == ArtifactOp::Create) != prev.is_none() { + return Err("prev required exactly on non-create revisions"); + } + let root = fields.get("root").map(|s| event_id(s)).transpose()?; + if op == ArtifactOp::Delete { + if fields.contains_key("title") || !event.content.is_empty() { + return Err("delete must omit title and have empty content"); + } + if event + .tags + .iter() + .any(|t| !names.contains(&t.as_slice()[0].as_str()) && t.as_slice()[0] != "auth") + { + return Err("delete allows only envelope and verified auth tags"); + } + } else { + let title = field("title")?; + if title.trim().is_empty() || title.len() > 512 { + return Err("title must be nonblank and at most 512 UTF-8 bytes"); + } + } + Ok(ArtifactEnvelope { + id, + home, + artifact_type: artifact_type.into(), + op, + prev, + root, + }) +} +#[cfg(test)] +mod tests { + use super::*; + use nostr::{EventBuilder, Keys, Kind, Tag}; + fn event(extra: Vec>, op: &str, content: &str) -> Event { + let mut tags = vec![ + vec!["ar".into(), "1".into()], + vec!["d".into(), Uuid::new_v4().to_string()], + vec!["h".into(), Uuid::new_v4().to_string()], + vec!["type".into(), "buzz.task".into()], + vec!["op".into(), op.into()], + ]; + if op != "delete" { + tags.push(vec!["title".into(), "Title".into()]); + } + if op != "create" { + tags.push(vec!["prev".into(), "a".repeat(64)]); + } + tags.extend(extra); + EventBuilder::new(Kind::Custom(45010), content) + .tags(tags.into_iter().map(|t| Tag::parse(t).unwrap())) + .sign_with_keys(&Keys::generate()) + .unwrap() + } + #[test] + fn lifecycle_envelopes() { + for op in ["create", "update", "move", "delete", "restore"] { + assert!(validate(&event(vec![], op, "")).is_ok(), "{op}"); + } + assert!(validate(&event( + vec![vec!["project".into(), "opaque".into(), "anything".into()]], + "create", + "not json" + )) + .is_ok()); + } + #[test] + fn malformed_and_delete_payloads() { + for (tags, op, content) in [ + (vec![vec!["title".into(), "duplicate".into()]], "create", ""), + (vec![vec!["ar".into(), "2".into()]], "create", ""), + (vec![vec!["prev".into(), "a".repeat(64)]], "create", ""), + (vec![vec!["root".into(), "A".repeat(64)]], "create", ""), + (vec![vec!["project".into(), "hidden".into()]], "delete", ""), + (vec![vec!["title".into(), "hidden".into()]], "delete", ""), + (vec![], "delete", "hidden"), + (vec![vec!["x".into(), "x".repeat(4097)]], "create", ""), + (vec![vec!["x".repeat(129), "x".into()]], "create", ""), + ] { + assert!(validate(&event(tags, op, content)).is_err()); + } + } + #[test] + fn filter_routes_and_views() { + use serde_json::json; + assert_eq!( + route_filter(&json!({"kinds":[30621],"#buzz-channel":["c"]})), + FilterRoute::Generic + ); + assert_eq!( + route_filter(&json!({"kinds":[45010,45011],"#h":["c"]})), + FilterRoute::Generic + ); + assert!(matches!( + route_filter(&json!({"kinds":[45010],"#project":["p"]})), + FilterRoute::Rejected(_) + )); + assert_eq!( + route_filter(&json!({"artifact":"current"})), + FilterRoute::Artifact + ); + let query = parse_query(&json!({"artifact":"current","#project":["p"]})).unwrap(); + assert_eq!(query.view, ArtifactView::Current); + for invalid in [ + json!({"artifact":"lookup","#d":["x"]}), + json!({"artifact":"history"}), + json!({"artifact":"current","kinds":[45011]}), + json!({"artifact":"current","search":"x"}), + json!({"artifact":"current","limit":1001}), + ] { + assert!(parse_query(&invalid).is_err(), "{invalid}"); + } + } + #[test] + fn canonical_identifiers() { + assert!(canonical_uuid("00000000-0000-0000-0000-000000000000").is_err()); + assert!(canonical_uuid(&Uuid::new_v4().simple().to_string()).is_err()); + assert!(event_id(&"f".repeat(64)).is_ok()); + assert!(event_id(&"F".repeat(64)).is_err()); + } +} + +/// Which artifact revisions an explicit query reads. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ArtifactView { + /// Each live artifact's current revision; deleted artifacts are omitted. + Current, + /// Every stored, unredacted revision of the artifacts named by `#d`. + History, +} +/// Explicit HTTP artifact query. All `#name` predicates compare the first +/// value of the same tag. +#[derive(Debug)] +pub struct ArtifactQuery { + /// Query view. + pub view: ArtifactView, + /// Exact tag-name/value predicates. + pub tags: Vec<(String, Vec)>, + /// Maximum result count. + pub limit: i64, + /// Bounded page offset. + pub offset: i64, +} +/// Maximum artifact page offset. +pub const MAX_OFFSET: u64 = 10000; +/// How a raw REQ/COUNT filter must be served. +#[derive(Debug, PartialEq, Eq)] +pub enum FilterRoute { + /// Standard Nostr filter handling. + Generic, + /// Explicit artifact query (`artifact` key present). + Artifact, + /// Would silently drop predicates on the generic path. + Rejected(&'static str), +} +/// Route a raw filter. Generic filters drop multi-character tag predicates, +/// so artifact-kind filters carrying them must use an explicit artifact query. +pub fn route_filter(value: &serde_json::Value) -> FilterRoute { + if value.get("artifact").is_some() { + return FilterRoute::Artifact; + } + let artifact_kind = value + .get("kinds") + .and_then(|k| k.as_array()) + .is_some_and(|ks| ks.iter().any(|k| matches!(k.as_u64(), Some(45010 | 45011)))); + let multi_character = value + .as_object() + .is_some_and(|o| o.keys().any(|k| k.starts_with('#') && k.len() > 2)); + if artifact_kind && multi_character { + return FilterRoute::Rejected( + "multi-character tag predicates on artifacts require an artifact query", + ); + } + FilterRoute::Generic +} +/// Parse an artifact filter, explicitly rejecting unsupported predicates. +pub fn parse_query(value: &serde_json::Value) -> Result { + let object = value + .as_object() + .ok_or("artifact filter must be an object")?; + let view = match object.get("artifact").and_then(|v| v.as_str()) { + Some("current") => ArtifactView::Current, + Some("history") => ArtifactView::History, + _ => return Err("artifact must be \"current\" or \"history\""), + }; + if let Some(kinds) = object.get("kinds") { + if kinds.as_array().map(Vec::as_slice) != Some(&[serde_json::json!(45010)]) { + return Err("artifact queries accept only kinds [45010]"); + } + } + let mut tags = Vec::new(); + let mut total = 0; + let mut bytes = 0; + for (name, value) in object { + if let Some(name) = name.strip_prefix('#') { + if name.is_empty() || name.len() > MAX_TAG_NAME_BYTES { + return Err("invalid predicate name size"); + } + let values = value.as_array().ok_or("predicate values must be arrays")?; + let mut parsed = Vec::new(); + for value in values { + let value = value.as_str().ok_or("predicate values must be strings")?; + if value.len() > MAX_TAG_VALUE_BYTES { + return Err("predicate value too large"); + } + bytes += name.len() + value.len(); + total += 1; + parsed.push(value.into()); + } + tags.push((name.into(), parsed)); + } else if !["artifact", "kinds", "limit", "offset"].contains(&name.as_str()) { + return Err("unsupported artifact predicate"); + } + } + if tags.len() > MAX_PREDICATES || total > MAX_QUERY_VALUES || bytes > MAX_TAG_BYTES { + return Err("artifact predicate limits exceeded"); + } + let number = |name: &str, default: u64| { + object + .get(name) + .map(|v| v.as_u64().ok_or("invalid limit or offset")) + .transpose() + .map(|v| v.unwrap_or(default)) + }; + let limit = number("limit", 100)?; + let offset = number("offset", 0)?; + if limit == 0 || limit > MAX_PAGE_SIZE as u64 || offset > MAX_OFFSET { + return Err("artifact page limit exceeded"); + } + if view == ArtifactView::History && !tags.iter().any(|(n, v)| n == "d" && !v.is_empty()) { + return Err("history requires #d"); + } + Ok(ArtifactQuery { + view, + tags, + limit: limit as i64, + offset: offset as i64, + }) +} diff --git a/crates/buzz-core/src/kind.rs b/crates/buzz-core/src/kind.rs index 45a7bf2e395..32d8959fb69 100644 --- a/crates/buzz-core/src/kind.rs +++ b/crates/buzz-core/src/kind.rs @@ -550,6 +550,10 @@ pub const KIND_AGENT_TURN_METRIC: u32 = 44200; // V1 used addressable range (30001–30003) — wrong. /// A forum post (thread root). pub const KIND_FORUM_POST: u32 = 45001; +/// NIP-AR complete artifact revision. +pub const KIND_ARTIFACT: u32 = 45010; +/// Relay-authenticated artifact removal. +pub const KIND_ARTIFACT_REMOVAL: u32 = 45011; /// A vote on a forum post. pub const KIND_FORUM_VOTE: u32 = 45002; /// A comment reply on a forum post. @@ -735,6 +739,8 @@ pub const ALL_KINDS: &[u32] = &[ KIND_USER_STATUS, KIND_READ_STATE, KIND_FORUM_POST, + KIND_ARTIFACT, + KIND_ARTIFACT_REMOVAL, KIND_FORUM_VOTE, KIND_FORUM_COMMENT, KIND_WORKFLOW_TRIGGER, @@ -836,7 +842,8 @@ pub const fn is_command_kind(kind: u32) -> bool { pub const fn is_relay_only_kind(kind: u32) -> bool { matches!( kind, - KIND_NIP43_MEMBERSHIP_LIST + KIND_ARTIFACT_REMOVAL + | KIND_NIP43_MEMBERSHIP_LIST | KIND_CHANNEL_SUMMARY | KIND_PRESENCE_SNAPSHOT | KIND_DM_VISIBILITY diff --git a/crates/buzz-core/src/lib.rs b/crates/buzz-core/src/lib.rs index ec11dc0bf02..e7669c874b0 100644 --- a/crates/buzz-core/src/lib.rs +++ b/crates/buzz-core/src/lib.rs @@ -81,3 +81,6 @@ pub mod test_helpers { StoredEvent::with_received_at(make_event(kind), Utc::now(), channel_id, true) } } + +/// NIP-AR channel artifact envelope and limits. +pub mod artifact; diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 3dac0b1ba2f..365d88e1a4e 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -62,11 +62,11 @@ pub(crate) use runtime::{ RoutePredicate, }; pub use store::{ - admin_moderation, allowlist, api_token, archived_identities, channel, channel_members, - community, deletion, dm, event, feed, git_repo, moderation, operator_listener, partition, - product_feedback, push, reaction, read_state, relay_admin_actions, relay_invite, relay_members, - relay_operators, reminder, replaceable, storage_accounting, thread, thread_window, usage, user, - workflow, + admin_moderation, allowlist, api_token, archived_identities, artifact, channel, + channel_members, community, deletion, dm, event, feed, git_repo, moderation, operator_listener, + partition, product_feedback, push, reaction, read_state, relay_admin_actions, relay_invite, + relay_members, relay_operators, reminder, replaceable, storage_accounting, thread, + thread_window, usage, user, workflow, }; pub use allowlist::AllowlistEntry; diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 9eae1dc72b7..9a003339da0 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -705,7 +705,7 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 51); + assert_eq!(migrations.len(), 52); assert_eq!(migrations[48].version, 49); assert_eq!(migrations[50].version, 51); assert!(migrations[48] @@ -1363,6 +1363,11 @@ mod postgres_tests { .contains("'rate_limit_violations', 'operator_listener_outbox'\n ]::TEXT[])"), "schema.sql must exclude the deployment-global listener outbox from tenant fencing" ); + assert_eq!(migrations[51].version, 52); + assert!(migrations[51] + .sql + .as_str() + .contains("CREATE TABLE artifact_heads")); } #[test] @@ -1934,6 +1939,7 @@ mod postgres_tests { let mut expected_fences = migration.fence_attachments.clone(); expected_fences.remove("product_feedback"); expected_fences.remove("rate_limit_violations"); + expected_fences.extend(["artifact_heads", "artifact_revisions"].map(str::to_owned)); assert_eq!( expected_fences, schema.fence_attachments, "write-fence attachment targets differ after recovery policy" @@ -2880,6 +2886,11 @@ mod postgres_tests { "all NIP-FI tables must be absent after migration 0044: {present:?}" ); + // Complete later additive migrations before comparing to the current + // binary's complete tenant-table inventory. + run_migrations(&pool) + .await + .expect("complete current migrations"); // The deletion catalog must validate with ledger relations gone. crate::deletion::DeletionStore::new(pool.clone()) .validate_catalog() diff --git a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs index 5d51f4348d5..82b7d2fb0a6 100644 --- a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs +++ b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs @@ -311,7 +311,7 @@ async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer() production_result.is_ok(), "production migrator must preserve ingestion progress: {production_result:?}" ); - assert_eq!(version, 51); + assert_eq!(version, 52); assert_eq!(final_oid, oid, "prebuild must not be replaced"); assert_eq!(count, 4, "all writer witnesses must persist"); } diff --git a/crates/buzz-db/src/store/artifact.rs b/crates/buzz-db/src/store/artifact.rs new file mode 100644 index 00000000000..c3e75847610 --- /dev/null +++ b/crates/buzz-db/src/store/artifact.rs @@ -0,0 +1,194 @@ +//! NIP-AR atomic artifact acceptance. Payloads remain opaque signed events. +//! +//! Channel write permission is checked by the relay's ordinary ingest gates +//! before this transaction; here the relay only serializes each identity and +//! compares the expected head. +use crate::{Db, DbError, Result}; +use buzz_core::artifact::{ArtifactEnvelope, ArtifactOp}; +use buzz_core::{CommunityId, StoredEvent}; +use nostr::{Event, EventBuilder, Keys, Kind, Tag}; +use sqlx::Row; +use uuid::Uuid; + +/// Acceptance outcome. Only a stale-`prev` conflict carries the current head, +/// which the caller discloses after authorizing its channel. +#[derive(Debug)] +pub enum ArtifactOutcome { + /// Accepted once; stored events to publish (the revision, plus the source + /// removal on a move). + Accepted(Vec), + /// Identical previously accepted event; no repeated side effects. + Duplicate, + /// `prev` or identity no longer matches; the client should reconcile. + Conflict(&'static str), + /// Protocol violation with a public-safe reason. + Rejected(&'static str), +} + +fn invalid(message: &str) -> DbError { + DbError::InvalidData(message.into()) +} + +fn removal_marker(keys: &Keys, artifact: Uuid, source: Uuid, position: i64) -> Result { + let tags = [ + ["ar", "1"].map(str::to_owned), + ["d".into(), artifact.to_string()], + ["h".into(), source.to_string()], + ["reason".into(), "moved".into()], + ["position".into(), position.to_string()], + ] + .into_iter() + .map(Tag::parse) + .collect::, _>>() + .map_err(|e| invalid(&e.to_string()))?; + EventBuilder::new(Kind::Custom(45011), "") + .tags(tags) + .sign_with_keys(keys) + .map_err(|e| invalid(&e.to_string())) +} + +impl Db { + /// Check the durable acceptance ledger, which survives redaction and retention. + pub async fn artifact_accepted(&self, community: CommunityId, id: &[u8]) -> Result { + let mut conn = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::EventWrite, + ) + .await?; + Ok(sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM artifact_revisions WHERE community_id=$1 AND event_id=$2)", + ) + .bind(community.as_uuid()) + .bind(id) + .fetch_one(&mut *conn) + .await?) + } + + /// Current home channel, read before authorizing a move's source. + pub async fn artifact_home( + &self, + community: CommunityId, + artifact: Uuid, + ) -> Result> { + let mut conn = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::EventWrite, + ) + .await?; + Ok(sqlx::query_scalar( + "SELECT channel_id FROM artifact_heads WHERE community_id=$1 AND artifact_id=$2", + ) + .bind(community.as_uuid()) + .bind(artifact) + .fetch_optional(&mut *conn) + .await?) + } + + /// Atomically compare the expected head, store the full revision, advance + /// the head, and on a move store the source removal. `authorized_source` is + /// the move source the caller authorized; a head that has since moved elsewhere + /// conflicts. Relay keys only sign removals. + pub async fn accept_artifact( + &self, + community: CommunityId, + event: &Event, + env: &ArtifactEnvelope, + authorized_source: Option, + relay_keys: &Keys, + ) -> Result { + let mut tx = self.begin_event_write_transaction().await?; + sqlx::query("SET LOCAL statement_timeout='5s'") + .execute(&mut *tx) + .await?; + // A coordinate lock handles the missing-row create race as well as edits. + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1,0))") + .bind(format!("artifact:{}:{}", community.as_uuid(), env.id)) + .execute(&mut *tx) + .await?; + let duplicate: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM artifact_revisions WHERE community_id=$1 AND event_id=$2)", + ) + .bind(community.as_uuid()) + .bind(event.id.as_bytes().as_slice()) + .fetch_one(&mut *tx) + .await?; + if duplicate { + return Ok(ArtifactOutcome::Duplicate); + } + let head = sqlx::query( + "SELECT event_id,channel_id,artifact_type,root,deleted FROM artifact_heads WHERE community_id=$1 AND artifact_id=$2 FOR UPDATE", + ) + .bind(community.as_uuid()) + .bind(env.id) + .fetch_optional(&mut *tx) + .await?; + let source = head.as_ref().map(|h| h.get::("channel_id")); + let old_root = head + .as_ref() + .and_then(|h| h.get::>, _>("root")); + match (&head, env.op) { + (Some(_), ArtifactOp::Create) => { + return Ok(ArtifactOutcome::Conflict("artifact identity is taken")) + } + (None, ArtifactOp::Create) => {} + (None, _) => return Ok(ArtifactOutcome::Conflict("artifact head unavailable")), + (Some(head), op) => { + let current: Vec = head.get("event_id"); + if env.prev.as_deref() != Some(current.as_slice()) { + return Ok(ArtifactOutcome::Conflict("artifact head changed")); + } + if env.artifact_type != head.get::("artifact_type") { + return Ok(ArtifactOutcome::Rejected("artifact type is immutable")); + } + if head.get::("deleted") != (op == ArtifactOp::Restore) { + return Ok(ArtifactOutcome::Rejected( + "deleted artifacts require restore; live artifacts cannot restore", + )); + } + if (op == ArtifactOp::Move) != (source != Some(env.home)) { + return Ok(ArtifactOutcome::Rejected( + "only move changes home and move must change home", + )); + } + if op == ArtifactOp::Move && authorized_source != source { + return Ok(ArtifactOutcome::Conflict("artifact home changed")); + } + if op == ArtifactOp::Delete && env.root != old_root { + return Ok(ArtifactOutcome::Rejected("delete preserves root")); + } + } + } + if head.is_none() || env.root != old_root || source != Some(env.home) { + if let Some(root) = &env.root { + let exists: bool = sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM events WHERE community_id=$1 AND id=$2 AND channel_id=$3 AND deleted_at IS NULL AND kind IN (9,40002,45001,45003))") + .bind(community.as_uuid()).bind(root).bind(env.home).fetch_one(&mut *tx).await?; + if !exists { + return Ok(ArtifactOutcome::Rejected( + "root must be an existing conversation anchor in home", + )); + } + } + } + let position: i64 = sqlx::query_scalar("INSERT INTO artifact_revisions (community_id,event_id,artifact_id) VALUES ($1,$2,$3) RETURNING position") + .bind(community.as_uuid()).bind(event.id.as_bytes().as_slice()).bind(env.id).fetch_one(&mut *tx).await?; + sqlx::query("INSERT INTO artifact_heads (community_id,artifact_id,event_id,channel_id,artifact_type,root,deleted) VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (community_id,artifact_id) DO UPDATE SET event_id=EXCLUDED.event_id,channel_id=EXCLUDED.channel_id,root=EXCLUDED.root,deleted=EXCLUDED.deleted") + .bind(community.as_uuid()).bind(env.id).bind(event.id.as_bytes().as_slice()).bind(env.home).bind(&env.artifact_type).bind(&env.root).bind(env.op == ArtifactOp::Delete).execute(&mut *tx).await?; + let (stored, _) = + crate::event::insert_event_in_transaction(&mut tx, community, event, Some(env.home)) + .await?; + let mut accepted = vec![stored]; + if let Some(source) = source.filter(|_| env.op == ArtifactOp::Move) { + let removal = removal_marker(relay_keys, env.id, source, position)?; + let (stored, _) = crate::event::insert_event_in_transaction( + &mut tx, + community, + &removal, + Some(source), + ) + .await?; + accepted.push(stored); + } + tx.commit().await?; + Ok(ArtifactOutcome::Accepted(accepted)) + } +} diff --git a/crates/buzz-db/src/store/artifact_postgres_tests.rs b/crates/buzz-db/src/store/artifact_postgres_tests.rs new file mode 100644 index 00000000000..97655ac1e48 --- /dev/null +++ b/crates/buzz-db/src/store/artifact_postgres_tests.rs @@ -0,0 +1,344 @@ +//! PostgreSQL regression tests for the production artifact transaction/query seams. +use super::artifact::ArtifactOutcome; +use crate::Db; +use buzz_core::{artifact, CommunityId}; +use nostr::{Event, EventBuilder, Keys, Kind, Tag}; +use sqlx::PgPool; +use uuid::Uuid; + +struct Fixture { + db: Db, + community: CommunityId, + a: Uuid, + b: Uuid, + owner: Keys, + peer: Keys, + relay: Keys, +} +impl Fixture { + async fn new() -> Self { + let pool = PgPool::connect(&crate::test_support::database_url()) + .await + .unwrap(); + if std::env::var("BUZZ_TEST_SCHEMA_MODE").as_deref() != Ok("desired") { + crate::migration::run_migrations(&pool).await.unwrap(); + } + let db = Db::from_pool(pool); + db.ensure_future_partitions(1).await.unwrap(); + let community = CommunityId::from_uuid(Uuid::new_v4()); + sqlx::query("INSERT INTO communities(id,host) VALUES($1,$2)") + .bind(community.as_uuid()) + .bind(format!("artifact-{}.test", community.as_uuid())) + .execute(&db.pool) + .await + .unwrap(); + let owner = Keys::generate(); + let peer = Keys::generate(); + let a = Uuid::new_v4(); + let b = Uuid::new_v4(); + for (id, visibility) in [(a, "open"), (b, "private")] { + sqlx::query("INSERT INTO channels(community_id,id,name,visibility,channel_type,created_by) VALUES($1,$2,$3,$4::channel_visibility,'stream',$5)") + .bind(community.as_uuid()).bind(id).bind(id.to_string()).bind(visibility).bind(owner.public_key().to_bytes().as_slice()).execute(&db.pool).await.unwrap(); + sqlx::query("INSERT INTO channel_members(community_id,channel_id,pubkey,role) VALUES($1,$2,$3,'owner')") + .bind(community.as_uuid()).bind(id).bind(owner.public_key().to_bytes().as_slice()).execute(&db.pool).await.unwrap(); + } + Self { + db, + community, + a, + b, + owner, + peer, + relay: Keys::generate(), + } + } + fn revision( + &self, + d: Uuid, + op: &str, + home: Uuid, + prev: Option<&Event>, + key: &Keys, + extra: Vec>, + ) -> Event { + let mut tags = vec![ + vec!["ar".into(), "1".into()], + vec!["d".into(), d.to_string()], + vec!["h".into(), home.to_string()], + vec!["type".into(), "buzz.task".into()], + vec!["op".into(), op.into()], + ]; + if op != "delete" { + tags.push(vec!["title".into(), "Test".into()]); + } + if let Some(e) = prev { + tags.push(vec!["prev".into(), e.id.to_hex()]); + } + tags.extend(extra); + EventBuilder::new( + Kind::Custom(45010), + if op == "delete" { "" } else { "secret" }, + ) + .tags(tags.into_iter().map(|t| Tag::parse(t).unwrap())) + .sign_with_keys(key) + .unwrap() + } + async fn accept(&self, e: &Event, source: Option) -> ArtifactOutcome { + let env = artifact::validate(e).unwrap(); + self.db + .accept_artifact(self.community, e, &env, source, &self.relay) + .await + .unwrap() + } + async fn count(&self, key: &Keys, value: serde_json::Value) -> i64 { + self.query(key, value, true).await.1 + } + async fn query( + &self, + key: &Keys, + value: serde_json::Value, + count: bool, + ) -> (Vec, i64) { + self.db + .query_artifacts( + self.community, + key.public_key().as_bytes(), + &artifact::parse_query(&value).unwrap(), + count, + ) + .await + .unwrap() + } +} +#[tokio::test] +#[ignore = "requires Postgres"] +async fn lifecycle_cas_queries_move_redaction_and_retention() { + let f = Fixture::new().await; + let d = Uuid::new_v4(); + let project = |v: &str| vec![vec!["project".to_string(), v.to_string()]]; + let create = f.revision(d, "create", f.a, None, &f.owner, project("A")); + assert!(matches!( + f.accept(&create, None).await, + ArtifactOutcome::Accepted(_) + )); + let x = f.revision(d, "update", f.a, Some(&create), &f.owner, project("B")); + let y = f.revision(d, "update", f.a, Some(&create), &f.peer, project("B")); + let (rx, ry) = tokio::join!(f.accept(&x, None), f.accept(&y, None)); + let head = match (rx, ry) { + (ArtifactOutcome::Accepted(_), ArtifactOutcome::Conflict(..)) => &x, + (ArtifactOutcome::Conflict(..), ArtifactOutcome::Accepted(_)) => &y, + other => panic!("exactly one CAS winner: {other:?}"), + }; + assert!(matches!( + f.accept(&create, None).await, + ArtifactOutcome::Duplicate + )); + let current = |p: &str| serde_json::json!({"artifact":"current","#d":[d],"#project":[p]}); + assert_eq!(f.count(&f.owner, current("A")).await, 0); + assert_eq!(f.count(&f.owner, current("B")).await, 1); + // First value in the same tag, never an annotation or reverse-name match. + let wrong = f.revision( + d, + "update", + f.a, + Some(head), + &f.owner, + vec![ + vec!["project".into(), "C".into(), "B".into()], + vec!["B".into(), "project".into()], + ], + ); + assert!(matches!( + f.accept(&wrong, None).await, + ArtifactOutcome::Accepted(_) + )); + assert_eq!(f.count(&f.owner, current("B")).await, 0); + + // A move commits only against the source the relay authorized. + let moved = f.revision(d, "move", f.b, Some(&wrong), &f.owner, vec![]); + assert!(matches!( + f.accept(&moved, Some(f.b)).await, + ArtifactOutcome::Conflict(..) + )); + let ArtifactOutcome::Accepted(stored) = f.accept(&moved, Some(f.a)).await else { + panic!("move accepted"); + }; + let removal = &stored[1]; + assert_eq!(removal.event.kind.as_u16(), 45011); + assert_eq!(removal.channel_id, Some(f.a)); + assert!(!serde_json::to_string(&removal.event) + .unwrap() + .contains(&f.b.to_string())); + let all = serde_json::json!({"artifact":"current","#d":[d]}); + let history = serde_json::json!({"artifact":"history","#d":[d]}); + assert_eq!(f.count(&f.peer, all.clone()).await, 0); + assert_eq!(f.count(&f.peer, history.clone()).await, 3); + assert_eq!(f.count(&f.owner, history.clone()).await, 4); + + // Retention removes old payloads but reserves identity and the current payload. + let retention = || { + sqlx::query("DELETE FROM events WHERE community_id=$1 AND id IN ($2,$3)") + .bind(f.community.as_uuid()) + .bind(create.id.as_bytes().as_slice()) + .bind(moved.id.as_bytes().as_slice()) + }; + assert_eq!( + retention() + .execute(&f.db.pool) + .await + .unwrap() + .rows_affected(), + 1 + ); + assert!(matches!( + f.accept(&create, None).await, + ArtifactOutcome::Duplicate + )); + + // Redaction is the generic soft delete: the head stays, its payload is hidden + // everywhere and may then be purged. + assert!(f + .db + .soft_delete_event_and_update_thread(f.community, moved.id.as_bytes(), None, None) + .await + .unwrap()); + assert_eq!(f.count(&f.owner, all.clone()).await, 0); + assert_eq!(f.count(&f.owner, history.clone()).await, 2); + assert_eq!( + retention() + .execute(&f.db.pool) + .await + .unwrap() + .rows_affected(), + 1 + ); + assert!(f + .db + .soft_delete_event_and_update_thread(f.community, removal.event.id.as_bytes(), None, None) + .await + .is_err()); + + let deleted = f.revision(d, "delete", f.b, Some(&moved), &f.owner, vec![]); + assert!(matches!( + f.accept(&deleted, None).await, + ArtifactOutcome::Accepted(_) + )); + assert_eq!(f.count(&f.owner, all.clone()).await, 0); + let update = f.revision(d, "update", f.b, Some(&deleted), &f.owner, vec![]); + assert!(matches!( + f.accept(&update, None).await, + ArtifactOutcome::Rejected(_) + )); + let restore = f.revision(d, "restore", f.b, Some(&deleted), &f.peer, vec![]); + assert!(matches!( + f.accept(&restore, None).await, + ArtifactOutcome::Accepted(_) + )); + assert_eq!(f.count(&f.owner, all.clone()).await, 1); + // Reads observe removed private membership. + sqlx::query( + "UPDATE channel_members SET removed_at=now() WHERE community_id=$1 AND channel_id=$2", + ) + .bind(f.community.as_uuid()) + .bind(f.b) + .execute(&f.db.pool) + .await + .unwrap(); + assert_eq!(f.count(&f.owner, all).await, 0); + f.db.validate_deletion_serving_catalog().await.unwrap(); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn root_anchor_and_tenant_boundaries() { + let f = Fixture::new().await; + let d = Uuid::new_v4(); + let invalid = f.revision( + d, + "create", + f.a, + None, + &f.owner, + vec![vec!["root".into(), "a".repeat(64)]], + ); + assert!(matches!( + f.accept(&invalid, None).await, + ArtifactOutcome::Rejected(_) + )); + let anchor = EventBuilder::new(Kind::Custom(9), "anchor") + .tags([Tag::parse(["h", &f.a.to_string()]).unwrap()]) + .sign_with_keys(&f.owner) + .unwrap(); + f.db.insert_event(f.community, &anchor, Some(f.a)) + .await + .unwrap(); + let root = || vec![vec!["root".to_string(), anchor.id.to_hex()]]; + let create = f.revision(d, "create", f.a, None, &f.owner, root()); + assert!(matches!( + f.accept(&create, None).await, + ArtifactOutcome::Accepted(_) + )); + // Later loss of the anchor does not block edits that keep it. + sqlx::query("UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=$2") + .bind(f.community.as_uuid()) + .bind(anchor.id.as_bytes().as_slice()) + .execute(&f.db.pool) + .await + .unwrap(); + let update = f.revision(d, "update", f.a, Some(&create), &f.peer, root()); + assert!(matches!( + f.accept(&update, None).await, + ArtifactOutcome::Accepted(_) + )); + let bad_delete = f.revision(d, "delete", f.a, Some(&update), &f.owner, vec![]); + assert!(matches!( + f.accept(&bad_delete, None).await, + ArtifactOutcome::Rejected(_) + )); + let other = Fixture::new().await; + assert!(!other + .db + .artifact_accepted(other.community, create.id.as_bytes()) + .await + .unwrap()); + let same_id = other.revision(d, "create", other.a, None, &other.owner, vec![]); + assert!(matches!( + other.accept(&same_id, None).await, + ArtifactOutcome::Accepted(_) + )); + assert_eq!( + other + .count( + &other.owner, + serde_json::json!({"artifact":"history","#d":[d]}) + ) + .await, + 1 + ); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn concurrent_artifacts_on_ttl_channel_commit() { + let f = Fixture::new().await; + sqlx::query("UPDATE channels SET ttl_seconds=60 WHERE community_id=$1 AND id=$2") + .bind(f.community.as_uuid()) + .bind(f.a) + .execute(&f.db.pool) + .await + .unwrap(); + let a = f.revision(Uuid::new_v4(), "create", f.a, None, &f.owner, vec![]); + let b = f.revision(Uuid::new_v4(), "create", f.a, None, &f.peer, vec![]); + let (a, b) = tokio::join!(f.accept(&a, None), f.accept(&b, None)); + assert!(matches!(a, ArtifactOutcome::Accepted(_))); + assert!(matches!(b, ArtifactOutcome::Accepted(_))); + let live: bool = sqlx::query_scalar( + "SELECT ttl_deadline>now() FROM channels WHERE community_id=$1 AND id=$2", + ) + .bind(f.community.as_uuid()) + .bind(f.a) + .fetch_one(&f.db.pool) + .await + .unwrap(); + assert!(live); +} diff --git a/crates/buzz-db/src/store/artifact_query.rs b/crates/buzz-db/src/store/artifact_query.rs new file mode 100644 index 00000000000..7ca2d7c308c --- /dev/null +++ b/crates/buzz-db/src/store/artifact_query.rs @@ -0,0 +1,72 @@ +//! Permission-scoped NIP-AR queries. Predicates precede pagination and counts. +use crate::{Db, Result}; +use buzz_core::{ + artifact::{ArtifactQuery, ArtifactView}, + CommunityId, StoredEvent, +}; +use sqlx::QueryBuilder; + +impl Db { + /// Execute a bounded query against current artifacts or revision history. + /// Permission is resolved in the query snapshot. + pub async fn query_artifacts( + &self, + community: CommunityId, + reader: &[u8], + query: &ArtifactQuery, + count: bool, + ) -> Result<(Vec, i64)> { + let mut tx = self.begin_event_write_transaction().await?; + sqlx::query("SET LOCAL statement_timeout='2s'") + .execute(&mut *tx) + .await?; + let mut qb = QueryBuilder::::new(if count { + "SELECT count(*) AS count FROM events e " + } else { + "SELECT e.id,e.pubkey,e.created_at,e.kind,e.tags,e.content,e.sig,e.received_at,e.channel_id FROM events e " + }); + if query.view == ArtifactView::Current { + qb.push("JOIN artifact_heads a ON a.community_id=e.community_id AND a.event_id=e.id AND NOT a.deleted "); + } + // The explicit kind lets the community/kind index exclude chat before + // the join, even for an unfiltered current-state page. + qb.push("WHERE e.community_id=").push_bind(community.as_uuid()).push(" AND e.kind=45010 AND e.deleted_at IS NULL AND EXISTS (SELECT 1 FROM channels c WHERE c.community_id=e.community_id AND c.id=e.channel_id AND c.deleted_at IS NULL AND (c.visibility='open' OR EXISTS (SELECT 1 FROM channel_members m WHERE m.community_id=c.community_id AND m.channel_id=c.id AND m.pubkey=") + .push_bind(reader.to_vec()).push(" AND m.removed_at IS NULL)))"); + for (name, values) in &query.tags { + // GIN-usable necessary prefilter; retain the positional recheck + // because JSON array containment alone also matches swapped values. + qb.push(" AND ("); + for (index, value) in values.iter().enumerate() { + if index != 0 { + qb.push(" OR "); + } + qb.push("e.tags @> ") + .push_bind(serde_json::json!([[name, value]])); + } + qb.push(")"); + qb.push(" AND EXISTS (SELECT 1 FROM jsonb_array_elements(e.tags) tag WHERE tag->>0=") + .push_bind(name) + .push(" AND tag->>1=ANY(") + .push_bind(values) + .push("))"); + } + if count { + let n: i64 = qb.build_query_scalar().fetch_one(&mut *tx).await?; + tx.commit().await?; + return Ok((vec![], n)); + } + qb.push(" ORDER BY e.received_at DESC,e.id ASC LIMIT ") + .push_bind(query.limit) + .push(" OFFSET ") + .push_bind(query.offset); + let rows = qb.build().fetch_all(&mut *tx).await?; + let mut events = vec![]; + for row in rows { + if let Some(event) = crate::event::row_to_stored_event(row)? { + events.push(event); + } + } + tx.commit().await?; + Ok((events, 0)) + } +} diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index cdd13d94c81..23b6d43049c 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -59,6 +59,8 @@ pub const CONTROL_PLANE_TABLES: &[&str] = &[ pub const EXPECTED_SCOPED_TABLES: &[&str] = &[ "api_tokens", "archived_identities", + "artifact_heads", + "artifact_revisions", "audit_log", "channel_members", "channels", @@ -107,6 +109,8 @@ pub const PURGE_SCOPED_TABLES: &[&str] = &[ "push_leases", "relay_invites", "delivery_log", + "artifact_heads", + "artifact_revisions", "events", "parameterized_event_watermarks", "git_repo_names", diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index 7031a60ca4c..3ec2cbf89dc 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -569,6 +569,12 @@ fn build_query_events_sql(q: &EventQuery) -> QueryBuilder { // Use unqualified column names when no join, qualified when joined. let col_prefix = if q.p_tag_hex.is_some() { "e." } else { "" }; + // Generic reads return only current artifact revisions, including delete + // tombstones; explicit revision IDs also read earlier revisions. + if q.ids.is_none() { + let table = if q.p_tag_hex.is_some() { "e" } else { "events" }; + qb.push(format!(" AND ({table}.kind <> 45010 OR EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); + } if let Some(ch) = q.channel_id { qb.push(format!(" AND {col_prefix}channel_id = ")) @@ -876,6 +882,12 @@ pub(crate) async fn count_events_on(conn: &mut sqlx::PgConnection, q: &EventQuer }; let col_prefix = if q.p_tag_hex.is_some() { "e." } else { "" }; + // Generic reads return only current artifact revisions, including delete + // tombstones; explicit revision IDs also read earlier revisions. + if q.ids.is_none() { + let table = if q.p_tag_hex.is_some() { "e" } else { "events" }; + qb.push(format!(" AND ({table}.kind <> 45010 OR EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); + } if let Some(ch) = q.channel_id { qb.push(format!(" AND {col_prefix}channel_id = ")) @@ -1109,6 +1121,13 @@ pub(crate) async fn soft_delete_event_and_update_thread_in_tx( .fetch_optional(&mut *tx) .await?; + // Relay-signed move removals are the source channel's only replay record. + if target.is_some_and(|(kind, _)| kind == 45011) { + return Err(DbError::InvalidData( + "artifact removal markers cannot be deleted".into(), + )); + } + if let Some((kind, Some(channel_id))) = target { if kind == KIND_CANVAS as i32 { let lock_key = event_replacement_lock_key( diff --git a/crates/buzz-db/src/store/mod.rs b/crates/buzz-db/src/store/mod.rs index 51b8717db16..79f0f078ab2 100644 --- a/crates/buzz-db/src/store/mod.rs +++ b/crates/buzz-db/src/store/mod.rs @@ -62,3 +62,11 @@ pub mod usage; pub mod user; /// Workflow, run, and approval persistence. pub mod workflow; + +/// NIP-AR artifact lifecycle and durable delivery. +pub mod artifact; + +mod artifact_query; + +#[cfg(test)] +mod artifact_postgres_tests; diff --git a/crates/buzz-relay/src/api/artifact.rs b/crates/buzz-relay/src/api/artifact.rs new file mode 100644 index 00000000000..bfe636c9086 --- /dev/null +++ b/crates/buzz-relay/src/api/artifact.rs @@ -0,0 +1,77 @@ +//! Raw artifact filters must be parsed before nostr::Filter discards extensions. +use super::api_error; +use crate::state::AppState; +use axum::{http::StatusCode, Json}; +use buzz_core::artifact::{route_filter, FilterRoute}; +use buzz_core::TenantContext; +use serde_json::Value; +use std::sync::Arc; + +/// Serve explicit artifact queries; `None` leaves the request to the generic +/// path (including existing extensions such as `#buzz-channel`). +pub(super) async fn query( + state: &Arc, + tenant: &TenantContext, + reader: &nostr::PublicKey, + filters: &[Value], + count: bool, +) -> Option, (StatusCode, Json)>> { + let mut artifact = false; + for filter in filters { + match route_filter(filter) { + FilterRoute::Generic => {} + FilterRoute::Artifact => artifact = true, + FilterRoute::Rejected(reason) => { + return Some(Err(api_error(StatusCode::BAD_REQUEST, reason))) + } + } + } + if !artifact { + return None; + } + Some(execute(state, tenant, reader, filters, count).await) +} + +async fn execute( + state: &Arc, + tenant: &TenantContext, + reader: &nostr::PublicKey, + filters: &[Value], + count: bool, +) -> Result, (StatusCode, Json)> { + // One filter avoids ambiguous OR-union count/page semantics. AND/OR inside + // exact predicates still supports cross-channel multi-value lists. + if filters.len() != 1 { + return Err(api_error( + StatusCode::BAD_REQUEST, + "artifact queries require exactly one filter", + )); + } + let query = buzz_core::artifact::parse_query(&filters[0]) + .map_err(|e| api_error(StatusCode::BAD_REQUEST, e))?; + // The query's statement_timeout is the enforced execution budget. + let (events, n) = state + .db + .query_artifacts(tenant.community(), reader.as_bytes(), &query, count) + .await + .map_err(|_| { + api_error( + StatusCode::SERVICE_UNAVAILABLE, + "artifact query failed or execution budget exceeded", + ) + })?; + if count { + return Ok(Json(serde_json::json!({"count":n}))); + } + let events = events + .into_iter() + .map(|e| serde_json::to_value(e.event)) + .collect::, _>>() + .map_err(|_| { + api_error( + StatusCode::INTERNAL_SERVER_ERROR, + "artifact serialization failed", + ) + })?; + Ok(Json(Value::Array(events))) +} diff --git a/crates/buzz-relay/src/api/artifact_postgres_tests.rs b/crates/buzz-relay/src/api/artifact_postgres_tests.rs new file mode 100644 index 00000000000..fda2e84ada5 --- /dev/null +++ b/crates/buzz-relay/src/api/artifact_postgres_tests.rs @@ -0,0 +1,289 @@ +// Exercises the production HTTP router and ingest gates for artifacts. +use super::postgres_tests::bridge_handler_test_state; +use super::*; +use nostr::{Event, EventBuilder, Keys, Kind, Tag}; +use serde_json::{json, Value}; +use uuid::Uuid; + +struct Fixture { + state: Arc, + pool: sqlx::PgPool, + host: String, + community: buzz_core::CommunityId, + home: Uuid, + private: Uuid, + owner: Keys, + peer: Keys, +} +impl Fixture { + async fn new() -> Self { + let state = bridge_handler_test_state() + .await + .expect("Postgres and Redis"); + let pool = sqlx::PgPool::connect(&crate::test_support::database_url()) + .await + .unwrap(); + let host = format!("artifact-review-{}.local", Uuid::new_v4()); + state.db.ensure_configured_community(&host).await.unwrap(); + state.db.ensure_future_partitions(1).await.unwrap(); + let tenant = crate::tenant::bind_community(&state.db, &host) + .await + .unwrap(); + let community = tenant.community(); + let home = Uuid::new_v4(); + let private = Uuid::new_v4(); + let owner = Keys::generate(); + for (id, visibility) in [(home, "open"), (private, "private")] { + sqlx::query("INSERT INTO channels(community_id,id,name,visibility,channel_type,created_by) VALUES($1,$2,$2::text,$3::channel_visibility,'stream',$4)") + .bind(community.as_uuid()).bind(id).bind(visibility).bind(owner.public_key().as_bytes().as_slice()).execute(&pool).await.unwrap(); + sqlx::query("INSERT INTO channel_members(community_id,channel_id,pubkey,role) VALUES($1,$2,$3,'owner')") + .bind(community.as_uuid()).bind(id).bind(owner.public_key().as_bytes().as_slice()).execute(&pool).await.unwrap(); + } + Self { + state, + pool, + host, + community, + home, + private, + owner, + peer: Keys::generate(), + } + } + fn revision(&self, id: Uuid, op: &str, prev: Option<&Event>) -> Event { + self.revision_as(&self.owner, self.home, id, op, prev) + } + fn revision_as( + &self, + key: &Keys, + home: Uuid, + id: Uuid, + op: &str, + prev: Option<&Event>, + ) -> Event { + let mut tags = vec![ + vec!["ar".into(), "1".into()], + vec!["d".into(), id.to_string()], + vec!["h".into(), home.to_string()], + vec!["type".into(), "buzz.task".into()], + vec!["op".into(), op.into()], + ]; + if op != "delete" { + tags.push(vec!["title".into(), "Review task".into()]); + } + if let Some(prev) = prev { + tags.push(vec!["prev".into(), prev.id.to_hex()]); + } + EventBuilder::new( + Kind::Custom(45010), + if op == "delete" { + "" + } else { + "review searchable" + }, + ) + .tags(tags.into_iter().map(|t| Tag::parse(t).unwrap())) + .sign_with_keys(key) + .unwrap() + } + async fn request(&self, path: &str, body: Value) -> (axum::http::StatusCode, Value) { + self.request_as(&self.owner, path, body).await + } + async fn request_as( + &self, + key: &Keys, + path: &str, + body: Value, + ) -> (axum::http::StatusCode, Value) { + use tower::ServiceExt; + let response = crate::router::build_router(self.state.clone()) + .oneshot( + axum::http::Request::builder() + .method("POST") + .uri(path) + .header("host", &self.host) + .header("content-type", "application/json") + .header("x-pubkey", key.public_key().to_hex()) + .body(axum::body::Body::from(body.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let status = response.status(); + let bytes = axum::body::to_bytes(response.into_body(), 1024 * 1024) + .await + .unwrap(); + (status, serde_json::from_slice(&bytes).unwrap()) + } + async fn publish(&self, event: &Event) { + let (status, body) = self.try_publish(event).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body["accepted"], true, "{body}"); + } + async fn try_publish(&self, event: &Event) -> (axum::http::StatusCode, Value) { + let key = if event.pubkey == self.peer.public_key() { + &self.peer + } else { + &self.owner + }; + self.request_as(key, "/events", json!(event)).await + } +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn project_filter_survives_artifact_query_hook() { + let f = Fixture::new().await; + let project = EventBuilder::new(Kind::Custom(30621), "") + .tags([ + Tag::parse(["d", "review-project"]).unwrap(), + Tag::parse(["buzz-channel", &f.home.to_string()]).unwrap(), + ]) + .sign_with_keys(&f.owner) + .unwrap(); + f.state + .db + .insert_event(f.community, &project, None) + .await + .unwrap(); + let filter = json!([{"kinds":[30621],"#buzz-channel":[f.home]}]); + let (status, body) = f.request("/query", filter.clone()).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body.as_array().unwrap().len(), 1); + assert_eq!(body[0]["id"], project.id.to_hex()); + let (status, body) = f.request("/count", filter).await; + assert!(status.is_success(), "{status}: {body}"); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn lifecycle_routes_duplicate_replay_delete_and_redaction() { + let f = Fixture::new().await; + let d = Uuid::new_v4(); + let create = f.revision(d, "create", None); + f.publish(&create).await; + let update = f.revision(d, "update", Some(&create)); + f.publish(&update).await; + f.publish(&create).await; + // The production ingest duplicate path must precede timestamp freshness. + let old = EventBuilder::new(Kind::Custom(45010), "old") + .tags(create.tags.iter().cloned()) + .custom_created_at(nostr::Timestamp::from(1)) + .sign_with_keys(&f.owner) + .unwrap(); + sqlx::query( + "INSERT INTO artifact_revisions(community_id,event_id,artifact_id) VALUES($1,$2,$3)", + ) + .bind(f.community.as_uuid()) + .bind(old.id.as_bytes().as_slice()) + .bind(Uuid::new_v4()) + .execute(&f.pool) + .await + .unwrap(); + f.publish(&old).await; + let delete = f.revision(d, "delete", Some(&update)); + f.publish(&delete).await; + // WebSocket replay's production filter-to-query and store seams. + let filter = serde_json::from_value(json!({"kinds":[45010],"#h":[f.home],"since":0})).unwrap(); + let query = crate::handlers::req::build_event_query_from_filter( + &filter, + f.owner.public_key().as_bytes(), + &f.state, + f.community, + ) + .await; + let replay = f.state.db.query_events(&query).await.unwrap(); + assert_eq!(replay.len(), 1); + assert_eq!(replay[0].event.id, delete.id); + for path in ["/query", "/count"] { + let (status, body) = f + .request(path, json!([{"artifact":"current","#d":[d]}])) + .await; + assert!(status.is_success(), "{body}"); + if path == "/query" { + assert_eq!(body, json!([])); + } else { + assert_eq!(body["count"], 0); + } + } + // Multi-character predicates are never silently dropped on the generic path. + let (status, _) = f + .request("/query", json!([{"kinds":[45010],"#project":["p"]}])) + .await; + assert_eq!(status, axum::http::StatusCode::BAD_REQUEST); + let (_, search) = f + .request( + "/query", + json!([{"kinds":[9],"search":"review searchable","#h":[f.home]}]), + ) + .await; + assert_eq!(search, json!([])); + let restore = f.revision(d, "restore", Some(&delete)); + f.publish(&restore).await; + let command = |kind| { + EventBuilder::new(Kind::Custom(kind), "") + .tags([ + Tag::parse(["e", &restore.id.to_hex()]).unwrap(), + Tag::parse(["h", &f.home.to_string()]).unwrap(), + ]) + .sign_with_keys(&f.owner) + .unwrap() + }; + let (status, body) = f.request("/events", json!(command(5))).await; + assert!(!status.is_success(), "{body}"); + assert!( + body.to_string() + .contains("artifacts cannot be deleted with kind 5"), + "{body}" + ); + // Redaction is the ordinary 9005 removal; the head stays editable. + f.publish(&command(9005)).await; + let current = json!([{"artifact":"current","#d":[d]}]); + let (_, body) = f.request("/query", current.clone()).await; + assert_eq!(body, json!([])); + let (_, body) = f + .request("/count", json!([{"artifact":"history","#d":[d]}])) + .await; + assert_eq!(body["count"], 3); + let edit = f.revision(d, "update", Some(&restore)); + f.publish(&edit).await; + let (_, body) = f.request("/query", current).await; + assert_eq!(body[0]["id"], edit.id.to_hex()); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn writes_use_channel_gates_for_home_and_move_source() { + let f = Fixture::new().await; + let d = Uuid::new_v4(); + let create = f.revision(d, "create", None); + f.publish(&create).await; + // Open-channel nonmembers may edit, as they may post kind 9. + let peer_edit = f.revision_as(&f.peer, f.home, d, "update", Some(&create)); + f.publish(&peer_edit).await; + let other = f.revision_as(&f.peer, f.private, Uuid::new_v4(), "create", None); + let (status, body) = f.try_publish(&other).await; + assert!(!status.is_success() || body["accepted"] == false, "{body}"); + let moved = f.revision_as(&f.owner, f.private, d, "move", Some(&peer_edit)); + f.publish(&moved).await; + let (_, removals) = f + .request("/query", json!([{"kinds":[45011],"#h":[f.home]}])) + .await; + assert_eq!(removals.as_array().unwrap().len(), 1, "{removals}"); + assert!(!removals.to_string().contains(&f.private.to_string())); + // The peer can write the destination but not the private source. + let back = f.revision_as(&f.peer, f.home, d, "move", Some(&moved)); + let (status, body) = f.try_publish(&back).await; + assert!(!status.is_success() || body["accepted"] == false, "{body}"); + assert!(body.to_string().contains("not a channel member"), "{body}"); + sqlx::query("UPDATE channels SET archived_at=now() WHERE community_id=$1 AND id=$2") + .bind(f.community.as_uuid()) + .bind(f.private) + .execute(&f.pool) + .await + .unwrap(); + let archived_source = f.revision_as(&f.owner, f.home, d, "move", Some(&moved)); + let (status, body) = f.try_publish(&archived_source).await; + assert!(!status.is_success() || body["accepted"] == false, "{body}"); + assert!(body.to_string().contains("archived"), "{body}"); +} diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index b5aa6d553b4..d7805693bef 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -1330,6 +1330,10 @@ async fn query_events_authed( // depth_limit, feed_types) that nostr::Filter silently drops. let raw_filters: Vec = serde_json::from_slice(body) .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; + if let Some(result) = super::artifact::query(state, tenant, &pubkey, &raw_filters, false).await + { + return result; + } let thread_windows = thread_window::parse(&raw_filters)?; let filters: Vec = raw_filters .iter() @@ -1951,6 +1955,11 @@ async fn count_events_authed( ) .await?; + let raw: Vec = serde_json::from_slice(body) + .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; + if let Some(result) = super::artifact::query(state, tenant, &pubkey, &raw, true).await { + return result; + } let filters: Vec = serde_json::from_slice(body) .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; crate::handlers::req::extract_channel_ids_from_filters_limited(&filters) @@ -2886,6 +2895,10 @@ fn ban_json(b: &buzz_db::moderation::BanRecord) -> Value { }) } +#[cfg(test)] +#[path = "artifact_postgres_tests.rs"] +mod artifact_postgres_tests; + #[cfg(test)] mod postgres_tests { use super::*; diff --git a/crates/buzz-relay/src/api/mod.rs b/crates/buzz-relay/src/api/mod.rs index 460f7f70b84..545954fb368 100644 --- a/crates/buzz-relay/src/api/mod.rs +++ b/crates/buzz-relay/src/api/mod.rs @@ -532,3 +532,5 @@ mod parse_query_tests { ); } } + +mod artifact; diff --git a/crates/buzz-relay/src/handlers/artifact.rs b/crates/buzz-relay/src/handlers/artifact.rs new file mode 100644 index 00000000000..e6d568777c3 --- /dev/null +++ b/crates/buzz-relay/src/handlers/artifact.rs @@ -0,0 +1,148 @@ +//! NIP-AR lifecycle entry point. +use super::ingest::{IngestAuth, IngestError, IngestResult}; +use crate::state::AppState; +use buzz_core::artifact::ArtifactOp; +use buzz_core::{StoredEvent, TenantContext}; +use buzz_db::artifact::ArtifactOutcome; +use nostr::Event; +use std::sync::Arc; + +/// Accept a revision whose home channel already passed ingest's kind-9 write +/// gates. A move must also be writable in its source channel. +pub(crate) async fn accept( + state: &Arc, + tenant: &TenantContext, + event: &Event, + auth: &IngestAuth, +) -> Result { + let env = buzz_core::artifact::validate(event) + .map_err(|e| IngestError::Rejected(format!("invalid: {e}")))?; + verify_revision_auth(event)?; + let source = if env.op == ArtifactOp::Move { + let source = state + .db + .artifact_home(tenant.community(), env.id) + .await + .map_err(|e| IngestError::Internal(format!("artifact storage: {e}")))? + .ok_or_else(|| { + IngestError::CanvasConflict("conflict: artifact head unavailable".into()) + })?; + super::ingest::check_channel_write(tenant, state, auth, source) + .await + .map_err(IngestError::Rejected)?; + Some(source) + } else { + None + }; + let outcome = state + .db + .accept_artifact( + tenant.community(), + event, + &env, + source, + &state.relay_keypair, + ) + .await + .map_err(|e| IngestError::Internal(format!("artifact storage: {e}")))?; + match outcome { + ArtifactOutcome::Accepted(events) => { + for stored in &events { + publish(state, tenant, stored).await; + } + } + ArtifactOutcome::Duplicate => {} + ArtifactOutcome::Conflict(message) => { + return Err(IngestError::CanvasConflict(format!("conflict: {message}"))) + } + ArtifactOutcome::Rejected(message) => { + return Err(IngestError::Rejected(format!("invalid: {message}"))) + } + } + Ok(IngestResult { + event_id: event.id.to_hex(), + accepted: true, + message: String::new(), + }) +} + +/// Live delivery only; unlike kind 9 there are no audit, workflow, or thread +/// side effects. Missed deliveries are recovered by replaying stored events. +async fn publish(state: &Arc, tenant: &TenantContext, stored: &StoredEvent) { + let Some(channel) = stored.channel_id else { + return; + }; + state.mark_local_event(tenant.community(), &stored.event.id); + if let Err(e) = state + .pubsub + .publish_event( + tenant, + buzz_pubsub::EventTopic::Channel(channel), + &stored.event, + ) + .await + { + state + .local_event_ids + .invalidate(&(tenant.community(), stored.event.id.to_bytes())); + tracing::warn!(error=%e, "artifact Redis publish failed"); + } + super::event::fan_out_event_to_local_subscribers(state, tenant.community(), stored).await; +} + +fn verify_revision_auth(event: &Event) -> Result<(), IngestError> { + for tag in event.tags.iter().filter(|t| t.as_slice()[0] == "auth") { + let json = serde_json::to_string(tag.as_slice()) + .map_err(|e| IngestError::Internal(e.to_string()))?; + buzz_sdk::nip_oa::verify_auth_tag_for_auth_event( + &json, + &event.pubkey, + event.created_at.as_secs(), + ) + .map_err(|_| IngestError::Rejected("invalid: unverified artifact auth tag".into()))?; + // Unlike connection admission, an attestation carried by a revision + // must satisfy its kind clauses against that revision. + if tag.as_slice()[2].split('&').any(|clause| { + clause + .strip_prefix("kind=") + .is_some_and(|value| value.parse::().ok() != Some(event.kind.as_u16())) + }) { + return Err(IngestError::Rejected( + "invalid: artifact auth kind mismatch".into(), + )); + } + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use nostr::{EventBuilder, Keys, Kind}; + + #[test] + fn revision_auth_checks_every_condition_not_just_admission() { + let owner = Keys::generate(); + let agent = Keys::generate(); + for (conditions, accepted) in [ + ("", true), + ("kind=45010", true), + ("kind=9", false), + ("kind=45010&kind=9", false), + ("created_at<1", false), + ] { + let json = buzz_sdk::nip_oa::compute_auth_tag(&owner, &agent.public_key(), conditions) + .unwrap(); + let tag = buzz_sdk::nip_oa::parse_auth_tag(&json).unwrap(); + let event = EventBuilder::new(Kind::Custom(45010), "") + .tags([tag]) + .sign_with_keys(&agent) + .unwrap(); + assert_eq!( + verify_revision_auth(&event).is_ok(), + accepted, + "{conditions}" + ); + } + } +} diff --git a/crates/buzz-relay/src/handlers/ingest.rs b/crates/buzz-relay/src/handlers/ingest.rs index 9b95218a026..794d768b46b 100644 --- a/crates/buzz-relay/src/handlers/ingest.rs +++ b/crates/buzz-relay/src/handlers/ingest.rs @@ -499,7 +499,7 @@ fn map_push_accept_error(error: super::push_lease::AcceptError) -> IngestError { fn required_scope_for_kind(kind: u32, event: &Event) -> Result { match kind { KIND_PROFILE => Ok(Scope::UsersWrite), - KIND_TEXT_NOTE | KIND_LONG_FORM => Ok(Scope::MessagesWrite), + KIND_TEXT_NOTE | KIND_LONG_FORM | buzz_core::kind::KIND_ARTIFACT => Ok(Scope::MessagesWrite), KIND_CONTACT_LIST | KIND_READ_STATE | KIND_USER_STATUS | KIND_AGENT_ENGRAM | KIND_EVENT_REMINDER | KIND_PERSONA | KIND_TEAM | KIND_MANAGED_AGENT | KIND_PRIVATE_MANAGED_AGENT | KIND_TEAM_CATALOG | super::push_lease::KIND_PUSH_LEASE => { @@ -836,6 +836,34 @@ pub(crate) async fn check_channel_membership( } } +/// The kind-9 channel write gates (token scope, membership or open channel, +/// not archived) for a channel other than the event's own `h`. +pub(crate) async fn check_channel_write( + tenant: &TenantContext, + state: &AppState, + auth: &IngestAuth, + ch_id: Uuid, +) -> Result<(), String> { + check_token_channel_access(auth, ch_id)?; + let channel = state + .db + .get_channel_for_event_write(tenant.community(), ch_id) + .await + .ok(); + check_channel_membership( + tenant, + state, + ch_id, + &auth.pubkey().to_bytes(), + channel.as_ref(), + ) + .await?; + if channel.is_some_and(|ch| ch.archived_at.is_some()) { + return Err("invalid: channel is archived".into()); + } + Ok(()) +} + fn check_token_channel_access(auth: &IngestAuth, channel_id: Uuid) -> Result<(), String> { if let Some(allowed) = auth.channel_ids() { if !allowed.contains(&channel_id) { @@ -2313,6 +2341,33 @@ async fn ingest_event_inner( } let event = std::sync::Arc::try_unwrap(event).unwrap_or_else(|arc| (*arc).clone()); + if kind_u32 == buzz_core::kind::KIND_ARTIFACT + && event.pubkey == *auth.pubkey() + && state + .db + .artifact_accepted(tenant.community(), event.id.as_bytes()) + .await + .map_err(|e| IngestError::Internal(e.to_string()))? + { + emit( + tracer, + TraceAction::WriteDuplicate { + msg_id: msg_id_label(event.id.as_bytes()), + channel: channel_label( + extract_channel_id(&event) + .ok_or_else(|| IngestError::Rejected("invalid: missing home".into()))?, + ), + claimed_community: claimed_community_from_event(&event), + }, + state_for_request(tenant, auth.pubkey()), + ); + return Ok(IngestResult { + event_id: event_id_hex, + accepted: true, + message: String::new(), + }); + } + const MAX_TIMESTAMP_DRIFT_SECS: i64 = 900; // ±15 minutes let now = chrono::Utc::now().timestamp(); let event_ts = event.created_at.as_secs() as i64; @@ -2797,6 +2852,24 @@ async fn ingest_event_inner( } } + // Artifact revisions passed the same home-channel write gates as kind 9 + // above; they are stored and published without conversation side effects. + if kind_u32 == buzz_core::kind::KIND_ARTIFACT { + let result = super::artifact::accept(state, tenant, &event, &auth).await?; + if let Some(ch_id) = channel_id { + emit( + tracer, + TraceAction::WriteInsert { + msg_id: msg_id_label(event.id.as_bytes()), + channel: channel_label(ch_id), + claimed_community: claimed_community_from_event(&event), + }, + state_for_request(tenant, auth.pubkey()), + ); + } + return Ok(result); + } + // NIP-09: kind:5 may reference targets via `e` tag (regular events) OR // `a` tag (addressable/parameterized-replaceable events like kind:30620). if kind_u32 == KIND_NIP29_DELETE_EVENT || kind_u32 == KIND_DELETION { diff --git a/crates/buzz-relay/src/handlers/mod.rs b/crates/buzz-relay/src/handlers/mod.rs index 2f4aa00b595..e4eea46906c 100644 --- a/crates/buzz-relay/src/handlers/mod.rs +++ b/crates/buzz-relay/src/handlers/mod.rs @@ -66,3 +66,6 @@ pub fn resolve_ttl(event: &nostr::Event, ephemeral_ttl_override: Option) -> (ttl, _) => ttl, } } + +/// NIP-AR artifact lifecycle. +pub mod artifact; diff --git a/crates/buzz-relay/src/handlers/side_effects.rs b/crates/buzz-relay/src/handlers/side_effects.rs index c2131b1a558..a2e260c9fc6 100644 --- a/crates/buzz-relay/src/handlers/side_effects.rs +++ b/crates/buzz-relay/src/handlers/side_effects.rs @@ -403,6 +403,11 @@ pub async fn validate_standard_deletion_event( .await? .ok_or_else(|| anyhow::anyhow!("target event not found"))?; + if matches!(target_event.event.kind.as_u16(), 45010 | 45011) { + anyhow::bail!( + "artifacts cannot be deleted with kind 5; use op=delete or kind 9005 redaction" + ); + } let target_author = effective_message_author(&target_event.event, &state.relay_keypair.public_key()); if target_author != actor_bytes @@ -739,6 +744,11 @@ pub async fn validate_admin_event( } _ => {} // Same channel — OK } + if target_event.event.kind.as_u16() == 45011 { + return Err(anyhow::anyhow!( + "artifact removal markers cannot be deleted" + )); + } // Check if actor is the event author. // For relay-signed REST messages, the real author is in the p tag. diff --git a/crates/buzz-relay/src/nip11.rs b/crates/buzz-relay/src/nip11.rs index 043458b4bca..218647f8625 100644 --- a/crates/buzz-relay/src/nip11.rs +++ b/crates/buzz-relay/src/nip11.rs @@ -34,6 +34,8 @@ pub struct RelayInfo { /// Host-bound atomic read-state snapshot capability; absent on unresolved hosts. #[serde(skip_serializing_if = "Option::is_none")] pub read_state_snapshot: Option, + /// NIP-AR artifact query transport and enforced resource limits. + pub artifacts: serde_json::Value, /// Relay operator's public key (hex), if published. pub pubkey: Option, /// Contact address for the relay operator. @@ -192,7 +194,7 @@ impl RelayInfo { supported_nips.push(NIP_RELAY_MEMBERSHIP); } - let mut supported_extensions = vec!["nip-er".to_string()]; + let mut supported_extensions = vec!["nip-er".to_string(), "nip-ar".to_string()]; let gif = gif_provider.map(|provider| { supported_extensions.push("buzz-gif".to_string()); GifDescriptor { @@ -207,6 +209,20 @@ impl RelayInfo { description: "Buzz — private team communication relay".to_string(), icon: icon.filter(|s| !s.is_empty()).map(|s| s.to_string()), read_state_snapshot: None, + artifacts: serde_json::json!({ + "version": 1, "revision_kind": 45010, "removal_kind": 45011, + "query": "/query", "count": "/count", + "modes": ["current", "history"], + "websocket_query_extensions": false, + "max_tags": buzz_core::artifact::MAX_TAGS, + "max_tag_name_bytes": buzz_core::artifact::MAX_TAG_NAME_BYTES, + "max_tag_value_bytes": buzz_core::artifact::MAX_TAG_VALUE_BYTES, + "max_tag_bytes": buzz_core::artifact::MAX_TAG_BYTES, + "max_predicates": buzz_core::artifact::MAX_PREDICATES, + "max_values": buzz_core::artifact::MAX_QUERY_VALUES, + "max_page_size": buzz_core::artifact::MAX_PAGE_SIZE, + "max_offset": buzz_core::artifact::MAX_OFFSET, "max_filters": 1, "query_timeout_ms": 2000 + }), pubkey: None, contact: None, supported_nips, diff --git a/crates/buzz-relay/src/protocol.rs b/crates/buzz-relay/src/protocol.rs index 5832a72a879..294e4f6f51f 100644 --- a/crates/buzz-relay/src/protocol.rs +++ b/crates/buzz-relay/src/protocol.rs @@ -38,6 +38,24 @@ pub enum ClientMessage { Auth(Event), } +/// Artifact queries are HTTP-only; reject rather than silently drop their +/// predicates on the generic WebSocket path. +fn reject_artifact_query_filters(filters: &[serde_json::Value]) -> Result<()> { + use buzz_core::artifact::{route_filter, FilterRoute}; + for filter in filters { + match route_filter(filter) { + FilterRoute::Generic => {} + FilterRoute::Artifact => { + return Err(RelayError::InvalidMessage( + "artifact queries require HTTP /query or /count".into(), + )) + } + FilterRoute::Rejected(reason) => return Err(RelayError::InvalidMessage(reason.into())), + } + } + Ok(()) +} + impl ClientMessage { /// Parse a raw JSON WebSocket frame into a [`ClientMessage`]. pub fn parse(raw: &str) -> Result { @@ -91,6 +109,8 @@ impl ClientMessage { ))); } let filter_values = &arr[2..]; + reject_artifact_query_filters(filter_values)?; + // Enforce NIP-11 advertised max_filters: 10 if filter_values.len() > MAX_FILTERS_PER_REQ { return Err(RelayError::InvalidMessage(format!( @@ -163,6 +183,8 @@ impl ClientMessage { ))); } let filter_values = &arr[2..]; + reject_artifact_query_filters(filter_values)?; + if filter_values.len() > MAX_FILTERS_PER_REQ { return Err(RelayError::InvalidMessage(format!( "COUNT contains {} filters, maximum is {MAX_FILTERS_PER_REQ}", @@ -253,6 +275,29 @@ impl RelayMessage { #[cfg(test)] mod tests { + #[test] + fn non_artifact_extensions_keep_existing_ws_behavior() { + for command in ["REQ", "COUNT"] { + let request = serde_json::json!([command, "projects", {"kinds":[30621],"#buzz-channel":["channel"]}]); + assert!(super::ClientMessage::parse(&request.to_string()).is_ok()); + } + } + + #[test] + fn artifact_live_filters_preserved_extensions_rejected() { + assert!(super::ClientMessage::parse( + r##"["REQ","ar",{"kinds":[45010,45011],"#h":["channel"]}]"## + ) + .is_ok()); + for request in [ + r##"["REQ","ar",{"artifact":"history","#d":["artifact"]}]"##, + r##"["REQ","ar",{"kinds":[45010],"#project":["project"]}]"##, + r##"["COUNT","ar",{"artifact":"current"}]"##, + ] { + assert!(super::ClientMessage::parse(request).is_err(), "{request}"); + } + } + use super::*; use buzz_core::test_helpers::make_event; use nostr::{EventBuilder, Keys, Kind}; diff --git a/crates/buzz-search/Cargo.toml b/crates/buzz-search/Cargo.toml index 6c6b9ada221..2d491e70565 100644 --- a/crates/buzz-search/Cargo.toml +++ b/crates/buzz-search/Cargo.toml @@ -17,4 +17,5 @@ tracing = { workspace = true } metrics = { workspace = true } [dev-dependencies] +buzz-db = { workspace = true } tokio = { workspace = true } diff --git a/crates/buzz-search/src/query.rs b/crates/buzz-search/src/query.rs index bd95e8cdbbc..a2bed24e22f 100644 --- a/crates/buzz-search/src/query.rs +++ b/crates/buzz-search/src/query.rs @@ -252,6 +252,9 @@ pub async fn search(pool: &PgPool, query: &SearchQuery) -> Result 45010 OR EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id=events.community_id AND ah.event_id=events.id AND NOT ah.deleted))"); + // Channel scope — see `ChannelScope` doc for the four-case mapping. The // emitted SQL fragments are identical to the legacy 2x2 tuple for the // three carry-over cases; `ChannelLessOnly` is the new fence that the diff --git a/crates/buzz-search/tests/postgres_fts_integration.rs b/crates/buzz-search/tests/postgres_fts_integration.rs index 175a01aaaa3..dcff99a188e 100644 --- a/crates/buzz-search/tests/postgres_fts_integration.rs +++ b/crates/buzz-search/tests/postgres_fts_integration.rs @@ -1,10 +1,9 @@ //! Integration tests for community-scoped Postgres FTS. //! -//! Run with a local PG: `BUZZ_TEST_DATABASE_URL=postgres://buzz:buzz_dev@localhost:5432/buzz cargo test -p buzz-search --tests -- --include-ignored` -//! -//! Each test creates a uniquely-named schema, applies every FTS-affecting -//! migration in order, exercises a scenario, and drops it. Tests are -//! parallel-safe. +//! Run through `scripts/postgres-test-run.sh -p buzz-search --tests` for a +//! fresh desired-state database per test. Outside that lane, point +//! BUZZ_TEST_DATABASE_URL at a disposable dedicated database; fixtures use +//! the production migrator and unique community IDs. use buzz_core::{ kind::{ @@ -14,96 +13,30 @@ use buzz_core::{ CommunityId, }; use buzz_search::{ChannelScope, SearchQuery, SearchService}; -use sqlx::{postgres::PgPoolOptions, Executor, PgPool}; +use sqlx::{postgres::PgPoolOptions, PgPool}; use uuid::Uuid; const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; -const MIGRATION_0001_SQL: &str = include_str!("../../../migrations/0001_initial_schema.sql"); -const MIGRATION_0002_SQL: &str = include_str!("../../../migrations/0002_git_repo_names.sql"); -const MIGRATION_0003_SQL: &str = include_str!("../../../migrations/0003_community_icon.sql"); -const MIGRATION_0004_SQL: &str = include_str!("../../../migrations/0004_events_tags_gin.sql"); -const MIGRATION_0005_SQL: &str = include_str!("../../../migrations/0005_agent_turn_metric_fts.sql"); -const MIGRATION_0006_SQL: &str = include_str!("../../../migrations/0006_moderation.sql"); -const MIGRATION_0007_SQL: &str = include_str!("../../../migrations/0007_nip_rs_retention.sql"); -const MIGRATION_0008_SQL: &str = - include_str!("../../../migrations/0008_fresh_install_search_allowlist.sql"); -const MIGRATION_0014_SQL: &str = include_str!("../../../migrations/0014_push_lease_fts.sql"); -const MIGRATION_0033_SQL: &str = - include_str!("../../../migrations/0033_private_managed_agent_fts.sql"); - async fn setup() -> (PgPool, String) { let url = std::env::var("BUZZ_TEST_DATABASE_URL").unwrap_or_else(|_| TEST_DB_URL.to_string()); - let schema = format!("fts_test_{}", Uuid::new_v4().simple()); - // Connect to the default schema first to create the test schema. - let admin_pool = PgPoolOptions::new() - .max_connections(1) - .connect(&url) - .await - .expect("connect"); - let create_sql = format!("CREATE SCHEMA \"{schema}\""); - sqlx::query(sqlx::AssertSqlSafe(create_sql)) - .execute(&admin_pool) - .await - .expect("create schema"); - admin_pool.close().await; - - // Connect with search_path set so the migration's CREATE TABLE lands here. - let url_with_search_path = format!("{url}?options=-c%20search_path%3D{schema}"); + // The PostgreSQL lane supplies an isolated desired-state database. Outside + // that lane use a dedicated database with the production migrations. let pool = PgPoolOptions::new() - .max_connections(2) - .connect(&url_with_search_path) - .await - .expect("connect with search_path"); - // Apply the full migration chain in order so the test schema exactly matches - // production. Future FTS-affecting migrations must be added here. - pool.execute(MIGRATION_0001_SQL) - .await - .expect("apply 0001 migration"); - pool.execute(MIGRATION_0002_SQL) - .await - .expect("apply 0002 migration"); - pool.execute(MIGRATION_0003_SQL) - .await - .expect("apply 0003 migration"); - pool.execute(MIGRATION_0004_SQL) - .await - .expect("apply 0004 migration"); - pool.execute(MIGRATION_0005_SQL) - .await - .expect("apply 0005 migration"); - pool.execute(MIGRATION_0006_SQL) - .await - .expect("apply 0006 migration"); - pool.execute(MIGRATION_0007_SQL) - .await - .expect("apply 0007 migration"); - pool.execute(MIGRATION_0008_SQL) - .await - .expect("apply 0008 migration"); - pool.execute(MIGRATION_0014_SQL) - .await - .expect("apply 0014 migration"); - pool.execute(MIGRATION_0033_SQL) + .max_connections(5) + .connect(&url) .await - .expect("apply 0033 migration"); - (pool, schema) + .expect("connect test database"); + if std::env::var("BUZZ_TEST_SCHEMA_MODE").as_deref() != Ok("desired") { + buzz_db::migration::run_migrations(&pool) + .await + .expect("apply production migrations"); + } + (pool, String::new()) } -async fn teardown(pool: PgPool, schema: &str) { +async fn teardown(pool: PgPool, _schema: &str) { pool.close().await; - let admin_pool = PgPoolOptions::new() - .max_connections(1) - .connect( - &std::env::var("BUZZ_TEST_DATABASE_URL").unwrap_or_else(|_| TEST_DB_URL.to_string()), - ) - .await - .expect("reconnect for drop"); - let drop_sql = format!("DROP SCHEMA \"{schema}\" CASCADE"); - sqlx::query(sqlx::AssertSqlSafe(drop_sql)) - .execute(&admin_pool) - .await - .expect("drop schema"); - admin_pool.close().await; + // The test runner owns database cleanup. } /// Insert a community row, return its id. @@ -1512,3 +1445,63 @@ async fn p_gated_persistent_kinds_have_storage_null_tsvector() { teardown(pool, &schema).await; } + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn artifact_search_uses_only_live_current_head() { + let (pool, schema) = setup().await; + let community = mk_community(&pool, "artifact-search.example").await; + let pk = rand_bytes32(); + let old = rand_bytes32(); + let head = rand_bytes32(); + let artifact = Uuid::new_v4(); + insert_event( + &pool, + community, + old, + pk, + 45010, + "artifactneedle old", + None, + 1700000000, + ) + .await; + insert_event( + &pool, + community, + head, + pk, + 45010, + "artifactneedle current", + None, + 1700000001, + ) + .await; + sqlx::query("INSERT INTO artifact_heads(community_id,artifact_id,event_id,channel_id,artifact_type) VALUES($1,$2,$3,$4,'buzz.task')") + .bind(community.as_uuid()).bind(artifact).bind(head.as_slice()).bind(Uuid::new_v4()).execute(&pool).await.unwrap(); + let svc = SearchService::new(pool.clone()); + let query = SearchQuery { + community, + q: "artifactneedle".into(), + channel_scope: ChannelScope::Any, + kinds: Some(vec![45010]), + authors: None, + since: None, + until: None, + page: 1, + per_page: 10, + mode: buzz_search::SearchMode::FullText, + }; + let result = svc.search(&query).await.unwrap(); + assert_eq!(result.hits.len(), 1); + assert_eq!(result.hits[0].event_id, head); + sqlx::query("UPDATE artifact_heads SET deleted=true WHERE community_id=$1 AND artifact_id=$2") + .bind(community.as_uuid()) + .bind(artifact) + .execute(&pool) + .await + .unwrap(); + let result = svc.search(&query).await.unwrap(); + assert!(result.hits.is_empty()); + teardown(pool, &schema).await; +} diff --git a/docs/nips/NIP-AR.md b/docs/nips/NIP-AR.md index 6277fb4396e..18e5e0d3a46 100644 --- a/docs/nips/NIP-AR.md +++ b/docs/nips/NIP-AR.md @@ -56,7 +56,7 @@ An artifact is identified by `(community, d)`, independently of its author or ho Creation requires an unused `d`, `op=create`, and no `prev`. Every later revision MUST name the current accepted event in `prev`. The relay checks authorization and advances the current revision atomically: two competing edits cannot both succeed. Timestamps do not choose the winner. -Resubmitting an already accepted event succeeds without applying it again, even if a later revision has since become current, so a client that lost an acknowledgement never sees a false conflict. The response reveals nothing the sender can no longer read. A different event whose `prev` is not the current revision is a conflict. A conflict response MAY name the current revision only when the sender can read it. On conflict, clients fetch the current accessible revision and reconcile their changes rather than blindly retrying. An `update` keeps `h`; a `move` changes it. +Resubmitting an already accepted event succeeds without applying it again, even if a later revision has since become current, so a client that lost an acknowledgement never sees a false conflict. The response reveals nothing the sender can no longer read. A different event whose `prev` is not the current revision is a conflict. On conflict, clients fetch the current revision and reconcile their changes rather than resubmitting with only `prev` replaced, which would overwrite the competing edit. An `update` keeps `h`; a `move` changes it. A `create` whose `d` is already in use fails. When the sender cannot read the existing artifact, the error reveals only that the identity is taken, never its home, type, title, or current revision. Clients MAY derive `d` deterministically, for example from the message an artifact tracks, and accept that disclosure. @@ -66,15 +66,14 @@ On creation or when changed from the previous accepted revision, `root` MUST ide Channel read permission governs every artifact read, including lookups, history, search, previews, counts, and live updates. Threads inherit their channel's audience. Membership and visibility changes take effect on subsequent reads and deliveries. -Write permission means permission to post a kind-9 message in the channel, including authentication, token restrictions, moderation, and archive checks. This includes agents and, where channel policy permits, nonmembers of open channels. Artifact writes use `messages:write`, the scope that admits kind-9 posting. +Write permission means permission to post a kind-9 message in the channel, including authentication, token restrictions, moderation, and archive checks. Relays apply these checks the same way, and at the same point, as for a kind-9 message. This includes agents and, where channel policy permits, nonmembers of open channels. Artifact writes use `messages:write`, the scope that admits kind-9 posting. -| Operation | Required permission at commit time | +| Operation | Required permission | | --- | --- | -| Create or update | Write in the home channel. | +| Create, update, delete, or restore | Write in the home channel. | | Move | Write in both source and destination. | -| Delete or restore | Write in the home channel, and either being the artifact's author (the signer of its `create` revision) or holding the channel's owner or admin role. In a DM, write alone suffices. | -These rules apply to every type. Apart from deletion and restoration, they do not depend on who created the artifact. DM participants are peers. The relay determines channel type from the stored channel, not from an artifact tag. +These rules apply to every type and do not depend on who created the artifact. DM participants are peers. The relay determines channel type from the stored channel, not from an artifact tag. Relationship tags organize work without changing access. Linking a task to a project does not move it or share it. References to repositories or other services retain those services' access and action permissions. Errors MUST NOT disclose inaccessible artifact details. @@ -82,7 +81,7 @@ Relationship tags organize work without changing access. Linking a task to a pro A move publishes the current snapshot into the destination under the same `d`. Its `root` must be absent or belong to the destination. Clients MUST show the destination audience and the information being shared before confirmation. -The relay MUST atomically advance the current revision and durably record both arrival in the destination and removal from the source. The source receives a separate relay-authenticated removal identifying the artifact, source scope, and replay position, without destination metadata or content. Both changes must survive disconnects and crashes. Relays unable to provide this MUST reject moves. +The relay MUST atomically advance the current revision and store both arrival in the destination and removal from the source. The source receives a separate relay-authenticated removal identifying the artifact, source scope, and replay position, without destination metadata or content. Both are stored events, so a client that misses live delivery recovers them by replaying the channel. Relays unable to provide this MUST reject moves. Earlier revisions remain under their original channels' access rules; a move never lets the destination read revisions from the source. The destination can load current state without reading earlier revisions. Conversation messages stay where they are. @@ -90,7 +89,7 @@ Deletion is a soft delete: a revision with `op=delete`, empty content, and no `t ## Moderation and retention -Relays MUST reject kind-5 deletion requests targeting artifact revisions with an explicit reason. A NIP-29 kind-9005 removal targeting a revision redacts it instead of deleting it: the relay withholds that revision's title, content, and client-defined tags from every surface, including history, search, and live delivery, and serves a relay-authenticated removal marker that keeps its `d`, `h`, `type`, `op`, and `prev`. Redaction never changes which revision is current, so a redacted current revision can still be edited, deleted, or restored with a new revision. Revision history is subject to the community's retention policy; expiring earlier revisions MUST NOT remove the current revision or break `prev` checks. +Relays MUST reject kind-5 deletion requests targeting artifact revisions with an explicit reason. A NIP-29 kind-9005 removal redacts a revision under the channel's ordinary message-removal rules: the relay withholds that revision from every surface, including lookups, history, and search. The relay still records the revision as accepted, so redaction never changes which revision is current, and its ID remains valid as `prev`. Redacting an artifact's current revision retires the artifact: it is absent from current-state queries, and clients that do not already hold that revision's ID cannot edit it further. Revision history is subject to the community's retention policy; expiring earlier revisions MUST NOT remove the current revision or break `prev` checks. ## Current state and queries diff --git a/migrations/0052_channel_artifacts.sql b/migrations/0052_channel_artifacts.sql new file mode 100644 index 00000000000..32380529b5c --- /dev/null +++ b/migrations/0052_channel_artifacts.sql @@ -0,0 +1,41 @@ +-- NIP-AR current heads and acceptance ledger are independent of event retention. +CREATE TABLE artifact_heads ( + community_id UUID NOT NULL REFERENCES communities(id), + artifact_id UUID NOT NULL, + event_id BYTEA NOT NULL CHECK (length(event_id) = 32), + channel_id UUID NOT NULL, + artifact_type TEXT NOT NULL, + root BYTEA, + deleted BOOLEAN NOT NULL DEFAULT false, + PRIMARY KEY (community_id, artifact_id) +); +CREATE INDEX artifact_heads_event ON artifact_heads (community_id, event_id); +-- Every accepted revision ID, so replays stay idempotent after redaction or +-- retention. `position` records acceptance order. +CREATE TABLE artifact_revisions ( + community_id UUID NOT NULL REFERENCES communities(id), + event_id BYTEA NOT NULL CHECK (length(event_id) = 32), + artifact_id UUID NOT NULL, + position BIGINT GENERATED ALWAYS AS IDENTITY, + PRIMARY KEY (community_id, event_id) +); + +SELECT attach_community_write_fence('artifact_heads'); +SELECT attach_community_write_fence('artifact_revisions'); + +-- Row retention may expire earlier payloads, never the live CAS head. Partition +-- retirement must also preserve head payloads; this relay does not drop partitions. +-- A redacted (soft-deleted) head payload may still be purged. +CREATE FUNCTION retain_current_artifact() RETURNS TRIGGER LANGUAGE plpgsql AS $$ +BEGIN + IF OLD.kind = 45010 AND OLD.deleted_at IS NULL AND EXISTS ( + SELECT 1 FROM artifact_heads h + WHERE h.community_id=OLD.community_id AND h.event_id=OLD.id + ) THEN + RETURN NULL; + END IF; + RETURN OLD; +END; +$$; +CREATE TRIGGER retain_current_artifact BEFORE DELETE ON events +FOR EACH ROW EXECUTE FUNCTION retain_current_artifact(); diff --git a/schema/schema.sql b/schema/schema.sql index 37112cc6702..af72f819271 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1998,3 +1998,45 @@ CREATE TABLE storage_accounting_snapshots ( INSERT INTO _operator_global_tables (table_name, reason) VALUES ('storage_accounting_snapshots', 'deployment-global completed media accounting handoff'); + +-- NIP-AR current heads and acceptance ledger are independent of event retention. +CREATE TABLE artifact_heads ( + community_id UUID NOT NULL REFERENCES communities(id), + artifact_id UUID NOT NULL, + event_id BYTEA NOT NULL CHECK (length(event_id) = 32), + channel_id UUID NOT NULL, + artifact_type TEXT NOT NULL, + root BYTEA, + deleted BOOLEAN NOT NULL DEFAULT false, + PRIMARY KEY (community_id, artifact_id) +); +CREATE INDEX artifact_heads_event ON artifact_heads (community_id, event_id); +-- Every accepted revision ID, so replays stay idempotent after redaction or +-- retention. `position` records acceptance order. +CREATE TABLE artifact_revisions ( + community_id UUID NOT NULL REFERENCES communities(id), + event_id BYTEA NOT NULL CHECK (length(event_id) = 32), + artifact_id UUID NOT NULL, + position BIGINT GENERATED ALWAYS AS IDENTITY, + PRIMARY KEY (community_id, event_id) +); + +SELECT attach_community_write_fence('artifact_heads'); +SELECT attach_community_write_fence('artifact_revisions'); + +-- Row retention may expire earlier payloads, never the live CAS head. Partition +-- retirement must also preserve head payloads; this relay does not drop partitions. +-- A redacted (soft-deleted) head payload may still be purged. +CREATE FUNCTION retain_current_artifact() RETURNS TRIGGER LANGUAGE plpgsql AS $$ +BEGIN + IF OLD.kind = 45010 AND OLD.deleted_at IS NULL AND EXISTS ( + SELECT 1 FROM artifact_heads h + WHERE h.community_id=OLD.community_id AND h.event_id=OLD.id + ) THEN + RETURN NULL; + END IF; + RETURN OLD; +END; +$$; +CREATE TRIGGER retain_current_artifact BEFORE DELETE ON events +FOR EACH ROW EXECUTE FUNCTION retain_current_artifact(); diff --git a/scripts/reconcile-schema-after-pgschema.sql b/scripts/reconcile-schema-after-pgschema.sql index 7c3a7be871a..def1558f8c5 100644 --- a/scripts/reconcile-schema-after-pgschema.sql +++ b/scripts/reconcile-schema-after-pgschema.sql @@ -19,6 +19,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p_past; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p_past; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p_past; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p_past; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p_past; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p_past; ALTER TABLE events ATTACH PARTITION events_p_past @@ -33,6 +34,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_01; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_01; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_01; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_01; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_01; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_01; ALTER TABLE events ATTACH PARTITION events_p2026_01 @@ -47,6 +49,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_02; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_02; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_02; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_02; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_02; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_02; ALTER TABLE events ATTACH PARTITION events_p2026_02 @@ -61,6 +64,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_03; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_03; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_03; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_03; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_03; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_03; ALTER TABLE events ATTACH PARTITION events_p2026_03 @@ -75,6 +79,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_04; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_04; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_04; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_04; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_04; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_04; ALTER TABLE events ATTACH PARTITION events_p2026_04 @@ -89,6 +94,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_05; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_05; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_05; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_05; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_05; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_05; ALTER TABLE events ATTACH PARTITION events_p2026_05 @@ -103,6 +109,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_06; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_06; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_06; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_06; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_06; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_06; ALTER TABLE events ATTACH PARTITION events_p2026_06 @@ -117,6 +124,7 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p_future; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p_future; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p_future; + DROP TRIGGER IF EXISTS retain_current_artifact ON events_p_future; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p_future; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p_future; ALTER TABLE events ATTACH PARTITION events_p_future From 4955e04aa2c7f6fbffcf76b50666fe703fb3a94a Mon Sep 17 00:00:00 2001 From: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> Date: Sun, 27 Sep 2026 14:56:23 -0700 Subject: [PATCH 3/6] fix(relay): close NIP-AR review gaps - Move markers tag the replaced revision (`prev`) instead of a relay-wide `position`, so they stay unique across repeated same-second moves without revealing activity outside the source. The unused `artifact_revisions.position` column is dropped. - Empty predicate arrays fail with 400 instead of generating invalid SQL. - Filters without `kinds` may match artifacts, so multi-character predicates on them are rejected on HTTP and WS rather than silently dropped. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> --- crates/buzz-core/src/artifact.rs | 38 +++++++++++++----- crates/buzz-db/src/store/artifact.rs | 20 +++++++--- .../src/store/artifact_postgres_tests.rs | 40 +++++++++++++++++++ .../src/api/artifact_postgres_tests.rs | 30 ++++++++++++++ crates/buzz-relay/src/protocol.rs | 22 ++++++---- docs/nips/NIP-AR.md | 2 +- migrations/0052_channel_artifacts.sql | 3 +- schema/schema.sql | 3 +- 8 files changed, 128 insertions(+), 30 deletions(-) diff --git a/crates/buzz-core/src/artifact.rs b/crates/buzz-core/src/artifact.rs index b4b8a27638a..74f33c08d8a 100644 --- a/crates/buzz-core/src/artifact.rs +++ b/crates/buzz-core/src/artifact.rs @@ -228,10 +228,20 @@ mod tests { route_filter(&json!({"kinds":[45010,45011],"#h":["c"]})), FilterRoute::Generic ); - assert!(matches!( - route_filter(&json!({"kinds":[45010],"#project":["p"]})), - FilterRoute::Rejected(_) - )); + for rejected in [ + json!({"kinds":[45010],"#project":["p"]}), + json!({"ids":["a"],"#project":["p"]}), + json!({"kinds":"x","#project":["p"]}), + ] { + assert!( + matches!(route_filter(&rejected), FilterRoute::Rejected(_)), + "{rejected}" + ); + } + assert_eq!( + route_filter(&json!({"ids":["a"],"#h":["c"]})), + FilterRoute::Generic + ); assert_eq!( route_filter(&json!({"artifact":"current"})), FilterRoute::Artifact @@ -244,6 +254,8 @@ mod tests { json!({"artifact":"current","kinds":[45011]}), json!({"artifact":"current","search":"x"}), json!({"artifact":"current","limit":1001}), + json!({"artifact":"current","#assignee":[]}), + json!({"artifact":"history","#d":[]}), ] { assert!(parse_query(&invalid).is_err(), "{invalid}"); } @@ -291,15 +303,16 @@ pub enum FilterRoute { Rejected(&'static str), } /// Route a raw filter. Generic filters drop multi-character tag predicates, -/// so artifact-kind filters carrying them must use an explicit artifact query. +/// so filters that may match artifacts (including those without `kinds`) +/// carrying them must use an explicit artifact query. pub fn route_filter(value: &serde_json::Value) -> FilterRoute { if value.get("artifact").is_some() { return FilterRoute::Artifact; } - let artifact_kind = value - .get("kinds") - .and_then(|k| k.as_array()) - .is_some_and(|ks| ks.iter().any(|k| matches!(k.as_u64(), Some(45010 | 45011)))); + let artifact_kind = value.get("kinds").is_none_or(|k| { + k.as_array() + .is_none_or(|ks| ks.iter().any(|k| matches!(k.as_u64(), Some(45010 | 45011)))) + }); let multi_character = value .as_object() .is_some_and(|o| o.keys().any(|k| k.starts_with('#') && k.len() > 2)); @@ -333,7 +346,10 @@ pub fn parse_query(value: &serde_json::Value) -> Result MAX_TAG_NAME_BYTES { return Err("invalid predicate name size"); } - let values = value.as_array().ok_or("predicate values must be arrays")?; + let values = value + .as_array() + .filter(|v| !v.is_empty()) + .ok_or("predicate values must be non-empty arrays")?; let mut parsed = Vec::new(); for value in values { let value = value.as_str().ok_or("predicate values must be strings")?; @@ -364,7 +380,7 @@ pub fn parse_query(value: &serde_json::Value) -> Result MAX_PAGE_SIZE as u64 || offset > MAX_OFFSET { return Err("artifact page limit exceeded"); } - if view == ArtifactView::History && !tags.iter().any(|(n, v)| n == "d" && !v.is_empty()) { + if view == ArtifactView::History && !tags.iter().any(|(n, _)| n == "d") { return Err("history requires #d"); } Ok(ArtifactQuery { diff --git a/crates/buzz-db/src/store/artifact.rs b/crates/buzz-db/src/store/artifact.rs index c3e75847610..653abe665e5 100644 --- a/crates/buzz-db/src/store/artifact.rs +++ b/crates/buzz-db/src/store/artifact.rs @@ -29,13 +29,15 @@ fn invalid(message: &str) -> DbError { DbError::InvalidData(message.into()) } -fn removal_marker(keys: &Keys, artifact: Uuid, source: Uuid, position: i64) -> Result { +/// The replaced revision (`prev`) was readable in the source, so it +/// distinguishes repeated moves without revealing other activity. +fn removal_marker(keys: &Keys, artifact: Uuid, source: Uuid, prev: &[u8]) -> Result { let tags = [ ["ar", "1"].map(str::to_owned), ["d".into(), artifact.to_string()], ["h".into(), source.to_string()], ["reason".into(), "moved".into()], - ["position".into(), position.to_string()], + ["prev".into(), hex::encode(prev)], ] .into_iter() .map(Tag::parse) @@ -169,16 +171,22 @@ impl Db { } } } - let position: i64 = sqlx::query_scalar("INSERT INTO artifact_revisions (community_id,event_id,artifact_id) VALUES ($1,$2,$3) RETURNING position") - .bind(community.as_uuid()).bind(event.id.as_bytes().as_slice()).bind(env.id).fetch_one(&mut *tx).await?; + sqlx::query( + "INSERT INTO artifact_revisions (community_id,event_id,artifact_id) VALUES ($1,$2,$3)", + ) + .bind(community.as_uuid()) + .bind(event.id.as_bytes().as_slice()) + .bind(env.id) + .execute(&mut *tx) + .await?; sqlx::query("INSERT INTO artifact_heads (community_id,artifact_id,event_id,channel_id,artifact_type,root,deleted) VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (community_id,artifact_id) DO UPDATE SET event_id=EXCLUDED.event_id,channel_id=EXCLUDED.channel_id,root=EXCLUDED.root,deleted=EXCLUDED.deleted") .bind(community.as_uuid()).bind(env.id).bind(event.id.as_bytes().as_slice()).bind(env.home).bind(&env.artifact_type).bind(&env.root).bind(env.op == ArtifactOp::Delete).execute(&mut *tx).await?; let (stored, _) = crate::event::insert_event_in_transaction(&mut tx, community, event, Some(env.home)) .await?; let mut accepted = vec![stored]; - if let Some(source) = source.filter(|_| env.op == ArtifactOp::Move) { - let removal = removal_marker(relay_keys, env.id, source, position)?; + if let (ArtifactOp::Move, Some(source), Some(prev)) = (env.op, source, &env.prev) { + let removal = removal_marker(relay_keys, env.id, source, prev)?; let (stored, _) = crate::event::insert_event_in_transaction( &mut tx, community, diff --git a/crates/buzz-db/src/store/artifact_postgres_tests.rs b/crates/buzz-db/src/store/artifact_postgres_tests.rs index 97655ac1e48..8fa06fd97dc 100644 --- a/crates/buzz-db/src/store/artifact_postgres_tests.rs +++ b/crates/buzz-db/src/store/artifact_postgres_tests.rs @@ -248,6 +248,46 @@ async fn lifecycle_cas_queries_move_redaction_and_retention() { f.db.validate_deletion_serving_catalog().await.unwrap(); } +#[tokio::test] +#[ignore = "requires Postgres"] +async fn repeated_moves_get_distinct_source_safe_markers() { + let f = Fixture::new().await; + let d = Uuid::new_v4(); + let create = f.revision(d, "create", f.a, None, &f.owner, vec![]); + let there = f.revision(d, "move", f.b, Some(&create), &f.owner, vec![]); + let back = f.revision(d, "move", f.a, Some(&there), &f.owner, vec![]); + let again = f.revision(d, "move", f.b, Some(&back), &f.owner, vec![]); + assert!(matches!( + f.accept(&create, None).await, + ArtifactOutcome::Accepted(_) + )); + let mut markers = Vec::new(); + for (event, source, replaced) in [ + (&there, f.a, &create), + (&back, f.b, &there), + (&again, f.a, &back), + ] { + let ArtifactOutcome::Accepted(stored) = f.accept(event, Some(source)).await else { + panic!("move accepted"); + }; + let marker = stored[1].event.clone(); + // Only `prev` distinguishes same-second markers for one source. + let tags: Vec<_> = marker.tags.iter().map(|t| t.as_slice().to_vec()).collect(); + assert_eq!( + tags, + [ + vec!["ar".to_string(), "1".into()], + vec!["d".into(), d.to_string()], + vec!["h".into(), source.to_string()], + vec!["reason".into(), "moved".into()], + vec!["prev".into(), replaced.id.to_hex()], + ] + ); + markers.push(marker.id); + } + assert_ne!(markers[0], markers[2]); +} + #[tokio::test] #[ignore = "requires Postgres"] async fn root_anchor_and_tenant_boundaries() { diff --git a/crates/buzz-relay/src/api/artifact_postgres_tests.rs b/crates/buzz-relay/src/api/artifact_postgres_tests.rs index fda2e84ada5..f2c14820e79 100644 --- a/crates/buzz-relay/src/api/artifact_postgres_tests.rs +++ b/crates/buzz-relay/src/api/artifact_postgres_tests.rs @@ -155,6 +155,31 @@ async fn project_filter_survives_artifact_query_hook() { assert!(status.is_success(), "{status}: {body}"); } +#[tokio::test] +#[ignore = "requires Postgres"] +async fn unsupported_predicates_fail_explicitly() { + let f = Fixture::new().await; + let create = f.revision(Uuid::new_v4(), "create", None); + f.publish(&create).await; + let (status, body) = f.request("/query", json!([{"ids":[create.id]}])).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body[0]["id"], create.id.to_hex()); + // A readable artifact must not match a predicate it does not carry. + for path in ["/query", "/count"] { + for filter in [ + json!({"ids":[create.id],"#project":["P"]}), + json!({"artifact":"current","#assignee":[]}), + ] { + let (status, body) = f.request(path, json!([filter])).await; + assert_eq!( + status, + axum::http::StatusCode::BAD_REQUEST, + "{filter}: {body}" + ); + } + } +} + #[tokio::test] #[ignore = "requires Postgres"] async fn lifecycle_routes_duplicate_replay_delete_and_redaction() { @@ -271,6 +296,11 @@ async fn writes_use_channel_gates_for_home_and_move_source() { .await; assert_eq!(removals.as_array().unwrap().len(), 1, "{removals}"); assert!(!removals.to_string().contains(&f.private.to_string())); + // The marker names only the source-readable revision it replaced. + assert!(removals[0]["tags"] + .as_array() + .unwrap() + .contains(&json!(["prev", peer_edit.id.to_hex()]))); // The peer can write the destination but not the private source. let back = f.revision_as(&f.peer, f.home, d, "move", Some(&moved)); let (status, body) = f.try_publish(&back).await; diff --git a/crates/buzz-relay/src/protocol.rs b/crates/buzz-relay/src/protocol.rs index 294e4f6f51f..8994a105c7c 100644 --- a/crates/buzz-relay/src/protocol.rs +++ b/crates/buzz-relay/src/protocol.rs @@ -285,16 +285,22 @@ mod tests { #[test] fn artifact_live_filters_preserved_extensions_rejected() { - assert!(super::ClientMessage::parse( - r##"["REQ","ar",{"kinds":[45010,45011],"#h":["channel"]}]"## - ) - .is_ok()); + let id = "f".repeat(64); for request in [ - r##"["REQ","ar",{"artifact":"history","#d":["artifact"]}]"##, - r##"["REQ","ar",{"kinds":[45010],"#project":["project"]}]"##, - r##"["COUNT","ar",{"artifact":"current"}]"##, + r##"["REQ","ar",{"kinds":[45010,45011],"#h":["channel"]}]"##.to_owned(), + format!(r##"["REQ","ar",{{"ids":["{id}"],"#h":["channel"]}}]"##), ] { - assert!(super::ClientMessage::parse(request).is_err(), "{request}"); + assert!(super::ClientMessage::parse(&request).is_ok(), "{request}"); + } + for request in [ + r##"["REQ","ar",{"artifact":"history","#d":["artifact"]}]"##.to_owned(), + r##"["REQ","ar",{"kinds":[45010],"#project":["project"]}]"##.to_owned(), + r##"["COUNT","ar",{"artifact":"current"}]"##.to_owned(), + // No `kinds` may match artifacts, so the predicate cannot be dropped. + format!(r##"["REQ","ar",{{"ids":["{id}"],"#project":["P"]}}]"##), + format!(r##"["COUNT","ar",{{"ids":["{id}"],"#project":["P"]}}]"##), + ] { + assert!(super::ClientMessage::parse(&request).is_err(), "{request}"); } } diff --git a/docs/nips/NIP-AR.md b/docs/nips/NIP-AR.md index 18e5e0d3a46..3f0a8874996 100644 --- a/docs/nips/NIP-AR.md +++ b/docs/nips/NIP-AR.md @@ -81,7 +81,7 @@ Relationship tags organize work without changing access. Linking a task to a pro A move publishes the current snapshot into the destination under the same `d`. Its `root` must be absent or belong to the destination. Clients MUST show the destination audience and the information being shared before confirmation. -The relay MUST atomically advance the current revision and store both arrival in the destination and removal from the source. The source receives a separate relay-authenticated removal identifying the artifact, source scope, and replay position, without destination metadata or content. Both are stored events, so a client that misses live delivery recovers them by replaying the channel. Relays unable to provide this MUST reject moves. +The relay MUST atomically advance the current revision and store both arrival in the destination and removal from the source. The source receives a separate relay-authenticated removal identifying the artifact, source scope, and the revision it replaced, without destination metadata or content. Both are stored events, so a client that misses live delivery recovers them by replaying the channel. Relays unable to provide this MUST reject moves. Earlier revisions remain under their original channels' access rules; a move never lets the destination read revisions from the source. The destination can load current state without reading earlier revisions. Conversation messages stay where they are. diff --git a/migrations/0052_channel_artifacts.sql b/migrations/0052_channel_artifacts.sql index 32380529b5c..6606c7d7af6 100644 --- a/migrations/0052_channel_artifacts.sql +++ b/migrations/0052_channel_artifacts.sql @@ -11,12 +11,11 @@ CREATE TABLE artifact_heads ( ); CREATE INDEX artifact_heads_event ON artifact_heads (community_id, event_id); -- Every accepted revision ID, so replays stay idempotent after redaction or --- retention. `position` records acceptance order. +-- retention. CREATE TABLE artifact_revisions ( community_id UUID NOT NULL REFERENCES communities(id), event_id BYTEA NOT NULL CHECK (length(event_id) = 32), artifact_id UUID NOT NULL, - position BIGINT GENERATED ALWAYS AS IDENTITY, PRIMARY KEY (community_id, event_id) ); diff --git a/schema/schema.sql b/schema/schema.sql index af72f819271..3ea3a354f4f 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -2012,12 +2012,11 @@ CREATE TABLE artifact_heads ( ); CREATE INDEX artifact_heads_event ON artifact_heads (community_id, event_id); -- Every accepted revision ID, so replays stay idempotent after redaction or --- retention. `position` records acceptance order. +-- retention. CREATE TABLE artifact_revisions ( community_id UUID NOT NULL REFERENCES communities(id), event_id BYTEA NOT NULL CHECK (length(event_id) = 32), artifact_id UUID NOT NULL, - position BIGINT GENERATED ALWAYS AS IDENTITY, PRIMARY KEY (community_id, event_id) ); From 9ff1400b67effb86955efe51c94e9918f67e6e8b Mon Sep 17 00:00:00 2001 From: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> Date: Sun, 27 Sep 2026 16:18:03 -0700 Subject: [PATCH 4/6] fix(relay): index artifact mentions and match #d before LIMIT - Artifact writes insert `p`-tag mentions in the same transaction as the revision, so generic `#p` replay and COUNT find what live subscribers already saw. - Filters over kind 45010 alone push `#d` into SQL via JSONB containment. Artifacts are not NIP-33, so `d_tag` is NULL and the Rust post-filter ran after LIMIT, dropping identity lookups. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> --- crates/buzz-db/src/store/artifact.rs | 1 + crates/buzz-db/src/store/event.rs | 33 +++++++++++ .../src/api/artifact_postgres_tests.rs | 57 +++++++++++++++++++ crates/buzz-relay/src/handlers/req.rs | 32 +++++++++++ 4 files changed, 123 insertions(+) diff --git a/crates/buzz-db/src/store/artifact.rs b/crates/buzz-db/src/store/artifact.rs index 653abe665e5..a4dbf429398 100644 --- a/crates/buzz-db/src/store/artifact.rs +++ b/crates/buzz-db/src/store/artifact.rs @@ -184,6 +184,7 @@ impl Db { let (stored, _) = crate::event::insert_event_in_transaction(&mut tx, community, event, Some(env.home)) .await?; + crate::insert_mentions_in_transaction(&mut tx, community, event, Some(env.home)).await?; let mut accepted = vec![stored]; if let (ArtifactOp::Move, Some(source), Some(prev)) = (env.op, source, &env.prev) { let removal = removal_marker(relay_keys, env.id, source, prev)?; diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index 3ec2cbf89dc..f137d2c5747 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -78,6 +78,10 @@ pub struct EventQuery { /// Restrict results to events with an `e` tag referencing any of these event IDs (hex). /// Uses JSONB containment (`tags @> ...`) against the `tags` column. pub e_tags: Option>, + /// Restrict results to events with a `d` tag matching any of these values, + /// via JSONB containment. For non-NIP-33 kinds whose `d_tag` column is NULL + /// (NIP-AR artifacts), so identity lookups match before SQL `LIMIT`. + pub d_tag_values: Option>, /// Restrict results to events with an exact custom tag pair. /// Uses JSONB containment against `tags` before SQL `LIMIT`. pub custom_tag: Option<(String, String)>, @@ -139,6 +143,7 @@ impl EventQuery { authors: None, ids: None, e_tags: None, + d_tag_values: None, custom_tag: None, channel_ids: None, channel_ids_include_global: true, @@ -658,6 +663,20 @@ fn build_query_events_sql(q: &EventQuery) -> QueryBuilder { } } + if let Some(ref values) = q.d_tag_values { + if !values.is_empty() { + qb.push(" AND ("); + for (i, value) in values.iter().enumerate() { + if i > 0 { + qb.push(" OR "); + } + qb.push(format!("{col_prefix}tags @> ")); + qb.push_bind(serde_json::json!([["d", value]])); + } + qb.push(")"); + } + } + if let Some((ref name, ref value)) = q.custom_tag { let containment = serde_json::json!([[name, value]]); qb.push(format!(" AND {col_prefix}tags @> ")) @@ -961,6 +980,20 @@ pub(crate) async fn count_events_on(conn: &mut sqlx::PgConnection, q: &EventQuer } } + if let Some(ref values) = q.d_tag_values { + if !values.is_empty() { + qb.push(" AND ("); + for (i, value) in values.iter().enumerate() { + if i > 0 { + qb.push(" OR "); + } + qb.push(format!("{col_prefix}tags @> ")); + qb.push_bind(serde_json::json!([["d", value]])); + } + qb.push(")"); + } + } + if let Some(s) = q.since { qb.push(format!(" AND {col_prefix}created_at >= ")) .push_bind(s); diff --git a/crates/buzz-relay/src/api/artifact_postgres_tests.rs b/crates/buzz-relay/src/api/artifact_postgres_tests.rs index f2c14820e79..640261f6090 100644 --- a/crates/buzz-relay/src/api/artifact_postgres_tests.rs +++ b/crates/buzz-relay/src/api/artifact_postgres_tests.rs @@ -317,3 +317,60 @@ async fn writes_use_channel_gates_for_home_and_move_source() { assert!(!status.is_success() || body["accepted"] == false, "{body}"); assert!(body.to_string().contains("archived"), "{body}"); } + +/// Re-sign an artifact revision with extra tags and an explicit timestamp. +fn resign(event: &Event, key: &Keys, extra: &[[&str; 2]], created_at: u64) -> Event { + EventBuilder::new(event.kind, event.content.clone()) + .tags( + event + .tags + .iter() + .cloned() + .chain(extra.iter().map(|t| Tag::parse(*t).unwrap())), + ) + .custom_created_at(nostr::Timestamp::from(created_at)) + .sign_with_keys(key) + .unwrap() +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn generic_p_reads_find_mentioned_artifacts() { + let f = Fixture::new().await; + let tagged = f.peer.public_key().to_hex(); + let create = f.revision(Uuid::new_v4(), "create", None); + let create = resign( + &create, + &f.owner, + &[["p", &tagged]], + create.created_at.as_secs(), + ); + f.publish(&create).await; + let filter = json!([{"kinds":[45010],"#h":[f.home],"#p":[tagged]}]); + let (status, body) = f.request("/query", filter.clone()).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body.as_array().unwrap().len(), 1, "{body}"); + assert_eq!(body[0]["id"], create.id.to_hex()); + let (status, body) = f.request("/count", filter).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body["count"], 1, "{body}"); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn generic_d_lookup_matches_before_limit() { + let f = Fixture::new().await; + let (a, b) = (Uuid::new_v4(), Uuid::new_v4()); + let now = nostr::Timestamp::now().as_secs(); + let older = resign(&f.revision(a, "create", None), &f.owner, &[], now - 60); + f.publish(&older).await; + f.publish(&f.revision(b, "create", None)).await; + let filter = json!([{"kinds":[45010],"#h":[f.home],"#d":[a],"limit":1}]); + let (status, body) = f.request("/query", filter.clone()).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body.as_array().unwrap().len(), 1, "{body}"); + assert_eq!(body[0]["id"], older.id.to_hex()); + let (status, body) = f.request("/count", filter).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body["count"], 1, "{body}"); +} diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index a428f1e821b..2ef7e85a614 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -1046,6 +1046,19 @@ fn filter_to_query_params( } else { (None, None) }; + // NIP-AR artifacts carry their stable identity in `d` but are not NIP-33, + // so `d_tag` stays NULL; match the tag itself before `LIMIT`. + let filter_is_artifact_only = kinds.as_ref().is_some_and(|ks| { + !ks.is_empty() + && ks + .iter() + .all(|&k| k as u32 == buzz_core::kind::KIND_ARTIFACT) + }); + let d_tag_values = filter_is_artifact_only + .then(|| filter.generic_tags.get(&d_tag_key)) + .flatten() + .filter(|values| !values.is_empty()) + .map(|values| values.iter().map(|v| v.to_string()).collect()); EventQuery { channel_id, @@ -1060,6 +1073,7 @@ fn filter_to_query_params( authors, ids, e_tags, + d_tag_values, ..EventQuery::for_community(community) } } @@ -2303,6 +2317,24 @@ mod tests { .expect("max_limit") as i64 } + #[test] + fn artifact_d_filter_is_pushed_before_limit() { + let community = buzz_core::tenant::CommunityId::from_uuid(uuid::Uuid::new_v4()); + let d = nostr::SingleLetterTag::lowercase(nostr::Alphabet::D); + let artifact = Filter::new() + .kind(nostr::Kind::Custom(45010)) + .custom_tag(d, "a") + .limit(1); + let q = filter_to_query_params(&artifact, None, community); + assert_eq!(q.d_tag_values, Some(vec!["a".to_string()])); + assert_eq!(q.d_tag, None); + + // Mixed kinds keep the generic post-filter path. + let mixed = artifact.clone().kind(nostr::Kind::Custom(9)); + let q = filter_to_query_params(&mixed, None, community); + assert_eq!(q.d_tag_values, None); + } + #[test] fn req_filter_limit_clamps_to_advertised_nip11_max_limit() { let advertised = advertised_max_limit(); From 1e95e6b2899ac1650171a54715c07cc9f5340a85 Mon Sep 17 00:00:00 2001 From: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> Date: Mon, 28 Sep 2026 10:00:21 -0700 Subject: [PATCH 5/6] fix(relay): match artifact #d on mixed and kindless filters - Push `#d` into SQL whenever a filter can select NIP-AR rows (no `kinds`, or kinds including 45010/45011), as `kind NOT IN (45010,45011) OR tags @> [["d", v]]`. A newer artifact no longer takes the `limit` slot from the requested one; other kinds keep the generic post-filter. - Drop the `retain_current_artifact` BEFORE DELETE trigger on events. The relay has no event retention; the community purge clears `artifact_heads` first. Future retention must skip head payloads. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> --- .../src/store/artifact_postgres_tests.rs | 13 ++--- crates/buzz-db/src/store/event.rs | 58 +++++++++++-------- .../src/api/artifact_postgres_tests.rs | 56 +++++++++++++++--- crates/buzz-relay/src/handlers/req.rs | 30 ++++++---- migrations/0052_channel_artifacts.sql | 19 +----- schema/schema.sql | 19 +----- scripts/reconcile-schema-after-pgschema.sql | 8 --- 7 files changed, 112 insertions(+), 91 deletions(-) diff --git a/crates/buzz-db/src/store/artifact_postgres_tests.rs b/crates/buzz-db/src/store/artifact_postgres_tests.rs index 8fa06fd97dc..d41fd0cb804 100644 --- a/crates/buzz-db/src/store/artifact_postgres_tests.rs +++ b/crates/buzz-db/src/store/artifact_postgres_tests.rs @@ -175,15 +175,14 @@ async fn lifecycle_cas_queries_move_redaction_and_retention() { assert_eq!(f.count(&f.peer, history.clone()).await, 3); assert_eq!(f.count(&f.owner, history.clone()).await, 4); - // Retention removes old payloads but reserves identity and the current payload. - let retention = || { - sqlx::query("DELETE FROM events WHERE community_id=$1 AND id IN ($2,$3)") + // Expiring an earlier payload keeps its identity reserved. + let retention = |id: nostr::EventId| { + sqlx::query("DELETE FROM events WHERE community_id=$1 AND id=$2") .bind(f.community.as_uuid()) - .bind(create.id.as_bytes().as_slice()) - .bind(moved.id.as_bytes().as_slice()) + .bind(id.as_bytes().to_vec()) }; assert_eq!( - retention() + retention(create.id) .execute(&f.db.pool) .await .unwrap() @@ -205,7 +204,7 @@ async fn lifecycle_cas_queries_move_redaction_and_retention() { assert_eq!(f.count(&f.owner, all.clone()).await, 0); assert_eq!(f.count(&f.owner, history.clone()).await, 2); assert_eq!( - retention() + retention(moved.id) .execute(&f.db.pool) .await .unwrap() diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index f137d2c5747..462a9f0e0c6 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -32,6 +32,12 @@ pub use crate::reminder::{ /// the advertised ceiling and the enforced one cannot drift. pub const DEFAULT_MAX_PAGE_LIMIT: i64 = 1_000; +/// NIP-AR revision and removal kinds, whose stable identity is a `d` tag. +pub const ARTIFACT_KINDS: [i32; 2] = [ + buzz_core::kind::KIND_ARTIFACT as i32, + buzz_core::kind::KIND_ARTIFACT_REMOVAL as i32, +]; + /// Optional filters for [`query_events`]. #[derive(Debug, Clone)] pub struct EventQuery { @@ -78,9 +84,10 @@ pub struct EventQuery { /// Restrict results to events with an `e` tag referencing any of these event IDs (hex). /// Uses JSONB containment (`tags @> ...`) against the `tags` column. pub e_tags: Option>, - /// Restrict results to events with a `d` tag matching any of these values, - /// via JSONB containment. For non-NIP-33 kinds whose `d_tag` column is NULL - /// (NIP-AR artifacts), so identity lookups match before SQL `LIMIT`. + /// Restrict artifact rows ([`ARTIFACT_KINDS`]) to those with a `d` tag + /// matching any of these values, via JSONB containment. Their `d_tag` + /// column is NULL (not NIP-33), so this lets identity lookups match before + /// SQL `LIMIT`. Rows of other kinds are left to the caller's post-filter. pub d_tag_values: Option>, /// Restrict results to events with an exact custom tag pair. /// Uses JSONB containment against `tags` before SQL `LIMIT`. @@ -664,17 +671,7 @@ fn build_query_events_sql(q: &EventQuery) -> QueryBuilder { } if let Some(ref values) = q.d_tag_values { - if !values.is_empty() { - qb.push(" AND ("); - for (i, value) in values.iter().enumerate() { - if i > 0 { - qb.push(" OR "); - } - qb.push(format!("{col_prefix}tags @> ")); - qb.push_bind(serde_json::json!([["d", value]])); - } - qb.push(")"); - } + push_artifact_d_tag_predicate(&mut qb, col_prefix, values); } if let Some((ref name, ref value)) = q.custom_tag { @@ -807,6 +804,27 @@ fn push_e_tag_filter(qb: &mut QueryBuilder, col_prefix: &str, e_ .push("::jsonb[])"); } +/// Match `#d` on artifact rows before `LIMIT` while leaving other kinds to the +/// caller's post-filter: `(kind NOT IN (artifact kinds) OR tags @> [["d", v]] ...)`. +fn push_artifact_d_tag_predicate( + qb: &mut QueryBuilder, + col_prefix: &str, + values: &[String], +) { + if values.is_empty() { + return; + } + let [revision, removal] = ARTIFACT_KINDS; + qb.push(format!( + " AND ({col_prefix}kind NOT IN ({revision}, {removal})" + )); + for value in values { + qb.push(format!(" OR {col_prefix}tags @> ")); + qb.push_bind(serde_json::json!([["d", value]])); + } + qb.push(")"); +} + pub(crate) fn row_to_stored_event(row: sqlx::postgres::PgRow) -> Result> { let id_bytes: Vec = row.try_get("id")?; let pubkey_bytes: Vec = row.try_get("pubkey")?; @@ -981,17 +999,7 @@ pub(crate) async fn count_events_on(conn: &mut sqlx::PgConnection, q: &EventQuer } if let Some(ref values) = q.d_tag_values { - if !values.is_empty() { - qb.push(" AND ("); - for (i, value) in values.iter().enumerate() { - if i > 0 { - qb.push(" OR "); - } - qb.push(format!("{col_prefix}tags @> ")); - qb.push_bind(serde_json::json!([["d", value]])); - } - qb.push(")"); - } + push_artifact_d_tag_predicate(&mut qb, col_prefix, values); } if let Some(s) = q.since { diff --git a/crates/buzz-relay/src/api/artifact_postgres_tests.rs b/crates/buzz-relay/src/api/artifact_postgres_tests.rs index 640261f6090..a78dca7eea1 100644 --- a/crates/buzz-relay/src/api/artifact_postgres_tests.rs +++ b/crates/buzz-relay/src/api/artifact_postgres_tests.rs @@ -362,15 +362,53 @@ async fn generic_d_lookup_matches_before_limit() { let f = Fixture::new().await; let (a, b) = (Uuid::new_v4(), Uuid::new_v4()); let now = nostr::Timestamp::now().as_secs(); - let older = resign(&f.revision(a, "create", None), &f.owner, &[], now - 60); + // Kindless reads are p-gated to the reader, so every row tags the peer. + let reader = f.peer.public_key().to_hex(); + let tag_reader = [["p", reader.as_str()]]; + let older = resign( + &f.revision(a, "create", None), + &f.owner, + &tag_reader, + now - 60, + ); f.publish(&older).await; - f.publish(&f.revision(b, "create", None)).await; - let filter = json!([{"kinds":[45010],"#h":[f.home],"#d":[a],"limit":1}]); - let (status, body) = f.request("/query", filter.clone()).await; - assert!(status.is_success(), "{status}: {body}"); - assert_eq!(body.as_array().unwrap().len(), 1, "{body}"); - assert_eq!(body[0]["id"], older.id.to_hex()); - let (status, body) = f.request("/count", filter).await; + let newer = resign(&f.revision(b, "create", None), &f.owner, &tag_reader, now); + f.publish(&newer).await; + for filter in [ + json!([{"kinds":[45010],"#h":[f.home],"#d":[a],"limit":1}]), + json!([{"kinds":[45010,45011],"#h":[f.home],"#d":[a],"limit":1}]), + json!([{"#h":[f.home],"#p":[reader],"#d":[a],"limit":1}]), + ] { + let (status, body) = f.request_as(&f.peer, "/query", filter.clone()).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body.as_array().unwrap().len(), 1, "{filter}: {body}"); + assert_eq!(body[0]["id"], older.id.to_hex(), "{filter}"); + let (status, body) = f.request_as(&f.peer, "/count", filter.clone()).await; + assert!(status.is_success(), "{status}: {body}"); + assert_eq!(body["count"], 1, "{filter}: {body}"); + } + + // Non-artifact rows still reach the generic `#d` post-filter. + let message = EventBuilder::new(Kind::Custom(9), "d-tagged chat") + .tags([ + Tag::parse(["h", &f.home.to_string()]).unwrap(), + Tag::parse(["d", &a.to_string()]).unwrap(), + Tag::parse(["p", &reader]).unwrap(), + ]) + .sign_with_keys(&f.owner) + .unwrap(); + f.publish(&message).await; + let filter = json!([{"#h":[f.home],"#p":[reader],"#d":[a]}]); + let (status, body) = f.request_as(&f.peer, "/query", filter).await; assert!(status.is_success(), "{status}: {body}"); - assert_eq!(body["count"], 1, "{body}"); + let mut ids: Vec<_> = body + .as_array() + .unwrap() + .iter() + .map(|e| e["id"].as_str().unwrap().to_owned()) + .collect(); + ids.sort(); + let mut expected = vec![older.id.to_hex(), message.id.to_hex()]; + expected.sort(); + assert_eq!(ids, expected, "{body}"); } diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index 2ef7e85a614..44efd5f309b 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -1046,15 +1046,14 @@ fn filter_to_query_params( } else { (None, None) }; - // NIP-AR artifacts carry their stable identity in `d` but are not NIP-33, - // so `d_tag` stays NULL; match the tag itself before `LIMIT`. - let filter_is_artifact_only = kinds.as_ref().is_some_and(|ks| { - !ks.is_empty() - && ks - .iter() - .all(|&k| k as u32 == buzz_core::kind::KIND_ARTIFACT) + // NIP-AR revisions and removals carry their stable identity in `d` but are + // not NIP-33, so `d_tag` stays NULL; match the tag on artifact rows before + // `LIMIT` whenever the filter can select them. + let filter_can_match_artifacts = kinds.as_ref().is_none_or(|ks| { + ks.iter() + .any(|k| buzz_db::event::ARTIFACT_KINDS.contains(k)) }); - let d_tag_values = filter_is_artifact_only + let d_tag_values = filter_can_match_artifacts .then(|| filter.generic_tags.get(&d_tag_key)) .flatten() .filter(|values| !values.is_empty()) @@ -2329,9 +2328,20 @@ mod tests { assert_eq!(q.d_tag_values, Some(vec!["a".to_string()])); assert_eq!(q.d_tag, None); - // Mixed kinds keep the generic post-filter path. - let mixed = artifact.clone().kind(nostr::Kind::Custom(9)); + // Mixed and kindless filters that can select artifact rows push it too. + let mixed = artifact.clone().kind(nostr::Kind::Custom(45011)); let q = filter_to_query_params(&mixed, None, community); + assert_eq!(q.d_tag_values, Some(vec!["a".to_string()])); + let kindless = Filter::new().custom_tag(d, "a").limit(1); + let q = filter_to_query_params(&kindless, None, community); + assert_eq!(q.d_tag_values, Some(vec!["a".to_string()])); + + // Filters that cannot select artifacts keep the generic post-filter path. + let other = Filter::new() + .kind(nostr::Kind::Custom(9)) + .custom_tag(d, "a") + .limit(1); + let q = filter_to_query_params(&other, None, community); assert_eq!(q.d_tag_values, None); } diff --git a/migrations/0052_channel_artifacts.sql b/migrations/0052_channel_artifacts.sql index 6606c7d7af6..11331139aa1 100644 --- a/migrations/0052_channel_artifacts.sql +++ b/migrations/0052_channel_artifacts.sql @@ -22,19 +22,6 @@ CREATE TABLE artifact_revisions ( SELECT attach_community_write_fence('artifact_heads'); SELECT attach_community_write_fence('artifact_revisions'); --- Row retention may expire earlier payloads, never the live CAS head. Partition --- retirement must also preserve head payloads; this relay does not drop partitions. --- A redacted (soft-deleted) head payload may still be purged. -CREATE FUNCTION retain_current_artifact() RETURNS TRIGGER LANGUAGE plpgsql AS $$ -BEGIN - IF OLD.kind = 45010 AND OLD.deleted_at IS NULL AND EXISTS ( - SELECT 1 FROM artifact_heads h - WHERE h.community_id=OLD.community_id AND h.event_id=OLD.id - ) THEN - RETURN NULL; - END IF; - RETURN OLD; -END; -$$; -CREATE TRIGGER retain_current_artifact BEFORE DELETE ON events -FOR EACH ROW EXECUTE FUNCTION retain_current_artifact(); +-- The relay does not expire events. Any future row retention or partition +-- retirement must skip payloads referenced by `artifact_heads.event_id` +-- (NIP-AR: expiring earlier revisions MUST NOT remove the current revision). diff --git a/schema/schema.sql b/schema/schema.sql index 3ea3a354f4f..c8f294c40ef 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -2023,19 +2023,6 @@ CREATE TABLE artifact_revisions ( SELECT attach_community_write_fence('artifact_heads'); SELECT attach_community_write_fence('artifact_revisions'); --- Row retention may expire earlier payloads, never the live CAS head. Partition --- retirement must also preserve head payloads; this relay does not drop partitions. --- A redacted (soft-deleted) head payload may still be purged. -CREATE FUNCTION retain_current_artifact() RETURNS TRIGGER LANGUAGE plpgsql AS $$ -BEGIN - IF OLD.kind = 45010 AND OLD.deleted_at IS NULL AND EXISTS ( - SELECT 1 FROM artifact_heads h - WHERE h.community_id=OLD.community_id AND h.event_id=OLD.id - ) THEN - RETURN NULL; - END IF; - RETURN OLD; -END; -$$; -CREATE TRIGGER retain_current_artifact BEFORE DELETE ON events -FOR EACH ROW EXECUTE FUNCTION retain_current_artifact(); +-- The relay does not expire events. Any future row retention or partition +-- retirement must skip payloads referenced by `artifact_heads.event_id` +-- (NIP-AR: expiring earlier revisions MUST NOT remove the current revision). diff --git a/scripts/reconcile-schema-after-pgschema.sql b/scripts/reconcile-schema-after-pgschema.sql index def1558f8c5..7c3a7be871a 100644 --- a/scripts/reconcile-schema-after-pgschema.sql +++ b/scripts/reconcile-schema-after-pgschema.sql @@ -19,7 +19,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p_past; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p_past; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p_past; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p_past; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p_past; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p_past; ALTER TABLE events ATTACH PARTITION events_p_past @@ -34,7 +33,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_01; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_01; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_01; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_01; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_01; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_01; ALTER TABLE events ATTACH PARTITION events_p2026_01 @@ -49,7 +47,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_02; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_02; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_02; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_02; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_02; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_02; ALTER TABLE events ATTACH PARTITION events_p2026_02 @@ -64,7 +61,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_03; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_03; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_03; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_03; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_03; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_03; ALTER TABLE events ATTACH PARTITION events_p2026_03 @@ -79,7 +75,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_04; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_04; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_04; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_04; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_04; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_04; ALTER TABLE events ATTACH PARTITION events_p2026_04 @@ -94,7 +89,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_05; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_05; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_05; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_05; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_05; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_05; ALTER TABLE events ATTACH PARTITION events_p2026_05 @@ -109,7 +103,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p2026_06; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p2026_06; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p2026_06; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p2026_06; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p2026_06; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p2026_06; ALTER TABLE events ATTACH PARTITION events_p2026_06 @@ -124,7 +117,6 @@ BEGIN DROP TRIGGER IF EXISTS events_enqueue_push_match ON events_p_future; DROP TRIGGER IF EXISTS events_refresh_channel_ttl ON events_p_future; DROP TRIGGER IF EXISTS events_created_at_floor ON events_p_future; - DROP TRIGGER IF EXISTS retain_current_artifact ON events_p_future; DROP TRIGGER IF EXISTS community_write_fence_events ON events_p_future; DROP TRIGGER IF EXISTS trg_events_guard_channel_roster_snapshot ON events_p_future; ALTER TABLE events ATTACH PARTITION events_p_future From 5a7d2db7f4561857465db8d8d8e235a3c10df45d Mon Sep 17 00:00:00 2001 From: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> Date: Mon, 28 Sep 2026 14:27:45 -0700 Subject: [PATCH 6/6] fix(db): express artifact-head filter without OR The event query shape test from #7854 asserts the e-tag query contains no ` OR `. Rewrite `kind <> 45010 OR EXISTS (...)` as the equivalent `NOT (kind = 45010 AND NOT EXISTS (...))`. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: murderbot <3754f8729004d95654c46dbab3129e4ab9ef05cc2534e2a3fbfc155983bd637b@buzz.block.builderlab.xyz> --- crates/buzz-db/src/store/event.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index 462a9f0e0c6..168c6d16131 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -585,7 +585,7 @@ fn build_query_events_sql(q: &EventQuery) -> QueryBuilder { // tombstones; explicit revision IDs also read earlier revisions. if q.ids.is_none() { let table = if q.p_tag_hex.is_some() { "e" } else { "events" }; - qb.push(format!(" AND ({table}.kind <> 45010 OR EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); + qb.push(format!(" AND NOT ({table}.kind = 45010 AND NOT EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); } if let Some(ch) = q.channel_id { @@ -923,7 +923,7 @@ pub(crate) async fn count_events_on(conn: &mut sqlx::PgConnection, q: &EventQuer // tombstones; explicit revision IDs also read earlier revisions. if q.ids.is_none() { let table = if q.p_tag_hex.is_some() { "e" } else { "events" }; - qb.push(format!(" AND ({table}.kind <> 45010 OR EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); + qb.push(format!(" AND NOT ({table}.kind = 45010 AND NOT EXISTS (SELECT 1 FROM artifact_heads ah WHERE ah.community_id={table}.community_id AND ah.event_id={table}.id))")); } if let Some(ch) = q.channel_id {