diff --git a/crates/buzz-db/src/admin_moderation.rs b/crates/buzz-db/src/admin_moderation.rs index 31efaca36..c255a8eb4 100644 --- a/crates/buzz-db/src/admin_moderation.rs +++ b/crates/buzz-db/src/admin_moderation.rs @@ -68,6 +68,61 @@ pub struct AdminReportedMessage { pub deleted_at: Option>, } +/// The `relay_admin_actions` enforcement record governing a report. +/// +/// Populated on the report detail read and enforcement resolve response. Carries +/// the durable state machine so the console can render enforcement progress or +/// terminal outcome without inventing a shape the relay never emits. +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AdminActionDto { + /// Action row identifier. + pub id: Uuid, + /// Client-generated idempotency key. + pub request_id: Uuid, + /// Principal who claimed the report. + pub actor_pubkey: String, + /// Role of the actor: `"operator"` | `"moderator"`. + pub actor_role: String, + /// Enforcement action name: `"delete"` | `"kick"` | `"ban"` | `"timeout"`. + pub action: String, + /// State machine: `"pending"` | `"enforcing"` | `"succeeded"` | `"failed"` | `"cancelled"`. + pub status: String, + /// Operator reason, if provided. + pub reason: Option, + /// Absolute timeout expiry for `timeout` actions; null otherwise. Absolute + /// (not remaining-seconds) so repeated reads never disagree; the client + /// computes remaining time. + pub expires_at: Option>, + /// Error from the last failure, if any. + pub error_message: Option, + /// Action creation time. + pub created_at: DateTime, + /// Action last-updated time. + pub updated_at: DateTime, +} + +impl AdminActionDto { + /// Build the wire DTO from a persistence record. Used to embed the + /// just-cancelled action in the cancel response — the last look at a record + /// that a subsequent detail read (report back to `open`) no longer surfaces. + pub fn from_record(record: &crate::relay_admin_actions::AdminActionRecord) -> Self { + Self { + id: record.id, + request_id: record.request_id, + actor_pubkey: hex::encode(&record.actor_pubkey), + actor_role: record.actor_role.clone(), + action: record.action.clone(), + status: record.state.clone(), + reason: record.reason.clone(), + expires_at: record.timeout_until, + error_message: record.error_message.clone(), + created_at: record.created_at, + updated_at: record.updated_at, + } + } +} + /// Deployment-global moderation report detail. #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase")] @@ -77,6 +132,9 @@ pub struct AdminReportDetail { pub report: AdminReport, /// Reported message when the report targets a stored event. pub message: Option, + /// Governing enforcement action, when one exists (live or terminal). Null + /// for reports never enforced via the HTTP admin plane. + pub active_action: Option, } /// Deployment-global product feedback with source-community provenance. @@ -85,10 +143,13 @@ pub struct AdminReportDetail { pub struct AdminFeedback { /// Feedback row identifier. pub id: Uuid, - /// Source community identifier. - pub community_id: Uuid, - /// Source community host. - pub community_host: String, + /// Source community identifier. `None` once the source community has been + /// purged: `product_feedback` is deployment-global operator evidence whose + /// `community_id` is severed to NULL on tenant purge, not cascade-deleted. + pub community_id: Option, + /// Source community host. `None` when `community_id` is severed (no row to + /// join) — the feedback is retained without its origin tenant. + pub community_host: Option, /// Signed feedback event identifier. pub event_id: String, /// Submitter public key. @@ -99,6 +160,8 @@ pub struct AdminFeedback { pub body: String, /// Full source tags, including attachment metadata. pub tags: serde_json::Value, + /// Operator-managed lifecycle status: `"new"` | `"reviewed"` | `"archived"`. + pub status: String, /// Timestamp signed into the feedback event. pub event_created_at: DateTime, /// Time accepted by this deployment. @@ -165,7 +228,18 @@ pub async fn get_report(pool: &PgPool, report_id: Uuid) -> Result Result Result Result> { + let Some(id) = row.try_get::, _>("action_id_admin")? else { + return Ok(None); + }; + Ok(Some(AdminActionDto { + id, + request_id: row.try_get("action_request_id")?, + actor_pubkey: hex::encode(row.try_get::, _>("action_actor_pubkey")?), + actor_role: row.try_get("action_actor_role")?, + action: row.try_get("action_name")?, + status: row.try_get("action_state")?, + reason: row.try_get("action_reason")?, + expires_at: row.try_get("action_timeout_until")?, + error_message: row.try_get("action_error_message")?, + created_at: row.try_get("action_created_at")?, + updated_at: row.try_get("action_updated_at")?, + })) +} + fn row_to_report(row: sqlx::postgres::PgRow) -> Result { let target_kind: String = row.try_get("target_kind")?; let target = match target_kind.as_str() { @@ -237,10 +347,10 @@ pub async fn list_feedback(pool: &PgPool, limit: i64) -> Result Result Result { category: row.try_get("category")?, body: row.try_get("body")?, tags: row.try_get("tags")?, + status: row.try_get("status")?, event_created_at: row.try_get("event_created_at")?, received_at: row.try_get("received_at")?, }) @@ -384,6 +495,14 @@ mod tests { } async fn delete_report_fixture(pool: &PgPool, community_id: Uuid) { + // relay_admin_actions FK-references (community_id, report_id), so clear + // any enforcement/audit rows before the reports they point at. A no-op + // for tests that never insert actions. + sqlx::query("DELETE FROM relay_admin_actions WHERE report_community_id = $1") + .bind(community_id) + .execute(pool) + .await + .expect("delete admin action fixture"); sqlx::query("DELETE FROM moderation_reports WHERE community_id = $1") .bind(community_id) .execute(pool) @@ -485,4 +604,288 @@ mod tests { delete_report_fixture(&pool, community_id).await; } + + // ── activeAction LATERAL join ───────────────────────────────────────────── + + #[allow(clippy::too_many_arguments)] + async fn insert_admin_action( + pool: &PgPool, + id: Uuid, + community_id: Uuid, + report_id: Uuid, + action: &str, + state: &str, + created_at: DateTime, + ) { + sqlx::query( + r#" + INSERT INTO relay_admin_actions ( + id, report_id, report_community_id, request_id, actor_pubkey, + actor_role, action, state, created_at, updated_at + ) VALUES ($1, $2, $3, $4, $5, 'operator', $6, $7, $8, $8) + "#, + ) + .bind(id) + .bind(report_id) + .bind(community_id) + .bind(Uuid::new_v4()) + .bind(vec![2_u8; 32]) + .bind(action) + .bind(state) + .bind(created_at) + .execute(pool) + .await + .expect("insert admin action"); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn report_detail_surfaces_succeeded_enforcement_on_a_dismissed_reopened_report() { + // Enforcement succeeded, report was later reopened and re-triaged to + // `dismissed`. active_action_id is NULL, but the succeeded enforcement + // row still matches `a.state='succeeded'` — the activeAction must surface + // that DTO: a later dismissal does not un-happen the executed ban. + let pool = setup_pool().await; + let community_id = insert_community(&pool, "dismissed-after-enforce").await; + let report_id = insert_pubkey_report(&pool, community_id).await; + let action_id = Uuid::new_v4(); + insert_admin_action( + &pool, + action_id, + community_id, + report_id, + "ban", + "succeeded", + Utc::now(), + ) + .await; + set_report_status(&pool, community_id, report_id, "dismissed").await; + + let detail = get_report(&pool, report_id) + .await + .expect("query report") + .expect("report exists"); + assert_eq!(detail.report.status, "dismissed"); + let action = detail + .active_action + .expect("succeeded enforcement DTO survives dismissal"); + assert_eq!(action.id, action_id); + assert_eq!(action.action, "ban"); + assert_eq!(action.status, "succeeded"); + + delete_report_fixture(&pool, community_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn report_detail_active_action_breaks_equal_timestamp_ties_by_id_desc() { + // Two succeeded enforcement rows (possible across reopen cycles) sharing + // an identical created_at: the `a.id DESC` tiebreaker must pick the + // greater id deterministically, never leave the choice to row order. + let pool = setup_pool().await; + let community_id = insert_community(&pool, "equal-ts-tiebreak").await; + let report_id = insert_pubkey_report(&pool, community_id).await; + let ts = Utc::now(); + let id_a = Uuid::new_v4(); + let id_b = Uuid::new_v4(); + insert_admin_action(&pool, id_a, community_id, report_id, "ban", "succeeded", ts).await; + insert_admin_action( + &pool, + id_b, + community_id, + report_id, + "kick", + "succeeded", + ts, + ) + .await; + set_report_status(&pool, community_id, report_id, "resolved").await; + + let detail = get_report(&pool, report_id) + .await + .expect("query report") + .expect("report exists"); + let action = detail.active_action.expect("an action surfaces"); + assert_eq!( + action.id, + id_a.max(id_b), + "equal timestamps must resolve to the greater id via a.id DESC" + ); + + delete_report_fixture(&pool, community_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn report_detail_active_action_excludes_reopen_audit_rows() { + // A `reopen` audit row is written state='succeeded'. The enforcement DTO + // join filters `action IN (delete,kick,ban,timeout)`, so a report whose + // only relay_admin_actions row is a reopen audit must surface no action. + let pool = setup_pool().await; + let community_id = insert_community(&pool, "reopen-audit-excluded").await; + let report_id = insert_pubkey_report(&pool, community_id).await; + insert_admin_action( + &pool, + Uuid::new_v4(), + community_id, + report_id, + "reopen", + "succeeded", + Utc::now(), + ) + .await; + + let detail = get_report(&pool, report_id) + .await + .expect("query report") + .expect("report exists"); + assert!( + detail.active_action.is_none(), + "reopen audit row must not surface as an enforcement action" + ); + + delete_report_fixture(&pool, community_id).await; + } + + async fn set_report_status(pool: &PgPool, community_id: Uuid, report_id: Uuid, status: &str) { + sqlx::query( + "UPDATE moderation_reports SET status = $3 WHERE community_id = $1 AND id = $2", + ) + .bind(community_id) + .bind(report_id) + .bind(status) + .execute(pool) + .await + .expect("set report status"); + } + + // ── Feedback severed-provenance survival ────────────────────────────────── + + async fn insert_feedback(pool: &PgPool, community_id: Uuid, status: &str) -> Uuid { + let id = Uuid::new_v4(); + let event_id: Vec = id + .as_bytes() + .iter() + .chain(id.as_bytes().iter()) + .copied() + .collect(); + sqlx::query( + r#" + INSERT INTO product_feedback ( + id, community_id, event_id, submitter_pubkey, category, body, + tags, status, event_created_at, received_at + ) VALUES ($1, $2, $3, $4, 'bug', 'reproduces on launch', '[]'::jsonb, + $5, now(), now()) + "#, + ) + .bind(id) + .bind(community_id) + .bind(event_id) + .bind(vec![9_u8; 32]) + .bind(status) + .execute(pool) + .await + .expect("insert feedback"); + id + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn feedback_survives_community_purge_with_null_provenance_in_list_and_detail() { + // purge_postgres severs tenant provenance without deleting the row: + // `UPDATE product_feedback SET community_id = NULL`. The LEFT JOIN must + // keep the row visible in both list and detail reads with null + // community fields and its operator-managed status intact. + let pool = setup_pool().await; + let community_id = insert_community(&pool, "severed-feedback").await; + let feedback_id = insert_feedback(&pool, community_id, "reviewed").await; + + // Sever provenance exactly as the community purge transaction does. + sqlx::query("UPDATE product_feedback SET community_id = NULL WHERE community_id = $1") + .bind(community_id) + .execute(&pool) + .await + .expect("sever provenance"); + + let detail = get_feedback(&pool, feedback_id) + .await + .expect("query feedback") + .expect("severed feedback still readable in detail"); + assert_eq!(detail.id, feedback_id); + assert!( + detail.community_id.is_none(), + "community_id severed to None" + ); + assert!( + detail.community_host.is_none(), + "community_host has no row to join" + ); + assert_eq!(detail.status, "reviewed", "operator status is retained"); + + let listed = list_feedback(&pool, MAX_PAGE_SIZE) + .await + .expect("list feedback"); + let row = listed + .iter() + .find(|f| f.id == feedback_id) + .expect("severed feedback still appears in the list read"); + assert!(row.community_id.is_none()); + assert!(row.community_host.is_none()); + assert_eq!(row.status, "reviewed"); + + sqlx::query("DELETE FROM product_feedback WHERE id = $1") + .bind(feedback_id) + .execute(&pool) + .await + .expect("delete feedback fixture"); + sqlx::query("DELETE FROM communities WHERE id = $1") + .bind(community_id) + .execute(&pool) + .await + .expect("delete community fixture"); + } + + // ── Wire contract: nullable action fields are required-nullable ─────────── + + #[test] + fn action_dto_emits_nullable_fields_as_json_null_not_absent() { + // The desktop console types must be required-nullable, not optional: + // `reason`, `expiresAt`, `errorMessage` are plain `Option` with no + // `skip_serializing_if`, so serde always emits the key (null when None). + // This test pins that contract so the seam can't silently drift. + let dto = AdminActionDto { + id: Uuid::nil(), + request_id: Uuid::nil(), + actor_pubkey: hex::encode([0_u8; 32]), + actor_role: "operator".to_string(), + action: "ban".to_string(), + status: "succeeded".to_string(), + reason: None, + expires_at: None, + error_message: None, + created_at: DateTime::::from_timestamp(0, 0).unwrap(), + updated_at: DateTime::::from_timestamp(0, 0).unwrap(), + }; + let value = serde_json::to_value(&dto).expect("serialize dto"); + let obj = value.as_object().expect("dto serializes to an object"); + for key in ["reason", "expiresAt", "errorMessage"] { + assert_eq!( + obj.get(key), + Some(&serde_json::Value::Null), + "{key} must be present and null, never absent" + ); + } + // Field names are camelCase on the wire. + for key in [ + "requestId", + "actorPubkey", + "actorRole", + "expiresAt", + "errorMessage", + "createdAt", + "updatedAt", + ] { + assert!(obj.contains_key(key), "missing camelCase key {key}"); + } + } } diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index c41393d52..dd70c97aa 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -4724,6 +4724,29 @@ impl Db { relay_admin_actions::cancel_action(&self.pool, action_id, community_id, report_id).await } + /// Reopen a terminal report (resolved|dismissed|escalated → open) with a + /// durable `reopen` audit row, keyed idempotent on `request_id`. + pub async fn reopen_report( + &self, + community_id: CommunityId, + report_id: uuid::Uuid, + request_id: uuid::Uuid, + actor_pubkey: &[u8], + actor_role: &str, + reason: Option<&str>, + ) -> Result { + relay_admin_actions::reopen_report( + &self.pool, + community_id, + report_id, + request_id, + actor_pubkey, + actor_role, + reason, + ) + .await + } + /// Fetch an action record by ID. pub async fn get_admin_action( &self, diff --git a/crates/buzz-db/src/relay_admin_actions.rs b/crates/buzz-db/src/relay_admin_actions.rs index 9f1364ea5..673308e57 100644 --- a/crates/buzz-db/src/relay_admin_actions.rs +++ b/crates/buzz-db/src/relay_admin_actions.rs @@ -1009,6 +1009,129 @@ pub async fn cancel_action( Ok(true) } +/// Result of attempting to reopen a terminal report. +#[derive(Debug)] +pub enum ReopenResult { + /// Report was terminal and is now `open`; a `reopen` audit row was inserted. + Reopened, + /// This exact `request_id` already reopened the report — idempotent replay. + /// No state changed; the earlier reopen stands. + AlreadyReopened, + /// The report is not in a terminal state. Carries its current status. + NotReopenable(String), + /// The report was not found globally. + NotFound, +} + +/// Reopen a terminal report (`resolved | dismissed | escalated` → `open`) in a +/// single transaction, recording a durable `reopen` audit row. +/// +/// The audit row is written to `relay_admin_actions` with `action = 'reopen'` +/// and `state = 'succeeded'`: `succeeded` keeps the stranded-action recovery +/// worker (which claims `state IN ('pending','enforcing')`) from ever driving +/// it, and the `action` value keeps it out of the enforcement DTO join (which +/// filters `action IN ('delete','kick','ban','timeout')`). +/// +/// Idempotency is keyed on `request_id`: a replay after the report has been +/// reopened (and possibly re-resolved) returns [`ReopenResult::AlreadyReopened`] +/// without mutating, so a client network retry never re-reopens a +/// freshly-resolved report. +pub async fn reopen_report( + pool: &PgPool, + community_id: CommunityId, + report_id: Uuid, + request_id: Uuid, + actor_pubkey: &[u8], + actor_role: &str, + reason: Option<&str>, +) -> Result { + let mut tx = pool.begin().await?; + + // Lock the report row to serialize concurrent reopen/resolve on it. + let report_row = sqlx::query( + r#" + SELECT status + FROM moderation_reports + WHERE community_id = $1 AND id = $2 + FOR UPDATE + "#, + ) + .bind(community_id.as_uuid()) + .bind(report_id) + .fetch_optional(&mut *tx) + .await?; + + let Some(report_row) = report_row else { + return Ok(ReopenResult::NotFound); + }; + + // Idempotent replay: this request_id already reopened the report. Checked + // before the terminal-status gate so a retry after a re-resolve still + // returns success rather than a spurious NotReopenable. + let existing = sqlx::query_scalar::<_, Uuid>( + r#" + SELECT id FROM relay_admin_actions + WHERE report_community_id = $1 AND report_id = $2 + AND request_id = $3 AND action = 'reopen' + "#, + ) + .bind(community_id.as_uuid()) + .bind(report_id) + .bind(request_id) + .fetch_optional(&mut *tx) + .await?; + + if existing.is_some() { + tx.rollback().await?; + return Ok(ReopenResult::AlreadyReopened); + } + + let status: String = report_row.try_get("status")?; + if !matches!(status.as_str(), "resolved" | "dismissed" | "escalated") { + tx.rollback().await?; + return Ok(ReopenResult::NotReopenable(status)); + } + + // Return the report to the queue. Clear the resolution stamp so an open + // report never carries a stale resolver/timestamp; active_action_id is + // already NULL on a terminal report but clear it defensively. + sqlx::query( + r#" + UPDATE moderation_reports + SET status = 'open', resolved_by = NULL, resolved_at = NULL, + active_action_id = NULL + WHERE community_id = $1 AND id = $2 + AND status IN ('resolved', 'dismissed', 'escalated') + "#, + ) + .bind(community_id.as_uuid()) + .bind(report_id) + .execute(&mut *tx) + .await?; + + // Durable audit row. Inserted as 'succeeded' so the recovery worker never + // claims it; 'reopen' keeps it out of the enforcement DTO join. + sqlx::query( + r#" + INSERT INTO relay_admin_actions ( + report_id, report_community_id, request_id, actor_pubkey, actor_role, + action, reason, state + ) VALUES ($1, $2, $3, $4, $5, 'reopen', $6, 'succeeded') + "#, + ) + .bind(report_id) + .bind(community_id.as_uuid()) + .bind(request_id) + .bind(actor_pubkey) + .bind(actor_role) + .bind(reason) + .execute(&mut *tx) + .await?; + + tx.commit().await?; + Ok(ReopenResult::Reopened) +} + /// Fetch an action record by ID. pub async fn get_action(pool: &PgPool, action_id: Uuid) -> Result> { let row = sqlx::query( @@ -2290,4 +2413,267 @@ mod tests { "second kick must return AlreadyGone" ); } + + // ── Reopen: terminal → open CAS + durable audit row ─────────────────────── + + async fn set_report_status(pool: &PgPool, community_id: Uuid, report_id: Uuid, status: &str) { + sqlx::query( + "UPDATE moderation_reports SET status = $3 WHERE community_id = $1 AND id = $2", + ) + .bind(community_id) + .bind(report_id) + .bind(status) + .execute(pool) + .await + .expect("set report status"); + } + + async fn report_status(pool: &PgPool, community_id: Uuid, report_id: Uuid) -> String { + sqlx::query_scalar( + "SELECT status FROM moderation_reports WHERE community_id = $1 AND id = $2", + ) + .bind(community_id) + .bind(report_id) + .fetch_one(pool) + .await + .expect("read report status") + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_terminal_report_returns_open_and_records_succeeded_audit_row() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + let cid = CommunityId::from_uuid(community_id); + + for terminal in ["resolved", "dismissed", "escalated"] { + let report_id = make_report(&pool, community_id).await; + set_report_status(&pool, community_id, report_id, terminal).await; + + let request_id = Uuid::new_v4(); + let result = reopen_report( + &pool, + cid, + report_id, + request_id, + &actor(), + "operator", + Some("re-triage"), + ) + .await + .expect("reopen_report"); + assert!( + matches!(result, ReopenResult::Reopened), + "reopen of {terminal} must return Reopened, got {result:?}" + ); + assert_eq!( + report_status(&pool, community_id, report_id).await, + "open", + "report must be open after reopen from {terminal}" + ); + + // The audit row is state='succeeded' and action='reopen' so the + // recovery worker never claims it and the DTO join never surfaces it. + let row = sqlx::query( + "SELECT action, state FROM relay_admin_actions WHERE report_id = $1 AND request_id = $2", + ) + .bind(report_id) + .bind(request_id) + .fetch_one(&pool) + .await + .expect("reopen audit row exists"); + assert_eq!(row.try_get::("action").unwrap(), "reopen"); + assert_eq!(row.try_get::("state").unwrap(), "succeeded"); + } + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_open_report_is_rejected_as_not_reopenable() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + let report_id = make_report(&pool, community_id).await; // starts 'open' + + let result = reopen_report( + &pool, + CommunityId::from_uuid(community_id), + report_id, + Uuid::new_v4(), + &actor(), + "moderator", + None, + ) + .await + .expect("reopen_report"); + assert!( + matches!(result, ReopenResult::NotReopenable(ref s) if s == "open"), + "reopen of an open report must return NotReopenable(open), got {result:?}" + ); + + // No audit row written on a rejected reopen. + let count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM relay_admin_actions WHERE report_id = $1") + .bind(report_id) + .fetch_one(&pool) + .await + .expect("count"); + assert_eq!(count, 0, "rejected reopen must not write an audit row"); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_processing_report_is_rejected_as_not_reopenable() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + let report_id = make_report(&pool, community_id).await; + // A live enforcement claim moves the report to 'processing'. + let _ = do_claim(&pool, community_id, report_id, Uuid::new_v4()).await; + + let result = reopen_report( + &pool, + CommunityId::from_uuid(community_id), + report_id, + Uuid::new_v4(), + &actor(), + "operator", + None, + ) + .await + .expect("reopen_report"); + assert!( + matches!(result, ReopenResult::NotReopenable(ref s) if s == "processing"), + "reopen of a processing report must return NotReopenable(processing), got {result:?}" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_is_idempotent_on_request_id_even_after_reresolve() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + let cid = CommunityId::from_uuid(community_id); + let report_id = make_report(&pool, community_id).await; + set_report_status(&pool, community_id, report_id, "resolved").await; + + let request_id = Uuid::new_v4(); + let first = reopen_report( + &pool, + cid, + report_id, + request_id, + &actor(), + "operator", + None, + ) + .await + .expect("first reopen"); + assert!(matches!(first, ReopenResult::Reopened)); + + // The report gets re-resolved (a fresh terminal cycle) before the client's + // network retry of the SAME reopen request lands. + set_report_status(&pool, community_id, report_id, "resolved").await; + + let replay = reopen_report( + &pool, + cid, + report_id, + request_id, + &actor(), + "operator", + None, + ) + .await + .expect("reopen replay"); + assert!( + matches!(replay, ReopenResult::AlreadyReopened), + "same request_id replay must return AlreadyReopened, got {replay:?}" + ); + + // The replay must NOT have re-reopened the freshly re-resolved report. + assert_eq!( + report_status(&pool, community_id, report_id).await, + "resolved", + "idempotent replay must not re-reopen a re-resolved report" + ); + + // Exactly one reopen audit row exists for this request_id. + let count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM relay_admin_actions WHERE report_id = $1 AND request_id = $2 AND action = 'reopen'", + ) + .bind(report_id) + .bind(request_id) + .fetch_one(&pool) + .await + .expect("count"); + assert_eq!( + count, 1, + "idempotent replay must not write a second audit row" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_missing_report_returns_not_found() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + + let result = reopen_report( + &pool, + CommunityId::from_uuid(community_id), + Uuid::new_v4(), // no such report + Uuid::new_v4(), + &actor(), + "operator", + None, + ) + .await + .expect("reopen_report"); + assert!( + matches!(result, ReopenResult::NotFound), + "reopen of a missing report must return NotFound, got {result:?}" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn reopen_audit_row_is_never_claimed_by_recovery_worker() { + let pool = setup_pool().await; + let community_id = make_community(&pool).await; + let cid = CommunityId::from_uuid(community_id); + let report_id = make_report(&pool, community_id).await; + set_report_status(&pool, community_id, report_id, "dismissed").await; + + let request_id = Uuid::new_v4(); + reopen_report( + &pool, + cid, + report_id, + request_id, + &actor(), + "operator", + None, + ) + .await + .expect("reopen"); + + // The stranded-action recovery worker claims state IN ('pending','enforcing'). + // A 'succeeded' reopen row must never appear in its batch — otherwise the + // worker would try to drive an enforcement mutation for a reopen. + let lease_until = chrono::Utc::now() + chrono::Duration::seconds(30); + let batch = claim_stranded_action_batch(&pool, "recovery-worker", lease_until, 1000) + .await + .expect("claim batch"); + let reopen_action_id: Uuid = sqlx::query_scalar( + "SELECT id FROM relay_admin_actions WHERE report_id = $1 AND request_id = $2", + ) + .bind(report_id) + .bind(request_id) + .fetch_one(&pool) + .await + .expect("reopen action id"); + assert!( + !batch.iter().any(|c| c.record.id == reopen_action_id), + "reopen audit row (state=succeeded) must never be claimed by the recovery worker" + ); + } } diff --git a/crates/buzz-relay/src/api/admin/mod.rs b/crates/buzz-relay/src/api/admin/mod.rs index c9451ac73..70fdb9d90 100644 --- a/crates/buzz-relay/src/api/admin/mod.rs +++ b/crates/buzz-relay/src/api/admin/mod.rs @@ -44,6 +44,8 @@ pub fn router(state: Arc) -> Router { .route("/reports", get(reports)) .route("/reports/{id}", get(report_detail)) .route("/reports/{id}/resolve", axum::routing::post(resolve_report)) + .route("/reports/{id}/reopen", axum::routing::post(reopen_report)) + .route("/reports/{id}/cancel", axum::routing::post(cancel_report)) .route("/feedback", get(feedback)) .route("/feedback/{id}", get(feedback_detail)) .route("/feedback/{id}", patch(update_feedback_status)) @@ -256,11 +258,15 @@ async fn report_detail( #[serde(rename_all = "camelCase")] struct FeedbackSummary { id: Uuid, - community_id: Uuid, - community_host: String, + /// `None` once the source community has been purged (provenance severed). + community_id: Option, + /// `None` when `community_id` is severed — feedback retained without origin. + community_host: Option, submitter_pubkey: String, category: Option, body_summary: String, + /// Operator-managed lifecycle status: `"new"` | `"reviewed"` | `"archived"`. + status: String, received_at: DateTime, } @@ -292,6 +298,7 @@ async fn feedback( submitter_pubkey: item.submitter_pubkey, category: item.category, body_summary, + status: item.status, received_at: item.received_at, } }) @@ -346,20 +353,30 @@ async fn feedback_attachment( .admin_get_feedback(id) .await? .ok_or_else(ApiError::not_found)?; - if !feedback_references_hash(&feedback.tags, &feedback.community_host, &sha256) { + + // A severed feedback row (source community purged, community_id NULL) has no + // tenant to bind and no tenant-scoped media to serve — its attachment bytes + // were purged with the community. Fail closed to 404. + let (Some(community_host), Some(community_id)) = + (feedback.community_host.as_deref(), feedback.community_id) + else { + return Err(ApiError::not_found()); + }; + + if !feedback_references_hash(&feedback.tags, community_host, &sha256) { return Err(ApiError::not_found()); } // Resolve the tenant from server-owned feedback provenance, then assert the // resolved row still agrees with the feedback FK. Client input never names // a community, host, object key, extension, or upstream URL. - let tenant = crate::tenant::bind_community(&state.db, &feedback.community_host) + let tenant = crate::tenant::bind_community(&state.db, community_host) .await .map_err(|_| ApiError::not_found())?; - if tenant.community().as_uuid() != &feedback.community_id { + if tenant.community().as_uuid() != &community_id { tracing::warn!( feedback_id = %feedback.id, - feedback_community_id = %feedback.community_id, + feedback_community_id = %community_id, resolved_community_id = %tenant.community(), "admin feedback attachment tenant provenance mismatch" ); @@ -374,7 +391,7 @@ async fn feedback_attachment( })?; tracing::info!( feedback_id = %feedback.id, - community_id = %feedback.community_id, + community_id = %community_id, attachment_sha256 = %sha256, "admin feedback attachment read" ); @@ -522,7 +539,11 @@ async fn resolve_report( .status(200) .header(header::CONTENT_TYPE, "application/json") .body(axum::body::Body::from( - serde_json::json!({"status": terminal_status}).to_string(), + serde_json::json!({ + "status": terminal_status, + "activeAction": serde_json::Value::Null, + }) + .to_string(), )) .unwrap()) } @@ -535,7 +556,7 @@ async fn resolve_report( ) })?; - let result = resolve_report_with_enforcement( + resolve_report_with_enforcement( &state, &tenant, &report_detail, @@ -565,13 +586,22 @@ async fn resolve_report( } })?; + // Re-read the report so the resolve response carries the same + // `status` + `activeAction` shape a later GET /reports/{id} returns — + // single source of truth for the enforcement DTO. + let detail = state + .db + .admin_get_report(report_id) + .await? + .ok_or_else(ApiError::internal)?; + Ok(axum::http::Response::builder() .status(200) .header(header::CONTENT_TYPE, "application/json") .body(axum::body::Body::from( serde_json::json!({ - "status": "resolved", - "actionId": result.action_id.to_string(), + "status": detail.report.status, + "activeAction": detail.active_action, }) .to_string(), )) @@ -580,6 +610,180 @@ async fn resolve_report( } } +/// Request body for POST /reports/{id}/reopen. +#[derive(Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +struct ReopenReportBody { + /// Client-generated idempotency key. A retry with the same key returns the + /// same success without re-reopening a report that has since been re-resolved. + request_id: Uuid, + /// Optional operator reason, recorded on the reopen audit row. + reason: Option, +} + +/// POST /reports/{id}/reopen +/// +/// Requires nip98 auth. Both Operator and Moderator may act. +/// +/// Returns a terminal report (`resolved | dismissed | escalated`) to `open` and +/// records a durable `reopen` audit row. `409` if the report is not terminal. +async fn reopen_report( + State(state): State>, + uri: Uri, + headers: HeaderMap, + Path(report_id): Path, + body_bytes: Bytes, +) -> Result, ApiError> { + use buzz_db::relay_admin_actions::ReopenResult; + + let principal_opt = authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "POST", + Some(&body_bytes), + ) + .await?; + + let principal = require_mutation_principal(principal_opt)?; + + let body: ReopenReportBody = serde_json::from_slice(&body_bytes) + .map_err(|_| ApiError::bad_request("invalid_body", "invalid JSON body"))?; + + // Load report globally to derive tenant provenance. + let report_detail = state + .db + .admin_get_report(report_id) + .await? + .ok_or_else(ApiError::not_found)?; + + // Bind tenant from server-owned report provenance (never from client input). + let tenant = crate::tenant::bind_community(&state.db, &report_detail.report.community_host) + .await + .map_err(|_| ApiError::internal())?; + + let actor_pubkey: Vec = principal.pubkey.to_vec(); + let actor_role_str = match principal.role { + AdminRole::Operator => "operator", + AdminRole::Moderator => "moderator", + }; + + let result = state + .db + .reopen_report( + tenant.community(), + report_id, + body.request_id, + &actor_pubkey, + actor_role_str, + body.reason.as_deref(), + ) + .await?; + + match result { + // AlreadyReopened returns the same success as the original reopen: the + // request_id identifies the reopen outcome, not a fresh status read. + ReopenResult::Reopened | ReopenResult::AlreadyReopened => { + Ok(Json(serde_json::json!({"status": "open"}))) + } + ReopenResult::NotReopenable(status) => Err(ApiError::conflict(&format!( + "report is not reopenable (current status: {status})" + ))), + ReopenResult::NotFound => Err(ApiError::not_found()), + } +} + +/// Request body for POST /reports/{id}/cancel. +#[derive(Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +struct CancelReportBody { + /// The failed action to cancel — the `activeAction.id` the client observed. + /// Fences the cancel to exactly that action: a mismatch (already cancelled, + /// superseded by a newer claim, or past the mutation point) resolves to 409. + action_id: Uuid, +} + +/// POST /reports/{id}/cancel +/// +/// Requires nip98 auth. Both Operator and Moderator may act. +/// +/// Cancels a pre-mutation `failed` enforcement action, returning the report to +/// `open`. Cancel is the only recovery path for a failed action (no composed +/// client-side retry). `409` if the action is not cancellable — treat as +/// "refresh detail" (someone else likely cancelled or the action advanced). +/// +/// The response embeds the just-cancelled action DTO: this is the last look at +/// that record, since a subsequent detail read (report back to `open`) serves +/// `activeAction: null`. +async fn cancel_report( + State(state): State>, + uri: Uri, + headers: HeaderMap, + Path(report_id): Path, + body_bytes: Bytes, +) -> Result, ApiError> { + let principal_opt = authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "POST", + Some(&body_bytes), + ) + .await?; + + let _principal = require_mutation_principal(principal_opt)?; + + let body: CancelReportBody = serde_json::from_slice(&body_bytes) + .map_err(|_| ApiError::bad_request("invalid_body", "invalid JSON body"))?; + + // Load report globally to derive tenant provenance. + let report_detail = state + .db + .admin_get_report(report_id) + .await? + .ok_or_else(ApiError::not_found)?; + + // Bind tenant from server-owned report provenance (never from client input). + let tenant = crate::tenant::bind_community(&state.db, &report_detail.report.community_host) + .await + .map_err(|_| ApiError::internal())?; + + let cancelled = state + .db + .cancel_admin_action(body.action_id, tenant.community(), report_id) + .await?; + + if !cancelled { + return Err(ApiError::conflict( + "action is not cancellable (already cancelled, superseded, or past the mutation point)", + )); + } + + // Re-read the just-cancelled action for the last-look DTO. The report is now + // `open`, so a detail read serves activeAction: null — this response is the + // only place the cancelled record surfaces. + let record = state + .db + .get_admin_action(body.action_id) + .await? + .ok_or_else(ApiError::internal)?; + let dto = buzz_db::admin_moderation::AdminActionDto::from_record(&record); + + Ok(axum::http::Response::builder() + .status(200) + .header(header::CONTENT_TYPE, "application/json") + .body(axum::body::Body::from( + serde_json::json!({ + "status": "open", + "activeAction": dto, + }) + .to_string(), + )) + .unwrap()) +} + /// PATCH /feedback/{id} /// /// Update product_feedback status. Requires nip98 auth. @@ -2930,6 +3134,262 @@ mod tests { ); } + // ── Item 9: HTTP → DB wiring for the reopen and cancel routes ───────────── + // + // reopen and cancel touch only `state.db` (no enforcement stack, Redis, or + // media), so a full `router().oneshot()` drive with nip98 auth exercises the + // real HTTP → handler → tenant-bind → DB path and reads the durable evidence + // back. This is the seeded action→HTTP→DB matrix for the two new routes. + + /// Seed a community whose host is `admin.example` (the nip98 test host) plus + /// one report in the given status. Returns the report id. + async fn seed_admin_host_report(pool: &sqlx::PgPool, status: &str) -> Uuid { + // The nip98 test host must resolve to a community, so bind_community in + // the handler succeeds. `communities.host` is uniquely indexed on + // lower(host), so reuse an existing row rather than racing an insert. + let existing: Option = + sqlx::query_scalar("SELECT id FROM communities WHERE lower(host) = 'admin.example'") + .fetch_optional(pool) + .await + .expect("lookup admin.example community"); + let community_id = match existing { + Some(id) => id, + None => sqlx::query_scalar( + "INSERT INTO communities (id, host) VALUES (gen_random_uuid(), 'admin.example') RETURNING id", + ) + .fetch_one(pool) + .await + .expect("seed admin.example community"), + }; + + let uid = Uuid::new_v4(); + let event_id: Vec = uid + .as_bytes() + .iter() + .chain(uid.as_bytes().iter()) + .copied() + .collect(); + let report_id: Uuid = sqlx::query_scalar( + r#" + INSERT INTO moderation_reports ( + community_id, report_event_id, reporter_pubkey, target_kind, + target_pubkey, report_type, status + ) VALUES ($1, $2, $3, 'pubkey', $4, 'harassment', $5) + RETURNING id + "#, + ) + .bind(community_id) + .bind(event_id) + .bind(vec![0u8; 32]) + .bind(vec![1u8; 32]) + .bind(status) + .fetch_one(pool) + .await + .expect("seed report"); + report_id + } + + #[tokio::test] + #[ignore = "requires Postgres — reopen HTTP route drives the DB"] + async fn reopen_route_returns_report_to_open_and_writes_audit_row() { + let keys = nostr::Keys::generate(); + let state = nip98_state(vec![keys.public_key().to_hex()]).await; + let pool = sqlx::PgPool::connect( + &std::env::var("BUZZ_TEST_DATABASE_URL") + .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()), + ) + .await + .expect("connect to test DB"); + let report_id = seed_admin_host_report(&pool, "resolved").await; + + let request_id = Uuid::new_v4(); + let body = serde_json::json!({ "requestId": request_id }).to_string(); + let path = format!("/reports/{report_id}/reopen"); + let auth = make_nostr_auth_post(&keys, &path, body.as_bytes()); + let response = status_for( + state, + Request::builder() + .method("POST") + .uri(&path) + .header(header::HOST, "admin.example") + .header(header::AUTHORIZATION, auth) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body)) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let bytes = axum::body::to_bytes(response.into_body(), 4096) + .await + .expect("body"); + let json: serde_json::Value = serde_json::from_slice(&bytes).expect("json"); + assert_eq!(json["status"], "open"); + + // DB evidence: report is open and a succeeded reopen audit row exists. + let status: String = + sqlx::query_scalar("SELECT status FROM moderation_reports WHERE id = $1") + .bind(report_id) + .fetch_one(&pool) + .await + .expect("read status"); + assert_eq!(status, "open"); + let (action, state_col): (String, String) = sqlx::query_as( + "SELECT action, state FROM relay_admin_actions WHERE report_id = $1 AND request_id = $2", + ) + .bind(report_id) + .bind(request_id) + .fetch_one(&pool) + .await + .expect("reopen audit row"); + assert_eq!( + (action.as_str(), state_col.as_str()), + ("reopen", "succeeded") + ); + + cleanup_admin_host_report(&pool, report_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres — reopen of an open report is 409"] + async fn reopen_route_rejects_non_terminal_report_with_409() { + let keys = nostr::Keys::generate(); + let state = nip98_state(vec![keys.public_key().to_hex()]).await; + let pool = sqlx::PgPool::connect( + &std::env::var("BUZZ_TEST_DATABASE_URL") + .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()), + ) + .await + .expect("connect to test DB"); + let report_id = seed_admin_host_report(&pool, "open").await; + + let body = serde_json::json!({ "requestId": Uuid::new_v4() }).to_string(); + let path = format!("/reports/{report_id}/reopen"); + let auth = make_nostr_auth_post(&keys, &path, body.as_bytes()); + let response = status_for( + state, + Request::builder() + .method("POST") + .uri(&path) + .header(header::HOST, "admin.example") + .header(header::AUTHORIZATION, auth) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body)) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::CONFLICT); + + cleanup_admin_host_report(&pool, report_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres — cancel HTTP route drives the DB"] + async fn cancel_route_returns_open_and_embeds_the_cancelled_action_dto() { + let keys = nostr::Keys::generate(); + let state = nip98_state(vec![keys.public_key().to_hex()]).await; + let pool = sqlx::PgPool::connect( + &std::env::var("BUZZ_TEST_DATABASE_URL") + .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()), + ) + .await + .expect("connect to test DB"); + let report_id = seed_admin_host_report(&pool, "open").await; + let community_id: Uuid = + sqlx::query_scalar("SELECT community_id FROM moderation_reports WHERE id = $1") + .bind(report_id) + .fetch_one(&pool) + .await + .expect("community id"); + let cid = buzz_core::CommunityId::from_uuid(community_id); + + // Claim → fail (pre-mutation) leaves a cancellable failed action. + let action_id = match buzz_db::relay_admin_actions::claim_report( + &pool, + cid, + report_id, + Uuid::new_v4(), + &[2u8; 32], + "operator", + "ban", + None, + None, + "resolve:ban", + "relay_operator", + Some(&[1u8; 32]), + None, + None, + ) + .await + .expect("claim") + { + buzz_db::relay_admin_actions::ClaimResult::Claimed(a) => a.id, + other => panic!("expected Claimed, got {other:?}"), + }; + // pending → enforcing → failed (pre-mutation): the only cancellable state. + buzz_db::relay_admin_actions::begin_enforcing(&pool, action_id) + .await + .expect("begin_enforcing"); + buzz_db::relay_admin_actions::record_failure(&pool, action_id, "boom") + .await + .expect("record_failure"); + + let body = serde_json::json!({ "actionId": action_id }).to_string(); + let path = format!("/reports/{report_id}/cancel"); + let auth = make_nostr_auth_post(&keys, &path, body.as_bytes()); + let response = status_for( + state, + Request::builder() + .method("POST") + .uri(&path) + .header(header::HOST, "admin.example") + .header(header::AUTHORIZATION, auth) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body)) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let bytes = axum::body::to_bytes(response.into_body(), 4096) + .await + .expect("body"); + let json: serde_json::Value = serde_json::from_slice(&bytes).expect("json"); + assert_eq!(json["status"], "open"); + // The last-look DTO embeds the just-cancelled action with status cancelled. + assert_eq!(json["activeAction"]["id"], action_id.to_string()); + assert_eq!(json["activeAction"]["status"], "cancelled"); + + // DB evidence: action is cancelled and the report is back to open. + let (state_col, report_status): (String, String) = sqlx::query_as( + r#" + SELECT a.state, r.status + FROM relay_admin_actions a + JOIN moderation_reports r ON r.id = a.report_id + WHERE a.id = $1 + "#, + ) + .bind(action_id) + .fetch_one(&pool) + .await + .expect("read action + report"); + assert_eq!(state_col, "cancelled"); + assert_eq!(report_status, "open"); + + cleanup_admin_host_report(&pool, report_id).await; + } + + async fn cleanup_admin_host_report(pool: &sqlx::PgPool, report_id: Uuid) { + sqlx::query("DELETE FROM relay_admin_actions WHERE report_id = $1") + .bind(report_id) + .execute(pool) + .await + .expect("delete actions"); + sqlx::query("DELETE FROM moderation_reports WHERE id = $1") + .bind(report_id) + .execute(pool) + .await + .expect("delete report"); + } + #[tokio::test] #[ignore = "requires Postgres — worker crash re-drive convergence"] async fn worker_crash_redrive_converges_to_exactly_one_enforcement() { @@ -3234,6 +3694,7 @@ mod tests { created_at: chrono::Utc::now(), }, message: None, + active_action: None, } } diff --git a/docs/admin/README.md b/docs/admin/README.md index 9d0773c6d..2879edd98 100644 --- a/docs/admin/README.md +++ b/docs/admin/README.md @@ -326,6 +326,24 @@ Feedback search and filters run over the bounded browser result set; the Body: `{"action": "delete|kick|ban|timeout|dismiss|escalate", "expirationSecs": , "reason": "", "requestId": ""}` `expirationSecs` required for `timeout`, rejected for all others. Target/channel are always derived from server-owned report provenance. + Response: `{"status": "", "activeAction": }`. Enforcement + actions return the governing action record; `dismiss`/`escalate` return `null`. +- `POST /api/admin/v1/reports/:id/reopen` + Body: `{"requestId": "", "reason": ""}` + Returns a terminal report (`resolved`, `dismissed`, or `escalated`) to `open` and + records a durable `reopen` audit row. Idempotent on `requestId`: a retry after the + report has been reopened (and possibly re-resolved) returns the same `200` without + re-reopening. Returns `{"status": "open"}` on success, `409` if the report is not + in a terminal state, `404` if it does not exist. +- `POST /api/admin/v1/reports/:id/cancel` + Body: `{"actionId": ""}` + Cancels a pre-mutation `failed` enforcement action, returning the report to `open`. + Cancel is the only recovery path for a failed action. `actionId` fences the cancel + to exactly the action the client observed. Returns `{"status": "open", "activeAction": + }` — the embedded record is the last look at that action, since a + subsequent detail read (report back to `open`) serves `activeAction: null`. Returns + `409` if the action is not cancellable (already cancelled, superseded, or past the + mutation point) — treat as "refresh detail". - `PATCH /api/admin/v1/feedback/:id` Body: `{"status": "new|reviewed|archived"}` @@ -358,6 +376,13 @@ sidecar before accessing the shared content-addressed blob. Unknown feedback, unreferenced hashes, malformed paths, and cross-community substitutions all collapse to `404`. +Product feedback is deployment-global operator evidence: when its source community +is purged, the row's `community_id` is severed to `NULL` (the row survives, its +provenance does not). List and detail reads still return such rows with +`communityId` and `communityHost` as `null`. Their attachments, however, were +purged with the tenant, so the attachment route fails closed to `404` for any +severed feedback — there is no tenant to bind and no tenant-scoped media to serve. + Only `GET` and `HEAD` are routed. Community `/media/*` reads always require Blossom authorization and relay membership; the browser receives no reusable signed URL. Responses are uncached, `nosniff`,