1use 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#[derive(Clone)]
68pub struct PdsAdminBridge {
69 pub policy: PdsAdminPolicy,
73 pub backend: Arc<dyn PdsAdminBackend>,
77}
78
79#[derive(Debug, Clone, Copy)]
84pub struct DispatchContext<'a> {
85 pub action_id: i64,
87 pub action_type: ActionType,
89 pub subject_did: &'a str,
92 pub reason_codes: &'a [String],
97 pub notes: Option<&'a str>,
102 pub duration_iso: Option<&'a str>,
109}
110
111pub 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 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 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 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 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
239async 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 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
279pub(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#[derive(Debug, Clone, Copy)]
312pub struct RevokeDispatchContext<'a> {
313 pub action_id: i64,
320 pub subject_did: &'a str,
322 pub revoke_reason: Option<&'a str>,
326}
327
328pub 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 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 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 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
441fn 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 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 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 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 #[test]
861 fn parse_duration_iso_to_days_handles_known_formats() {
862 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 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()); assert!(parse_duration_iso_to_days("P").is_err()); assert!(parse_duration_iso_to_days("P1Y").is_err()); }
890
891 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 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 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 #[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 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 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 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 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 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 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}