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..74f33c08d8a --- /dev/null +++ b/crates/buzz-core/src/artifact.rs @@ -0,0 +1,392 @@ +//! 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 + ); + 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 + ); + 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}), + json!({"artifact":"current","#assignee":[]}), + json!({"artifact":"history","#d":[]}), + ] { + 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 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").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)); + 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() + .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")?; + 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, _)| n == "d") { + 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..a4dbf429398 --- /dev/null +++ b/crates/buzz-db/src/store/artifact.rs @@ -0,0 +1,203 @@ +//! 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()) +} + +/// 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()], + ["prev".into(), hex::encode(prev)], + ] + .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", + )); + } + } + } + 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?; + 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)?; + 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..d41fd0cb804 --- /dev/null +++ b/crates/buzz-db/src/store/artifact_postgres_tests.rs @@ -0,0 +1,383 @@ +//! 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); + + // 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(id.as_bytes().to_vec()) + }; + assert_eq!( + retention(create.id) + .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(moved.id) + .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 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() { + 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 23082102e9f..168c6d16131 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,6 +84,11 @@ 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 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`. pub custom_tag: Option<(String, String)>, @@ -139,6 +150,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, @@ -569,6 +581,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 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 { qb.push(format!(" AND {col_prefix}channel_id = ")) @@ -652,6 +670,10 @@ fn build_query_events_sql(q: &EventQuery) -> QueryBuilder { } } + if let Some(ref values) = q.d_tag_values { + push_artifact_d_tag_predicate(&mut qb, col_prefix, values); + } + if let Some((ref name, ref value)) = q.custom_tag { let containment = serde_json::json!([[name, value]]); qb.push(format!(" AND {col_prefix}tags @> ")) @@ -782,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")?; @@ -876,6 +919,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 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 { qb.push(format!(" AND {col_prefix}channel_id = ")) @@ -949,6 +998,10 @@ pub(crate) async fn count_events_on(conn: &mut sqlx::PgConnection, q: &EventQuer } } + if let Some(ref values) = q.d_tag_values { + push_artifact_d_tag_predicate(&mut qb, col_prefix, values); + } + if let Some(s) = q.since { qb.push(format!(" AND {col_prefix}created_at >= ")) .push_bind(s); @@ -1109,6 +1162,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( @@ -1732,7 +1792,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( 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..a78dca7eea1 --- /dev/null +++ b/crates/buzz-relay/src/api/artifact_postgres_tests.rs @@ -0,0 +1,414 @@ +// 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 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() { + 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 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; + 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}"); +} + +/// 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(); + // 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; + 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}"); + 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/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/req.rs b/crates/buzz-relay/src/handlers/req.rs index a428f1e821b..44efd5f309b 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -1046,6 +1046,18 @@ fn filter_to_query_params( } else { (None, None) }; + // 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_can_match_artifacts + .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 +1072,7 @@ fn filter_to_query_params( authors, ids, e_tags, + d_tag_values, ..EventQuery::for_community(community) } } @@ -2303,6 +2316,35 @@ 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 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); + } + #[test] fn req_filter_limit_clamps_to_advertised_nip11_max_limit() { let advertised = advertised_max_limit(); 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..8994a105c7c 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,35 @@ 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() { + let id = "f".repeat(64); + for request in [ + r##"["REQ","ar",{"kinds":[45010,45011],"#h":["channel"]}]"##.to_owned(), + format!(r##"["REQ","ar",{{"ids":["{id}"],"#h":["channel"]}}]"##), + ] { + 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}"); + } + } + 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..3f0a8874996 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 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. @@ -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..11331139aa1 --- /dev/null +++ b/migrations/0052_channel_artifacts.sql @@ -0,0 +1,27 @@ +-- 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. +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, + PRIMARY KEY (community_id, event_id) +); + +SELECT attach_community_write_fence('artifact_heads'); +SELECT attach_community_write_fence('artifact_revisions'); + +-- 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 37112cc6702..c8f294c40ef 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1998,3 +1998,31 @@ 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. +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, + PRIMARY KEY (community_id, event_id) +); + +SELECT attach_community_write_fence('artifact_heads'); +SELECT attach_community_write_fence('artifact_revisions'); + +-- 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).