Skip to main content

cairn_mod/pds_admin/
dispatch.rs

1//! Post-recordAction PDS-admin dispatch (#87, v1.7).
2//!
3//! Bridges cairn-mod's recordAction pipeline (§F20-F22) to the
4//! configured [`PdsAdminBackend`] (v1.7 = `OzoneBackend` only).
5//! Called by the writer task in [`crate::writer`] after a
6//! `subject_actions` row has been committed and labels have been
7//! emitted; this module decides whether to dispatch a backend
8//! call, fires it, and audit-logs the result.
9//!
10//! # Failure semantics (per §A13)
11//!
12//! - The recordAction transaction has **already committed** by
13//!   the time this runs. Backend-call failures do **not** roll
14//!   back the cairn-mod-side action.
15//! - On backend-call failure, log loudly (`WARN` for transient
16//!   variants, `ERROR` for operator-actionable variants like
17//!   `Auth` / `Validation`), record the failure in
18//!   `pds_admin_audit`, and return.
19//! - On audit-insert failure, log at `ERROR` and return — do
20//!   NOT propagate to the writer task; the cairn-mod-side action
21//!   stays committed regardless. The audit divergence is a
22//!   discoverability issue, not a correctness issue at this
23//!   layer.
24//! - cairn-mod does NOT retry. Retry policy lands in v1.8 with
25//!   operator feedback informing it.
26//!
27//! # Method-selection logic
28//!
29//! 1. If [`PdsAdminPolicy::enabled`] is false → no-op.
30//! 2. Look up the action_type in
31//!    [`PdsAdminPolicy::action_map`]. Missing key → no-op
32//!    (defensive — #83's resolver guarantees coverage but a
33//!    future config-shape change shouldn't crash here).
34//! 3. [`ActionMapEntry::Skip`] → no-op.
35//! 4. [`ActionMapEntry::Method`] → dispatch:
36//!    - `TakedownAccount` (v1.7 implemented) → call the trait
37//!      method.
38//!    - `SuspendAccount` (lands in #88) → call the trait method;
39//!      `OzoneBackend`'s body currently `unimplemented!()`s,
40//!      which would panic. Acceptable in-cycle; #88 fills it in
41//!      before any tagged release.
42//!    - `RestoreAccount` → log error and skip. recordAction
43//!      doesn't carry a prior backend action id; restore is
44//!      reachable only from the revokeAction path (separate
45//!      flow, future cycle).
46//!    - `ApplyLabel` / `NegateLabel` → log warn and skip. #83's
47//!      action_map validation already warned at config-load;
48//!      this is defense-in-depth.
49
50use std::sync::Arc;
51
52use sqlx::{Pool, Sqlite};
53
54use crate::moderation::types::ActionType;
55use crate::pds_admin::audit::record_pds_admin_call;
56use crate::pds_admin::backend::{BackendActionId, BackendError, PdsAdminBackend};
57use crate::pds_admin::config::{ActionMapEntry, BackendMethod, PdsAdminPolicy};
58
59/// Bundled v1.7 PDS-admin runtime: the resolved
60/// [`PdsAdminPolicy`] from #83 + the trait-object backend the
61/// dispatch fires against.
62///
63/// `Option<PdsAdminBridge>` is what the writer task carries. When
64/// `None` (operator declared no `[pds_admin]` block, or
65/// `enabled = false`), the dispatch is short-circuited at the
66/// writer-task layer without even constructing this struct.
67#[derive(Clone)]
68pub struct PdsAdminBridge {
69    /// Resolved policy. The `enabled` flag is checked again on
70    /// every call (defensive — could become stale relative to a
71    /// future hot-reload feature).
72    pub policy: PdsAdminPolicy,
73    /// Trait-object backend. v1.7 wires `OzoneBackend`; v1.8
74    /// will add `LocusBackend` selectable per
75    /// [`PdsAdminPolicy::backend`].
76    pub backend: Arc<dyn PdsAdminBackend>,
77}
78
79/// Snapshot of the just-committed recordAction's fields the
80/// dispatch needs. Borrowed from the writer's
81/// `RecordActionRequest` to avoid cloning the whole struct
82/// across the post-commit boundary.
83#[derive(Debug, Clone, Copy)]
84pub struct DispatchContext<'a> {
85    /// `subject_actions(id)` of the just-committed row.
86    pub action_id: i64,
87    /// Cairn-mod-side action-type discriminator.
88    pub action_type: ActionType,
89    /// Subject DID (extracted from the request via
90    /// `route_subject` in writer.rs).
91    pub subject_did: &'a str,
92    /// Operator-vocabulary reason identifiers. The first entry
93    /// is propagated to the backend's `ref` / `reason` field;
94    /// remaining entries are recorded in the cairn-mod-side
95    /// audit log only.
96    pub reason_codes: &'a [String],
97    /// Optional moderator-facing free text. Currently NOT
98    /// mirrored to the backend (it's a cairn-mod-internal
99    /// artifact; bsky-PDS's `takedown.ref` field is for
100    /// cross-system tracking, not narrative content).
101    pub notes: Option<&'a str>,
102    /// ISO-8601 duration string from the recordAction request
103    /// (e.g. `P7D`, `PT24H`). Required for `temp_suspension`,
104    /// rejected for everything else by the recorder. The
105    /// dispatch parses this into `duration_days` for
106    /// `suspend_account`'s wire encoding (#89). `None` for
107    /// non-temp_suspension actions.
108    pub duration_iso: Option<&'a str>,
109}
110
111/// Dispatch the post-recordAction PDS-admin call (if any).
112///
113/// Idempotent on a `None` bridge: callers can pass `&None`
114/// without checking `enabled` themselves. Returns no error
115/// because failures are logged + audited, not propagated; the
116/// recordAction transaction has already committed.
117pub async fn dispatch_after_record_action(
118    bridge: Option<&PdsAdminBridge>,
119    pool: &Pool<Sqlite>,
120    ctx: DispatchContext<'_>,
121) {
122    let Some(bridge) = bridge else {
123        return;
124    };
125    if !bridge.policy.enabled {
126        return;
127    }
128    let entry = bridge
129        .policy
130        .action_map
131        .get(&ctx.action_type)
132        .copied()
133        .unwrap_or(ActionMapEntry::Skip);
134    let method = match entry {
135        ActionMapEntry::Skip => return,
136        ActionMapEntry::Method(m) => m,
137    };
138
139    // Defense-in-depth gating for methods that don't make sense
140    // from the recordAction path. #83's config-load validation
141    // catches the label cases with a warning; the runtime gate
142    // ensures we never actually touch those trait methods.
143    match method {
144        BackendMethod::ApplyLabel | BackendMethod::NegateLabel => {
145            tracing::warn!(
146                action_id = ctx.action_id,
147                method = method.as_wire_str(),
148                "pds_admin action_map routes recordAction to a label method; \
149                 cairn-mod's subscribeLabels (§F4) is the label distribution surface (§A5). \
150                 Skipping at runtime; #83's config validation already warned at startup."
151            );
152            return;
153        }
154        BackendMethod::RestoreAccount => {
155            tracing::error!(
156                action_id = ctx.action_id,
157                "pds_admin action_map routes recordAction to restore_account; \
158                 recordAction has no prior backend action id (restore is via revokeAction). \
159                 Operator config bug. Skipping."
160            );
161            return;
162        }
163        BackendMethod::TakedownAccount | BackendMethod::SuspendAccount => {}
164    }
165
166    let reason = ctx.reason_codes.first().map(String::as_str).unwrap_or("");
167
168    // Duration plumbing for SuspendAccount (#89). For other
169    // methods, duration_days is meaningless and ignored.
170    // Parse failures here are logged and dispatched as None
171    // (the suspension still goes through, but bsky-PDS's
172    // operator-visible ref field will say `duration_days=indef`
173    // which is wrong-but-safe — the cairn-mod-side action is
174    // still a temp_suspension; reconciliation is the operator's).
175    let duration_days = match (method, ctx.duration_iso) {
176        (BackendMethod::SuspendAccount, Some(iso)) => match parse_duration_iso_to_days(iso) {
177            Ok(days) => Some(days),
178            Err(e) => {
179                tracing::error!(
180                    action_id = ctx.action_id,
181                    duration_iso = iso,
182                    error = %e,
183                    "pds_admin: failed to parse temp_suspension duration; dispatching with duration_days=indef (cairn-mod-side action is still temp_suspension)"
184                );
185                None
186            }
187        },
188        _ => None,
189    };
190
191    let started_at = crate::writer::epoch_ms_now();
192    let call_result = invoke_backend_method(
193        bridge.backend.as_ref(),
194        method,
195        ctx.subject_did,
196        reason,
197        ctx.notes,
198        duration_days,
199        ctx.action_id,
200    )
201    .await;
202    let completed_at = crate::writer::epoch_ms_now();
203
204    log_call_outcome(method, ctx.action_id, ctx.subject_did, &call_result);
205
206    // Project the per-method success into the unified
207    // `Option<BackendActionId>` shape that
208    // `record_pds_admin_call` accepts. Per #87:
209    // `BackendMethod::returns_action_id()` single-sources the
210    // convention.
211    let unified: std::result::Result<Option<BackendActionId>, BackendError> =
212        match (method.returns_action_id(), call_result) {
213            (true, Ok(id)) => Ok(Some(id)),
214            (false, Ok(_)) => Ok(None),
215            (_, Err(e)) => Err(e),
216        };
217
218    if let Err(e) = record_pds_admin_call(
219        pool,
220        ctx.action_id,
221        method,
222        unified,
223        started_at,
224        completed_at,
225    )
226    .await
227    {
228        // Audit-row failure does NOT propagate — the
229        // cairn-mod-side action remains committed regardless.
230        tracing::error!(
231            error = %e,
232            action_id = ctx.action_id,
233            method = method.as_wire_str(),
234            "pds_admin audit insert failed; cairn-mod-side action remains committed"
235        );
236    }
237}
238
239/// Dispatch the trait method whose shape matches `method`. The
240/// match arms unify on `Result<BackendActionId, BackendError>`
241/// even for unit-result methods — those return a synthesized
242/// "ignored" id that the caller drops via `returns_action_id()`.
243///
244/// `duration_days` is honored only by `SuspendAccount`; the
245/// other methods ignore it (TakedownAccount has no duration
246/// concept; the `RestoreAccount` / label methods are filtered
247/// out by the caller).
248async fn invoke_backend_method(
249    backend: &dyn PdsAdminBackend,
250    method: BackendMethod,
251    did: &str,
252    reason: &str,
253    notes: Option<&str>,
254    duration_days: Option<u32>,
255    action_id: i64,
256) -> std::result::Result<BackendActionId, BackendError> {
257    match method {
258        BackendMethod::TakedownAccount => {
259            backend
260                .takedown_account(did, reason, notes, action_id)
261                .await
262        }
263        BackendMethod::SuspendAccount => {
264            backend
265                .suspend_account(did, reason, duration_days, notes, action_id)
266                .await
267        }
268        BackendMethod::RestoreAccount | BackendMethod::ApplyLabel | BackendMethod::NegateLabel => {
269            // Filtered out by the dispatch caller; reaching this
270            // arm would be a bug in this module.
271            unreachable!(
272                "invoke_backend_method dispatched non-record-action method {method:?}; \
273                 dispatch_after_record_action should have filtered it"
274            )
275        }
276    }
277}
278
279/// Parse an ISO-8601 duration string into days, rounded down.
280///
281/// Reuses [`crate::writer::parse_iso8601_duration`] (the same
282/// parser the recorder validates `RecordActionRequest.duration_iso`
283/// with at the input boundary) so cairn-mod has one source of
284/// truth for "what duration shapes are acceptable."
285///
286/// Sub-day suspensions (e.g. `PT12H`) round down to 0 days. The
287/// bsky-PDS `ref` field then encodes `duration_days=0`, which
288/// is operator-visible signal but not protocol-meaningful (the
289/// suspension lift is still cairn-mod-driven, not bsky-PDS-driven,
290/// so the day-rounding doesn't affect when the lift fires).
291/// Non-day-aligned operator config is uncommon enough to punt on.
292///
293/// Returns `Err(crate::error::Error::Signing(_))` for malformed
294/// input (matches the parser's existing error type), or for
295/// durations exceeding `u32::MAX` days (~11.7 million years —
296/// nonsense in practice but handled cleanly).
297pub(crate) fn parse_duration_iso_to_days(iso: &str) -> crate::error::Result<u32> {
298    let secs = crate::writer::parse_iso8601_duration(iso)?;
299    let days = secs / 86_400;
300    u32::try_from(days).map_err(|_| {
301        crate::error::Error::Signing(format!("duration {iso:?}: {days} days exceeds u32::MAX"))
302    })
303}
304
305// ===========================================================================
306// Revoke-action dispatch (#89)
307// ===========================================================================
308
309/// Snapshot of a revoke-action commit the PDS-admin restore
310/// dispatch needs.
311#[derive(Debug, Clone, Copy)]
312pub struct RevokeDispatchContext<'a> {
313    /// `subject_actions(id)` of the action being revoked. Used
314    /// both for the audit-row FK (the restore call's
315    /// `precipitating_action_id` is the original action's id —
316    /// "show me everything cairn-mod tried to do on the PDS for
317    /// this action" returns the takedown AND the restore in
318    /// chain order) AND for looking up the prior backend call.
319    pub action_id: i64,
320    /// Subject DID of the action being revoked.
321    pub subject_did: &'a str,
322    /// Optional revocation rationale from the moderator. Encoded
323    /// in the backend's `ref` field for cross-system traceability;
324    /// empty string when absent.
325    pub revoke_reason: Option<&'a str>,
326}
327
328/// Dispatch the post-revokeAction PDS-admin call (if any).
329///
330/// Looks up the prior `pds_admin_audit` row for the action
331/// being revoked — specifically, the most recent `success`-outcome
332/// row with a non-NULL `backend_action_id`. If found, fires
333/// [`PdsAdminBackend::restore_account`] against that
334/// `BackendActionId`. If not found (action was never propagated
335/// to the PDS, or the original call failed), logs a warning and
336/// skips — there's nothing on the PDS side to undo.
337///
338/// Same failure semantics as
339/// [`dispatch_after_record_action`]: log loudly, audit-record,
340/// don't propagate. The cairn-mod-side revocation has already
341/// committed by this point.
342pub async fn dispatch_after_revoke_action(
343    bridge: Option<&PdsAdminBridge>,
344    pool: &Pool<Sqlite>,
345    ctx: RevokeDispatchContext<'_>,
346) {
347    let Some(bridge) = bridge else {
348        return;
349    };
350    if !bridge.policy.enabled {
351        return;
352    }
353
354    // Look up prior PDS-side calls for this action. The
355    // list_pds_admin_audit_for_action API returns rows in chain
356    // order (call_completed_at ASC, ties broken on id ASC); we
357    // want the most recent success.
358    let prior_calls =
359        match crate::pds_admin::list_pds_admin_audit_for_action(pool, ctx.action_id).await {
360            Ok(rows) => rows,
361            Err(e) => {
362                tracing::error!(
363                    error = %e,
364                    action_id = ctx.action_id,
365                    "pds_admin revoke dispatch: pds_admin_audit lookup failed; \
366                     skipping restore (cairn-mod-side revocation is committed)"
367                );
368                return;
369            }
370        };
371    let prior = prior_calls
372        .iter()
373        .rev()
374        .find(|r| {
375            r.outcome == crate::pds_admin::AuditOutcome::Success && r.backend_action_id.is_some()
376        })
377        .cloned();
378
379    let Some(prior) = prior else {
380        tracing::warn!(
381            action_id = ctx.action_id,
382            "pds_admin revoke dispatch: no prior successful PDS call for this action; \
383             skipping restore (action was never propagated to PDS, or original call \
384             failed). cairn-mod-side revocation remains committed."
385        );
386        return;
387    };
388
389    // SAFETY: the find predicate above guaranteed Some.
390    let prior_action_id = prior
391        .backend_action_id
392        .expect("filter retains only rows with backend_action_id = Some");
393    let reason = ctx.revoke_reason.unwrap_or("");
394
395    let started_at = crate::writer::epoch_ms_now();
396    let call_result = bridge
397        .backend
398        .restore_account(ctx.subject_did, &prior_action_id, reason)
399        .await;
400    let completed_at = crate::writer::epoch_ms_now();
401
402    log_call_outcome(
403        BackendMethod::RestoreAccount,
404        ctx.action_id,
405        ctx.subject_did,
406        &call_result,
407    );
408
409    // Project Result<(), BackendError> into the unified
410    // Result<Option<BackendActionId>, BackendError> shape.
411    // restore_account is a unit-result method, so success carries
412    // no new BackendActionId — the audit row's backend_action_id
413    // column is None (the prior_action_id is preserved in the
414    // ref field on the wire, but the audit row's column refers to
415    // the call's RETURN value, not its input).
416    let unified: std::result::Result<Option<BackendActionId>, BackendError> = match call_result {
417        Ok(()) => Ok(None),
418        Err(e) => Err(e),
419    };
420
421    if let Err(e) = record_pds_admin_call(
422        pool,
423        ctx.action_id,
424        BackendMethod::RestoreAccount,
425        unified,
426        started_at,
427        completed_at,
428    )
429    .await
430    {
431        tracing::error!(
432            error = %e,
433            action_id = ctx.action_id,
434            method = BackendMethod::RestoreAccount.as_wire_str(),
435            "pds_admin audit insert failed for restore call; cairn-mod-side revocation \
436             remains committed"
437        );
438    }
439}
440
441/// Emit a structured tracing log line summarizing the call
442/// outcome. `WARN` for transient/network variants, `ERROR` for
443/// operator-actionable variants, `INFO` for success. Severity
444/// matches the `outcome` column the audit row will record so
445/// log-level subscribers and audit-table queries agree.
446fn log_call_outcome<T>(
447    method: BackendMethod,
448    action_id: i64,
449    subject_did: &str,
450    result: &std::result::Result<T, BackendError>,
451) {
452    match result {
453        Ok(_) => tracing::info!(
454            action_id,
455            subject_did,
456            method = method.as_wire_str(),
457            "pds_admin backend call succeeded"
458        ),
459        Err(BackendError::Unsupported(msg)) => tracing::error!(
460            action_id,
461            subject_did,
462            method = method.as_wire_str(),
463            error = %msg,
464            "pds_admin backend rejected method as unsupported (operator config issue)"
465        ),
466        Err(BackendError::Network(e)) => tracing::warn!(
467            action_id,
468            subject_did,
469            method = method.as_wire_str(),
470            error = %e,
471            "pds_admin backend call failed at the network layer (transient; not retried in v1.7)"
472        ),
473        Err(BackendError::Auth(e)) => tracing::error!(
474            action_id,
475            subject_did,
476            method = method.as_wire_str(),
477            error = %e,
478            "pds_admin backend rejected our admin auth (operator must rotate credentials)"
479        ),
480        Err(BackendError::RateLimited {
481            message,
482            retry_after_seconds,
483        }) => tracing::warn!(
484            action_id,
485            subject_did,
486            method = method.as_wire_str(),
487            retry_after_seconds = ?retry_after_seconds,
488            error = %message,
489            "pds_admin backend rate-limited the call (not retried in v1.7)"
490        ),
491        Err(BackendError::Conflict(e)) => tracing::warn!(
492            action_id,
493            subject_did,
494            method = method.as_wire_str(),
495            error = %e,
496            "pds_admin backend reported state conflict"
497        ),
498        Err(BackendError::RemoteError { code, message }) => tracing::warn!(
499            action_id,
500            subject_did,
501            method = method.as_wire_str(),
502            error_code = %code,
503            error = %message,
504            "pds_admin backend returned an unrecognized error envelope"
505        ),
506        Err(BackendError::Validation(e)) => tracing::error!(
507            action_id,
508            subject_did,
509            method = method.as_wire_str(),
510            error = %e,
511            "pds_admin backend rejected our request as malformed (cairn-mod-side bug)"
512        ),
513    }
514}
515
516#[cfg(test)]
517mod tests {
518    use super::*;
519    use crate::pds_admin::types::Subject;
520    use std::collections::BTreeMap;
521    use std::sync::Mutex;
522
523    /// Test backend that records every call and returns a
524    /// canned response. Lets us verify the dispatch fires the
525    /// right method without making any HTTP calls.
526    struct RecordingBackend {
527        calls: Mutex<Vec<RecordedCall>>,
528        takedown_response: Mutex<Option<std::result::Result<BackendActionId, BackendError>>>,
529    }
530
531    #[derive(Debug, Clone, PartialEq, Eq)]
532    struct RecordedCall {
533        method: &'static str,
534        did: String,
535        reason: String,
536        action_id: i64,
537    }
538
539    impl RecordingBackend {
540        fn new() -> Arc<Self> {
541            Arc::new(Self {
542                calls: Mutex::new(Vec::new()),
543                takedown_response: Mutex::new(None),
544            })
545        }
546
547        fn with_takedown_ok(self: &Arc<Self>, id: &str) {
548            *self.takedown_response.lock().unwrap() = Some(Ok(BackendActionId::new(id)));
549        }
550
551        fn with_takedown_err(self: &Arc<Self>, err: BackendError) {
552            *self.takedown_response.lock().unwrap() = Some(Err(err));
553        }
554
555        fn calls(&self) -> Vec<RecordedCall> {
556            self.calls.lock().unwrap().clone()
557        }
558    }
559
560    #[async_trait::async_trait]
561    impl PdsAdminBackend for RecordingBackend {
562        async fn takedown_account(
563            &self,
564            did: &str,
565            reason: &str,
566            _notes: Option<&str>,
567            precipitating_action_id: i64,
568        ) -> std::result::Result<BackendActionId, BackendError> {
569            self.calls.lock().unwrap().push(RecordedCall {
570                method: "takedown_account",
571                did: did.to_string(),
572                reason: reason.to_string(),
573                action_id: precipitating_action_id,
574            });
575            self.takedown_response
576                .lock()
577                .unwrap()
578                .take()
579                .unwrap_or_else(|| {
580                    Ok(BackendActionId::new(format!(
581                        "test:{did}:{precipitating_action_id}"
582                    )))
583                })
584        }
585
586        async fn suspend_account(
587            &self,
588            _did: &str,
589            _reason: &str,
590            _duration_days: Option<u32>,
591            _notes: Option<&str>,
592            _precipitating_action_id: i64,
593        ) -> std::result::Result<BackendActionId, BackendError> {
594            unimplemented!("test backend does not stub suspend_account")
595        }
596
597        async fn restore_account(
598            &self,
599            _did: &str,
600            _prior_action_id: &BackendActionId,
601            _reason: &str,
602        ) -> std::result::Result<(), BackendError> {
603            unimplemented!()
604        }
605
606        async fn apply_label(
607            &self,
608            _subject: &Subject,
609            _val: &str,
610            _expires_days: Option<u32>,
611        ) -> std::result::Result<(), BackendError> {
612            self.calls.lock().unwrap().push(RecordedCall {
613                method: "apply_label",
614                did: String::new(),
615                reason: String::new(),
616                action_id: 0,
617            });
618            Err(BackendError::Unsupported("test"))
619        }
620
621        async fn negate_label(
622            &self,
623            _subject: &Subject,
624            _val: &str,
625        ) -> std::result::Result<(), BackendError> {
626            unimplemented!()
627        }
628
629        async fn probe(
630            &self,
631        ) -> std::result::Result<crate::pds_admin::backend::ProbeReport, BackendError> {
632            unimplemented!("RecordingBackend test stub: probe not exercised by dispatch tests")
633        }
634    }
635
636    fn policy_with_action_map(
637        enabled: bool,
638        action_map: BTreeMap<ActionType, ActionMapEntry>,
639    ) -> PdsAdminPolicy {
640        // `backend: None` is a deliberate test-only shape. Dispatch
641        // doesn't read `policy.backend` — the trait-object backend on
642        // the `PdsAdminBridge` is what actually fires. v1.7's
643        // resolver always populates `policy.backend = Some(_)` when
644        // `enabled = true`, but for these unit tests we don't need
645        // it; the runtime invariant lives in #83's resolver, not
646        // here.
647        PdsAdminPolicy {
648            enabled,
649            backend: None,
650            action_map,
651        }
652    }
653
654    async fn fresh_pool() -> Pool<Sqlite> {
655        let dir = tempfile::tempdir().unwrap();
656        let path = dir.path().join("dispatch-test.db");
657        let pool = crate::storage::open(&path).await.unwrap();
658        Box::leak(Box::new(dir));
659        pool
660    }
661
662    async fn fixture_subject_action(pool: &Pool<Sqlite>) -> i64 {
663        sqlx::query_scalar!(
664            r#"INSERT INTO subject_actions (
665                subject_did, subject_uri, actor_did, action_type, reason_codes,
666                duration, effective_at, expires_at, notes, report_ids,
667                strike_value_base, strike_value_applied, was_dampened,
668                strikes_at_time_of_action, audit_log_id, created_at,
669                actor_kind, triggered_by_policy_rule
670             ) VALUES ('did:plc:s', NULL, 'did:plc:m', 'takedown', '["spam"]',
671                       NULL, ?1, NULL, NULL, NULL, 1, 1, 0, 1, NULL, ?1,
672                       'moderator', NULL)
673             RETURNING id AS "id!""#,
674            1_700_000_000_000_i64
675        )
676        .fetch_one(pool)
677        .await
678        .unwrap()
679    }
680
681    fn ctx<'a>(action_id: i64, did: &'a str) -> DispatchContext<'a> {
682        static REASONS: std::sync::OnceLock<Vec<String>> = std::sync::OnceLock::new();
683        let reasons = REASONS.get_or_init(|| vec!["spam".into()]);
684        DispatchContext {
685            action_id,
686            action_type: ActionType::Takedown,
687            subject_did: did,
688            reason_codes: reasons,
689            notes: None,
690            duration_iso: None,
691        }
692    }
693
694    #[tokio::test]
695    async fn dispatch_noop_when_bridge_none() {
696        let pool = fresh_pool().await;
697        let action_id = fixture_subject_action(&pool).await;
698        // Should not error, should not insert any audit row.
699        dispatch_after_record_action(None, &pool, ctx(action_id, "did:plc:s")).await;
700        let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
701            .fetch_one(&pool)
702            .await
703            .unwrap();
704        assert_eq!(count, 0);
705    }
706
707    #[tokio::test]
708    async fn dispatch_noop_when_policy_disabled() {
709        let pool = fresh_pool().await;
710        let action_id = fixture_subject_action(&pool).await;
711        let backend = RecordingBackend::new();
712        let bridge = PdsAdminBridge {
713            policy: policy_with_action_map(false, BTreeMap::new()),
714            backend: backend.clone(),
715        };
716        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
717        assert!(backend.calls().is_empty(), "no backend call when disabled");
718        let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
719            .fetch_one(&pool)
720            .await
721            .unwrap();
722        assert_eq!(count, 0);
723    }
724
725    #[tokio::test]
726    async fn dispatch_skips_when_action_type_is_skip() {
727        let pool = fresh_pool().await;
728        let action_id = fixture_subject_action(&pool).await;
729        let backend = RecordingBackend::new();
730        let mut map = BTreeMap::new();
731        map.insert(ActionType::Takedown, ActionMapEntry::Skip);
732        let bridge = PdsAdminBridge {
733            policy: policy_with_action_map(true, map),
734            backend: backend.clone(),
735        };
736        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
737        assert!(backend.calls().is_empty());
738        let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
739            .fetch_one(&pool)
740            .await
741            .unwrap();
742        assert_eq!(count, 0);
743    }
744
745    #[tokio::test]
746    async fn dispatch_takedown_records_audit_on_success() {
747        let pool = fresh_pool().await;
748        let action_id = fixture_subject_action(&pool).await;
749        let backend = RecordingBackend::new();
750        backend.with_takedown_ok("ozone:did:plc:s:42");
751        let mut map = BTreeMap::new();
752        map.insert(
753            ActionType::Takedown,
754            ActionMapEntry::Method(BackendMethod::TakedownAccount),
755        );
756        let bridge = PdsAdminBridge {
757            policy: policy_with_action_map(true, map),
758            backend: backend.clone(),
759        };
760        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
761
762        let calls = backend.calls();
763        assert_eq!(calls.len(), 1);
764        assert_eq!(calls[0].method, "takedown_account");
765        assert_eq!(calls[0].did, "did:plc:s");
766        assert_eq!(calls[0].action_id, action_id);
767
768        let audit_rows = crate::pds_admin::audit::list_pds_admin_audit_for_action(&pool, action_id)
769            .await
770            .unwrap();
771        assert_eq!(audit_rows.len(), 1);
772        assert_eq!(
773            audit_rows[0].outcome,
774            crate::pds_admin::AuditOutcome::Success
775        );
776        assert_eq!(
777            audit_rows[0].backend_action_id.as_ref().unwrap().as_str(),
778            "ozone:did:plc:s:42"
779        );
780    }
781
782    #[tokio::test]
783    async fn dispatch_takedown_records_audit_on_failure() {
784        let pool = fresh_pool().await;
785        let action_id = fixture_subject_action(&pool).await;
786        let backend = RecordingBackend::new();
787        backend.with_takedown_err(BackendError::Network("connection refused".into()));
788        let mut map = BTreeMap::new();
789        map.insert(
790            ActionType::Takedown,
791            ActionMapEntry::Method(BackendMethod::TakedownAccount),
792        );
793        let bridge = PdsAdminBridge {
794            policy: policy_with_action_map(true, map),
795            backend: backend.clone(),
796        };
797        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
798
799        let audit_rows = crate::pds_admin::audit::list_pds_admin_audit_for_action(&pool, action_id)
800            .await
801            .unwrap();
802        assert_eq!(audit_rows.len(), 1);
803        assert_eq!(
804            audit_rows[0].outcome,
805            crate::pds_admin::AuditOutcome::Network
806        );
807        assert_eq!(
808            audit_rows[0].error_message.as_deref(),
809            Some("connection refused")
810        );
811        assert!(audit_rows[0].backend_action_id.is_none());
812    }
813
814    #[tokio::test]
815    async fn dispatch_label_method_skips_with_warn() {
816        let pool = fresh_pool().await;
817        let action_id = fixture_subject_action(&pool).await;
818        let backend = RecordingBackend::new();
819        let mut map = BTreeMap::new();
820        map.insert(
821            ActionType::Takedown,
822            ActionMapEntry::Method(BackendMethod::ApplyLabel),
823        );
824        let bridge = PdsAdminBridge {
825            policy: policy_with_action_map(true, map),
826            backend: backend.clone(),
827        };
828        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
829        assert!(
830            backend.calls().is_empty(),
831            "label method must NOT actually be called from recordAction dispatch"
832        );
833        let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
834            .fetch_one(&pool)
835            .await
836            .unwrap();
837        assert_eq!(count, 0);
838    }
839
840    #[tokio::test]
841    async fn dispatch_restore_method_skips_with_error_log() {
842        let pool = fresh_pool().await;
843        let action_id = fixture_subject_action(&pool).await;
844        let backend = RecordingBackend::new();
845        let mut map = BTreeMap::new();
846        map.insert(
847            ActionType::Takedown,
848            ActionMapEntry::Method(BackendMethod::RestoreAccount),
849        );
850        let bridge = PdsAdminBridge {
851            policy: policy_with_action_map(true, map),
852            backend: backend.clone(),
853        };
854        dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
855        assert!(backend.calls().is_empty());
856    }
857
858    // ===== ISO-8601 → days parsing (#89) =====
859
860    #[test]
861    fn parse_duration_iso_to_days_handles_known_formats() {
862        // ISO-8601 day-aligned forms cover the v1.7 production
863        // range. Reuse cases the recorder's parser already
864        // exercises so any divergence between the two sites
865        // surfaces here.
866        assert_eq!(parse_duration_iso_to_days("P7D").unwrap(), 7);
867        assert_eq!(parse_duration_iso_to_days("P1W").unwrap(), 7);
868        assert_eq!(parse_duration_iso_to_days("P14D").unwrap(), 14);
869        assert_eq!(parse_duration_iso_to_days("P30D").unwrap(), 30);
870    }
871
872    #[test]
873    fn parse_duration_iso_to_days_rounds_sub_day_down() {
874        // Sub-day suspensions round down to 0. Matches the
875        // doc-comment's "sub-day → 0 days" contract; operators
876        // wanting hour-resolution suspensions are an edge case
877        // v1.7 punts on.
878        assert_eq!(parse_duration_iso_to_days("PT12H").unwrap(), 0);
879        assert_eq!(parse_duration_iso_to_days("PT30M").unwrap(), 0);
880        assert_eq!(parse_duration_iso_to_days("PT1S").unwrap(), 0);
881    }
882
883    #[test]
884    fn parse_duration_iso_to_days_rejects_malformed() {
885        assert!(parse_duration_iso_to_days("not a duration").is_err());
886        assert!(parse_duration_iso_to_days("7D").is_err()); // missing P prefix
887        assert!(parse_duration_iso_to_days("P").is_err()); // empty
888        assert!(parse_duration_iso_to_days("P1Y").is_err()); // years not supported per recorder
889    }
890
891    // ===== SuspendAccount dispatch with duration plumbing (#89) =====
892
893    fn temp_suspension_ctx<'a>(
894        action_id: i64,
895        did: &'a str,
896        iso: Option<&'a str>,
897    ) -> DispatchContext<'a> {
898        static REASONS: std::sync::OnceLock<Vec<String>> = std::sync::OnceLock::new();
899        let reasons = REASONS.get_or_init(|| vec!["spam".into()]);
900        DispatchContext {
901            action_id,
902            action_type: ActionType::TempSuspension,
903            subject_did: did,
904            reason_codes: reasons,
905            notes: None,
906            duration_iso: iso,
907        }
908    }
909
910    #[tokio::test]
911    async fn dispatch_suspend_with_duration_passes_days_to_backend() {
912        let pool = fresh_pool().await;
913        let action_id = fixture_subject_action(&pool).await;
914        let backend = SuspensionRecordingBackend::new();
915        let mut map = BTreeMap::new();
916        map.insert(
917            ActionType::TempSuspension,
918            ActionMapEntry::Method(BackendMethod::SuspendAccount),
919        );
920        let bridge = PdsAdminBridge {
921            policy: policy_with_action_map(true, map),
922            backend: backend.clone(),
923        };
924
925        dispatch_after_record_action(
926            Some(&bridge),
927            &pool,
928            temp_suspension_ctx(action_id, "did:plc:s", Some("P7D")),
929        )
930        .await;
931
932        let calls = backend.suspend_calls.lock().unwrap();
933        assert_eq!(calls.len(), 1);
934        assert_eq!(calls[0].duration_days, Some(7));
935        assert_eq!(calls[0].did, "did:plc:s");
936    }
937
938    #[tokio::test]
939    async fn dispatch_suspend_with_no_duration_passes_none() {
940        // IndefSuspension or unparseable → None passed through.
941        // (The recorder requires temp_suspension to have a
942        // duration, but the ActionMapEntry could route an
943        // IndefSuspension here in practice.)
944        let pool = fresh_pool().await;
945        let action_id = fixture_subject_action(&pool).await;
946        let backend = SuspensionRecordingBackend::new();
947        let mut map = BTreeMap::new();
948        map.insert(
949            ActionType::IndefSuspension,
950            ActionMapEntry::Method(BackendMethod::SuspendAccount),
951        );
952        let bridge = PdsAdminBridge {
953            policy: policy_with_action_map(true, map),
954            backend: backend.clone(),
955        };
956
957        let mut indef_ctx = ctx(action_id, "did:plc:s");
958        indef_ctx.action_type = ActionType::IndefSuspension;
959        dispatch_after_record_action(Some(&bridge), &pool, indef_ctx).await;
960
961        let calls = backend.suspend_calls.lock().unwrap();
962        assert_eq!(calls.len(), 1);
963        assert_eq!(calls[0].duration_days, None);
964    }
965
966    #[tokio::test]
967    async fn dispatch_suspend_with_malformed_duration_logs_and_passes_none() {
968        // Per #89's design: parse failure falls through with
969        // duration_days=None rather than aborting the dispatch.
970        // The cairn-mod-side action is still committed; the
971        // bsky-PDS ref will say `duration_days=indef` which is
972        // wrong-but-safe.
973        let pool = fresh_pool().await;
974        let action_id = fixture_subject_action(&pool).await;
975        let backend = SuspensionRecordingBackend::new();
976        let mut map = BTreeMap::new();
977        map.insert(
978            ActionType::TempSuspension,
979            ActionMapEntry::Method(BackendMethod::SuspendAccount),
980        );
981        let bridge = PdsAdminBridge {
982            policy: policy_with_action_map(true, map),
983            backend: backend.clone(),
984        };
985
986        dispatch_after_record_action(
987            Some(&bridge),
988            &pool,
989            temp_suspension_ctx(action_id, "did:plc:s", Some("not-a-duration")),
990        )
991        .await;
992
993        let calls = backend.suspend_calls.lock().unwrap();
994        assert_eq!(calls.len(), 1);
995        assert_eq!(
996            calls[0].duration_days, None,
997            "malformed duration → None (logged at error level)"
998        );
999    }
1000
1001    // ===== Revoke dispatch (#89) =====
1002
1003    #[tokio::test]
1004    async fn revoke_dispatch_noop_when_bridge_none() {
1005        let pool = fresh_pool().await;
1006        let action_id = fixture_subject_action(&pool).await;
1007        dispatch_after_revoke_action(
1008            None,
1009            &pool,
1010            RevokeDispatchContext {
1011                action_id,
1012                subject_did: "did:plc:s",
1013                revoke_reason: None,
1014            },
1015        )
1016        .await;
1017        let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
1018            .fetch_one(&pool)
1019            .await
1020            .unwrap();
1021        assert_eq!(count, 0);
1022    }
1023
1024    #[tokio::test]
1025    async fn revoke_dispatch_skips_when_no_prior_pds_call() {
1026        // The action was never propagated to the PDS (no
1027        // pds_admin_audit row); revoke has nothing to undo.
1028        let pool = fresh_pool().await;
1029        let action_id = fixture_subject_action(&pool).await;
1030        let backend = SuspensionRecordingBackend::new();
1031        let bridge = PdsAdminBridge {
1032            policy: policy_with_action_map(true, BTreeMap::new()),
1033            backend: backend.clone(),
1034        };
1035        dispatch_after_revoke_action(
1036            Some(&bridge),
1037            &pool,
1038            RevokeDispatchContext {
1039                action_id,
1040                subject_did: "did:plc:s",
1041                revoke_reason: Some("oops"),
1042            },
1043        )
1044        .await;
1045        assert!(
1046            backend.restore_calls.lock().unwrap().is_empty(),
1047            "no prior call → no restore"
1048        );
1049    }
1050
1051    #[tokio::test]
1052    async fn revoke_dispatch_calls_restore_with_prior_action_id() {
1053        // Setup: pretend a takedown succeeded for action_id by
1054        // writing the audit row directly. Then dispatch_revoke
1055        // should pick up that BackendActionId and pass it to
1056        // restore_account.
1057        let pool = fresh_pool().await;
1058        let action_id = fixture_subject_action(&pool).await;
1059
1060        crate::pds_admin::record_pds_admin_call(
1061            &pool,
1062            action_id,
1063            BackendMethod::TakedownAccount,
1064            Ok(Some(BackendActionId::new("ozone:did:plc:s:42"))),
1065            10,
1066            20,
1067        )
1068        .await
1069        .unwrap();
1070
1071        let backend = SuspensionRecordingBackend::new();
1072        let bridge = PdsAdminBridge {
1073            policy: policy_with_action_map(true, BTreeMap::new()),
1074            backend: backend.clone(),
1075        };
1076
1077        dispatch_after_revoke_action(
1078            Some(&bridge),
1079            &pool,
1080            RevokeDispatchContext {
1081                action_id,
1082                subject_did: "did:plc:s",
1083                revoke_reason: Some("manual lift"),
1084            },
1085        )
1086        .await;
1087
1088        // Snapshot + drop the mutex guard before the next await
1089        // (clippy::await_holding_lock is enforced as -D warnings).
1090        let snapshot = {
1091            let restores = backend.restore_calls.lock().unwrap();
1092            restores
1093                .iter()
1094                .map(|r| (r.did.clone(), r.prior_id.clone(), r.reason.clone()))
1095                .collect::<Vec<_>>()
1096        };
1097        assert_eq!(snapshot.len(), 1);
1098        assert_eq!(snapshot[0].0, "did:plc:s");
1099        assert_eq!(snapshot[0].1, "ozone:did:plc:s:42");
1100        assert_eq!(snapshot[0].2, "manual lift");
1101
1102        // The restore call's audit row landed.
1103        let rows = crate::pds_admin::list_pds_admin_audit_for_action(&pool, action_id)
1104            .await
1105            .unwrap();
1106        assert_eq!(rows.len(), 2, "takedown + restore audit rows");
1107        assert_eq!(rows[1].backend_method, BackendMethod::RestoreAccount);
1108        assert_eq!(rows[1].outcome, crate::pds_admin::AuditOutcome::Success);
1109        assert!(rows[1].backend_action_id.is_none());
1110    }
1111
1112    #[tokio::test]
1113    async fn revoke_dispatch_skips_when_prior_call_failed() {
1114        // Prior call recorded as Network failure (no
1115        // backend_action_id) — revoke has no id to refer to,
1116        // logs a warning and skips.
1117        let pool = fresh_pool().await;
1118        let action_id = fixture_subject_action(&pool).await;
1119        crate::pds_admin::record_pds_admin_call(
1120            &pool,
1121            action_id,
1122            BackendMethod::TakedownAccount,
1123            Err(BackendError::Network("dns".into())),
1124            10,
1125            20,
1126        )
1127        .await
1128        .unwrap();
1129
1130        let backend = SuspensionRecordingBackend::new();
1131        let bridge = PdsAdminBridge {
1132            policy: policy_with_action_map(true, BTreeMap::new()),
1133            backend: backend.clone(),
1134        };
1135        dispatch_after_revoke_action(
1136            Some(&bridge),
1137            &pool,
1138            RevokeDispatchContext {
1139                action_id,
1140                subject_did: "did:plc:s",
1141                revoke_reason: None,
1142            },
1143        )
1144        .await;
1145        assert!(
1146            backend.restore_calls.lock().unwrap().is_empty(),
1147            "no successful prior → no restore"
1148        );
1149    }
1150
1151    // Test backend that records suspend AND restore calls. Used
1152    // by the SuspendAccount + revoke tests above.
1153    struct SuspensionRecordingBackend {
1154        suspend_calls: Mutex<Vec<RecordedSuspend>>,
1155        restore_calls: Mutex<Vec<RecordedRestore>>,
1156    }
1157
1158    #[derive(Debug)]
1159    struct RecordedSuspend {
1160        did: String,
1161        duration_days: Option<u32>,
1162    }
1163
1164    #[derive(Debug)]
1165    struct RecordedRestore {
1166        did: String,
1167        prior_id: String,
1168        reason: String,
1169    }
1170
1171    impl SuspensionRecordingBackend {
1172        fn new() -> Arc<Self> {
1173            Arc::new(Self {
1174                suspend_calls: Mutex::new(Vec::new()),
1175                restore_calls: Mutex::new(Vec::new()),
1176            })
1177        }
1178    }
1179
1180    #[async_trait::async_trait]
1181    impl PdsAdminBackend for SuspensionRecordingBackend {
1182        async fn takedown_account(
1183            &self,
1184            _did: &str,
1185            _reason: &str,
1186            _notes: Option<&str>,
1187            _id: i64,
1188        ) -> std::result::Result<BackendActionId, BackendError> {
1189            unreachable!("test backend does not stub takedown")
1190        }
1191
1192        async fn suspend_account(
1193            &self,
1194            did: &str,
1195            _reason: &str,
1196            duration_days: Option<u32>,
1197            _notes: Option<&str>,
1198            id: i64,
1199        ) -> std::result::Result<BackendActionId, BackendError> {
1200            self.suspend_calls.lock().unwrap().push(RecordedSuspend {
1201                did: did.to_string(),
1202                duration_days,
1203            });
1204            Ok(BackendActionId::new(format!("ozone:{did}:{id}")))
1205        }
1206
1207        async fn restore_account(
1208            &self,
1209            did: &str,
1210            prior_action_id: &BackendActionId,
1211            reason: &str,
1212        ) -> std::result::Result<(), BackendError> {
1213            self.restore_calls.lock().unwrap().push(RecordedRestore {
1214                did: did.to_string(),
1215                prior_id: prior_action_id.as_str().to_string(),
1216                reason: reason.to_string(),
1217            });
1218            Ok(())
1219        }
1220
1221        async fn apply_label(
1222            &self,
1223            _subject: &Subject,
1224            _val: &str,
1225            _expires_days: Option<u32>,
1226        ) -> std::result::Result<(), BackendError> {
1227            unreachable!()
1228        }
1229
1230        async fn negate_label(
1231            &self,
1232            _subject: &Subject,
1233            _val: &str,
1234        ) -> std::result::Result<(), BackendError> {
1235            unreachable!()
1236        }
1237
1238        async fn probe(
1239            &self,
1240        ) -> std::result::Result<crate::pds_admin::backend::ProbeReport, BackendError> {
1241            unreachable!("SuspensionRecordingBackend test stub: probe not exercised here")
1242        }
1243    }
1244}