use std::sync::Arc;
use sqlx::{Pool, Sqlite};
use crate::moderation::types::ActionType;
use crate::pds_admin::audit::record_pds_admin_call;
use crate::pds_admin::backend::{BackendActionId, BackendError, PdsAdminBackend};
use crate::pds_admin::config::{ActionMapEntry, BackendMethod, PdsAdminPolicy};
#[derive(Clone)]
pub struct PdsAdminBridge {
pub policy: PdsAdminPolicy,
pub backend: Arc<dyn PdsAdminBackend>,
}
#[derive(Debug, Clone, Copy)]
pub struct DispatchContext<'a> {
pub action_id: i64,
pub action_type: ActionType,
pub subject_did: &'a str,
pub reason_codes: &'a [String],
pub notes: Option<&'a str>,
pub duration_iso: Option<&'a str>,
}
pub async fn dispatch_after_record_action(
bridge: Option<&PdsAdminBridge>,
pool: &Pool<Sqlite>,
ctx: DispatchContext<'_>,
) {
let Some(bridge) = bridge else {
return;
};
if !bridge.policy.enabled {
return;
}
let entry = bridge
.policy
.action_map
.get(&ctx.action_type)
.copied()
.unwrap_or(ActionMapEntry::Skip);
let method = match entry {
ActionMapEntry::Skip => return,
ActionMapEntry::Method(m) => m,
};
match method {
BackendMethod::ApplyLabel | BackendMethod::NegateLabel => {
tracing::warn!(
action_id = ctx.action_id,
method = method.as_wire_str(),
"pds_admin action_map routes recordAction to a label method; \
cairn-mod's subscribeLabels (§F4) is the label distribution surface (§A5). \
Skipping at runtime; #83's config validation already warned at startup."
);
return;
}
BackendMethod::RestoreAccount => {
tracing::error!(
action_id = ctx.action_id,
"pds_admin action_map routes recordAction to restore_account; \
recordAction has no prior backend action id (restore is via revokeAction). \
Operator config bug. Skipping."
);
return;
}
BackendMethod::TakedownAccount | BackendMethod::SuspendAccount => {}
}
let reason = ctx.reason_codes.first().map(String::as_str).unwrap_or("");
let duration_days = match (method, ctx.duration_iso) {
(BackendMethod::SuspendAccount, Some(iso)) => match parse_duration_iso_to_days(iso) {
Ok(days) => Some(days),
Err(e) => {
tracing::error!(
action_id = ctx.action_id,
duration_iso = iso,
error = %e,
"pds_admin: failed to parse temp_suspension duration; dispatching with duration_days=indef (cairn-mod-side action is still temp_suspension)"
);
None
}
},
_ => None,
};
let started_at = crate::writer::epoch_ms_now();
let call_result = invoke_backend_method(
bridge.backend.as_ref(),
method,
ctx.subject_did,
reason,
ctx.notes,
duration_days,
ctx.action_id,
)
.await;
let completed_at = crate::writer::epoch_ms_now();
log_call_outcome(method, ctx.action_id, ctx.subject_did, &call_result);
let unified: std::result::Result<Option<BackendActionId>, BackendError> =
match (method.returns_action_id(), call_result) {
(true, Ok(id)) => Ok(Some(id)),
(false, Ok(_)) => Ok(None),
(_, Err(e)) => Err(e),
};
if let Err(e) = record_pds_admin_call(
pool,
ctx.action_id,
method,
unified,
started_at,
completed_at,
)
.await
{
tracing::error!(
error = %e,
action_id = ctx.action_id,
method = method.as_wire_str(),
"pds_admin audit insert failed; cairn-mod-side action remains committed"
);
}
}
async fn invoke_backend_method(
backend: &dyn PdsAdminBackend,
method: BackendMethod,
did: &str,
reason: &str,
notes: Option<&str>,
duration_days: Option<u32>,
action_id: i64,
) -> std::result::Result<BackendActionId, BackendError> {
match method {
BackendMethod::TakedownAccount => {
backend
.takedown_account(did, reason, notes, action_id)
.await
}
BackendMethod::SuspendAccount => {
backend
.suspend_account(did, reason, duration_days, notes, action_id)
.await
}
BackendMethod::RestoreAccount | BackendMethod::ApplyLabel | BackendMethod::NegateLabel => {
unreachable!(
"invoke_backend_method dispatched non-record-action method {method:?}; \
dispatch_after_record_action should have filtered it"
)
}
}
}
pub(crate) fn parse_duration_iso_to_days(iso: &str) -> crate::error::Result<u32> {
let secs = crate::writer::parse_iso8601_duration(iso)?;
let days = secs / 86_400;
u32::try_from(days).map_err(|_| {
crate::error::Error::Signing(format!("duration {iso:?}: {days} days exceeds u32::MAX"))
})
}
#[derive(Debug, Clone, Copy)]
pub struct RevokeDispatchContext<'a> {
pub action_id: i64,
pub subject_did: &'a str,
pub revoke_reason: Option<&'a str>,
}
pub async fn dispatch_after_revoke_action(
bridge: Option<&PdsAdminBridge>,
pool: &Pool<Sqlite>,
ctx: RevokeDispatchContext<'_>,
) {
let Some(bridge) = bridge else {
return;
};
if !bridge.policy.enabled {
return;
}
let prior_calls =
match crate::pds_admin::list_pds_admin_audit_for_action(pool, ctx.action_id).await {
Ok(rows) => rows,
Err(e) => {
tracing::error!(
error = %e,
action_id = ctx.action_id,
"pds_admin revoke dispatch: pds_admin_audit lookup failed; \
skipping restore (cairn-mod-side revocation is committed)"
);
return;
}
};
let prior = prior_calls
.iter()
.rev()
.find(|r| {
r.outcome == crate::pds_admin::AuditOutcome::Success && r.backend_action_id.is_some()
})
.cloned();
let Some(prior) = prior else {
tracing::warn!(
action_id = ctx.action_id,
"pds_admin revoke dispatch: no prior successful PDS call for this action; \
skipping restore (action was never propagated to PDS, or original call \
failed). cairn-mod-side revocation remains committed."
);
return;
};
let prior_action_id = prior
.backend_action_id
.expect("filter retains only rows with backend_action_id = Some");
let reason = ctx.revoke_reason.unwrap_or("");
let started_at = crate::writer::epoch_ms_now();
let call_result = bridge
.backend
.restore_account(ctx.subject_did, &prior_action_id, reason)
.await;
let completed_at = crate::writer::epoch_ms_now();
log_call_outcome(
BackendMethod::RestoreAccount,
ctx.action_id,
ctx.subject_did,
&call_result,
);
let unified: std::result::Result<Option<BackendActionId>, BackendError> = match call_result {
Ok(()) => Ok(None),
Err(e) => Err(e),
};
if let Err(e) = record_pds_admin_call(
pool,
ctx.action_id,
BackendMethod::RestoreAccount,
unified,
started_at,
completed_at,
)
.await
{
tracing::error!(
error = %e,
action_id = ctx.action_id,
method = BackendMethod::RestoreAccount.as_wire_str(),
"pds_admin audit insert failed for restore call; cairn-mod-side revocation \
remains committed"
);
}
}
fn log_call_outcome<T>(
method: BackendMethod,
action_id: i64,
subject_did: &str,
result: &std::result::Result<T, BackendError>,
) {
match result {
Ok(_) => tracing::info!(
action_id,
subject_did,
method = method.as_wire_str(),
"pds_admin backend call succeeded"
),
Err(BackendError::Unsupported(msg)) => tracing::error!(
action_id,
subject_did,
method = method.as_wire_str(),
error = %msg,
"pds_admin backend rejected method as unsupported (operator config issue)"
),
Err(BackendError::Network(e)) => tracing::warn!(
action_id,
subject_did,
method = method.as_wire_str(),
error = %e,
"pds_admin backend call failed at the network layer (transient; not retried in v1.7)"
),
Err(BackendError::Auth(e)) => tracing::error!(
action_id,
subject_did,
method = method.as_wire_str(),
error = %e,
"pds_admin backend rejected our admin auth (operator must rotate credentials)"
),
Err(BackendError::RateLimited {
message,
retry_after_seconds,
}) => tracing::warn!(
action_id,
subject_did,
method = method.as_wire_str(),
retry_after_seconds = ?retry_after_seconds,
error = %message,
"pds_admin backend rate-limited the call (not retried in v1.7)"
),
Err(BackendError::Conflict(e)) => tracing::warn!(
action_id,
subject_did,
method = method.as_wire_str(),
error = %e,
"pds_admin backend reported state conflict"
),
Err(BackendError::RemoteError { code, message }) => tracing::warn!(
action_id,
subject_did,
method = method.as_wire_str(),
error_code = %code,
error = %message,
"pds_admin backend returned an unrecognized error envelope"
),
Err(BackendError::Validation(e)) => tracing::error!(
action_id,
subject_did,
method = method.as_wire_str(),
error = %e,
"pds_admin backend rejected our request as malformed (cairn-mod-side bug)"
),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pds_admin::types::Subject;
use std::collections::BTreeMap;
use std::sync::Mutex;
struct RecordingBackend {
calls: Mutex<Vec<RecordedCall>>,
takedown_response: Mutex<Option<std::result::Result<BackendActionId, BackendError>>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct RecordedCall {
method: &'static str,
did: String,
reason: String,
action_id: i64,
}
impl RecordingBackend {
fn new() -> Arc<Self> {
Arc::new(Self {
calls: Mutex::new(Vec::new()),
takedown_response: Mutex::new(None),
})
}
fn with_takedown_ok(self: &Arc<Self>, id: &str) {
*self.takedown_response.lock().unwrap() = Some(Ok(BackendActionId::new(id)));
}
fn with_takedown_err(self: &Arc<Self>, err: BackendError) {
*self.takedown_response.lock().unwrap() = Some(Err(err));
}
fn calls(&self) -> Vec<RecordedCall> {
self.calls.lock().unwrap().clone()
}
}
#[async_trait::async_trait]
impl PdsAdminBackend for RecordingBackend {
async fn takedown_account(
&self,
did: &str,
reason: &str,
_notes: Option<&str>,
precipitating_action_id: i64,
) -> std::result::Result<BackendActionId, BackendError> {
self.calls.lock().unwrap().push(RecordedCall {
method: "takedown_account",
did: did.to_string(),
reason: reason.to_string(),
action_id: precipitating_action_id,
});
self.takedown_response
.lock()
.unwrap()
.take()
.unwrap_or_else(|| {
Ok(BackendActionId::new(format!(
"test:{did}:{precipitating_action_id}"
)))
})
}
async fn suspend_account(
&self,
_did: &str,
_reason: &str,
_duration_days: Option<u32>,
_notes: Option<&str>,
_precipitating_action_id: i64,
) -> std::result::Result<BackendActionId, BackendError> {
unimplemented!("test backend does not stub suspend_account")
}
async fn restore_account(
&self,
_did: &str,
_prior_action_id: &BackendActionId,
_reason: &str,
) -> std::result::Result<(), BackendError> {
unimplemented!()
}
async fn apply_label(
&self,
_subject: &Subject,
_val: &str,
_expires_days: Option<u32>,
) -> std::result::Result<(), BackendError> {
self.calls.lock().unwrap().push(RecordedCall {
method: "apply_label",
did: String::new(),
reason: String::new(),
action_id: 0,
});
Err(BackendError::Unsupported("test"))
}
async fn negate_label(
&self,
_subject: &Subject,
_val: &str,
) -> std::result::Result<(), BackendError> {
unimplemented!()
}
async fn probe(
&self,
) -> std::result::Result<crate::pds_admin::backend::ProbeReport, BackendError> {
unimplemented!("RecordingBackend test stub: probe not exercised by dispatch tests")
}
}
fn policy_with_action_map(
enabled: bool,
action_map: BTreeMap<ActionType, ActionMapEntry>,
) -> PdsAdminPolicy {
PdsAdminPolicy {
enabled,
backend: None,
action_map,
}
}
async fn fresh_pool() -> Pool<Sqlite> {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("dispatch-test.db");
let pool = crate::storage::open(&path).await.unwrap();
Box::leak(Box::new(dir));
pool
}
async fn fixture_subject_action(pool: &Pool<Sqlite>) -> i64 {
sqlx::query_scalar!(
r#"INSERT INTO subject_actions (
subject_did, subject_uri, actor_did, action_type, reason_codes,
duration, effective_at, expires_at, notes, report_ids,
strike_value_base, strike_value_applied, was_dampened,
strikes_at_time_of_action, audit_log_id, created_at,
actor_kind, triggered_by_policy_rule
) VALUES ('did:plc:s', NULL, 'did:plc:m', 'takedown', '["spam"]',
NULL, ?1, NULL, NULL, NULL, 1, 1, 0, 1, NULL, ?1,
'moderator', NULL)
RETURNING id AS "id!""#,
1_700_000_000_000_i64
)
.fetch_one(pool)
.await
.unwrap()
}
fn ctx<'a>(action_id: i64, did: &'a str) -> DispatchContext<'a> {
static REASONS: std::sync::OnceLock<Vec<String>> = std::sync::OnceLock::new();
let reasons = REASONS.get_or_init(|| vec!["spam".into()]);
DispatchContext {
action_id,
action_type: ActionType::Takedown,
subject_did: did,
reason_codes: reasons,
notes: None,
duration_iso: None,
}
}
#[tokio::test]
async fn dispatch_noop_when_bridge_none() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
dispatch_after_record_action(None, &pool, ctx(action_id, "did:plc:s")).await;
let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn dispatch_noop_when_policy_disabled() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
let bridge = PdsAdminBridge {
policy: policy_with_action_map(false, BTreeMap::new()),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
assert!(backend.calls().is_empty(), "no backend call when disabled");
let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn dispatch_skips_when_action_type_is_skip() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(ActionType::Takedown, ActionMapEntry::Skip);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
assert!(backend.calls().is_empty());
let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn dispatch_takedown_records_audit_on_success() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
backend.with_takedown_ok("ozone:did:plc:s:42");
let mut map = BTreeMap::new();
map.insert(
ActionType::Takedown,
ActionMapEntry::Method(BackendMethod::TakedownAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
let calls = backend.calls();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].method, "takedown_account");
assert_eq!(calls[0].did, "did:plc:s");
assert_eq!(calls[0].action_id, action_id);
let audit_rows = crate::pds_admin::audit::list_pds_admin_audit_for_action(&pool, action_id)
.await
.unwrap();
assert_eq!(audit_rows.len(), 1);
assert_eq!(
audit_rows[0].outcome,
crate::pds_admin::AuditOutcome::Success
);
assert_eq!(
audit_rows[0].backend_action_id.as_ref().unwrap().as_str(),
"ozone:did:plc:s:42"
);
}
#[tokio::test]
async fn dispatch_takedown_records_audit_on_failure() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
backend.with_takedown_err(BackendError::Network("connection refused".into()));
let mut map = BTreeMap::new();
map.insert(
ActionType::Takedown,
ActionMapEntry::Method(BackendMethod::TakedownAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
let audit_rows = crate::pds_admin::audit::list_pds_admin_audit_for_action(&pool, action_id)
.await
.unwrap();
assert_eq!(audit_rows.len(), 1);
assert_eq!(
audit_rows[0].outcome,
crate::pds_admin::AuditOutcome::Network
);
assert_eq!(
audit_rows[0].error_message.as_deref(),
Some("connection refused")
);
assert!(audit_rows[0].backend_action_id.is_none());
}
#[tokio::test]
async fn dispatch_label_method_skips_with_warn() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(
ActionType::Takedown,
ActionMapEntry::Method(BackendMethod::ApplyLabel),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
assert!(
backend.calls().is_empty(),
"label method must NOT actually be called from recordAction dispatch"
);
let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn dispatch_restore_method_skips_with_error_log() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = RecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(
ActionType::Takedown,
ActionMapEntry::Method(BackendMethod::RestoreAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(Some(&bridge), &pool, ctx(action_id, "did:plc:s")).await;
assert!(backend.calls().is_empty());
}
#[test]
fn parse_duration_iso_to_days_handles_known_formats() {
assert_eq!(parse_duration_iso_to_days("P7D").unwrap(), 7);
assert_eq!(parse_duration_iso_to_days("P1W").unwrap(), 7);
assert_eq!(parse_duration_iso_to_days("P14D").unwrap(), 14);
assert_eq!(parse_duration_iso_to_days("P30D").unwrap(), 30);
}
#[test]
fn parse_duration_iso_to_days_rounds_sub_day_down() {
assert_eq!(parse_duration_iso_to_days("PT12H").unwrap(), 0);
assert_eq!(parse_duration_iso_to_days("PT30M").unwrap(), 0);
assert_eq!(parse_duration_iso_to_days("PT1S").unwrap(), 0);
}
#[test]
fn parse_duration_iso_to_days_rejects_malformed() {
assert!(parse_duration_iso_to_days("not a duration").is_err());
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()); }
fn temp_suspension_ctx<'a>(
action_id: i64,
did: &'a str,
iso: Option<&'a str>,
) -> DispatchContext<'a> {
static REASONS: std::sync::OnceLock<Vec<String>> = std::sync::OnceLock::new();
let reasons = REASONS.get_or_init(|| vec!["spam".into()]);
DispatchContext {
action_id,
action_type: ActionType::TempSuspension,
subject_did: did,
reason_codes: reasons,
notes: None,
duration_iso: iso,
}
}
#[tokio::test]
async fn dispatch_suspend_with_duration_passes_days_to_backend() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = SuspensionRecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(
ActionType::TempSuspension,
ActionMapEntry::Method(BackendMethod::SuspendAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(
Some(&bridge),
&pool,
temp_suspension_ctx(action_id, "did:plc:s", Some("P7D")),
)
.await;
let calls = backend.suspend_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].duration_days, Some(7));
assert_eq!(calls[0].did, "did:plc:s");
}
#[tokio::test]
async fn dispatch_suspend_with_no_duration_passes_none() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = SuspensionRecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(
ActionType::IndefSuspension,
ActionMapEntry::Method(BackendMethod::SuspendAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
let mut indef_ctx = ctx(action_id, "did:plc:s");
indef_ctx.action_type = ActionType::IndefSuspension;
dispatch_after_record_action(Some(&bridge), &pool, indef_ctx).await;
let calls = backend.suspend_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].duration_days, None);
}
#[tokio::test]
async fn dispatch_suspend_with_malformed_duration_logs_and_passes_none() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = SuspensionRecordingBackend::new();
let mut map = BTreeMap::new();
map.insert(
ActionType::TempSuspension,
ActionMapEntry::Method(BackendMethod::SuspendAccount),
);
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, map),
backend: backend.clone(),
};
dispatch_after_record_action(
Some(&bridge),
&pool,
temp_suspension_ctx(action_id, "did:plc:s", Some("not-a-duration")),
)
.await;
let calls = backend.suspend_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(
calls[0].duration_days, None,
"malformed duration → None (logged at error level)"
);
}
#[tokio::test]
async fn revoke_dispatch_noop_when_bridge_none() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
dispatch_after_revoke_action(
None,
&pool,
RevokeDispatchContext {
action_id,
subject_did: "did:plc:s",
revoke_reason: None,
},
)
.await;
let count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM pds_admin_audit")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn revoke_dispatch_skips_when_no_prior_pds_call() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
let backend = SuspensionRecordingBackend::new();
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, BTreeMap::new()),
backend: backend.clone(),
};
dispatch_after_revoke_action(
Some(&bridge),
&pool,
RevokeDispatchContext {
action_id,
subject_did: "did:plc:s",
revoke_reason: Some("oops"),
},
)
.await;
assert!(
backend.restore_calls.lock().unwrap().is_empty(),
"no prior call → no restore"
);
}
#[tokio::test]
async fn revoke_dispatch_calls_restore_with_prior_action_id() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
crate::pds_admin::record_pds_admin_call(
&pool,
action_id,
BackendMethod::TakedownAccount,
Ok(Some(BackendActionId::new("ozone:did:plc:s:42"))),
10,
20,
)
.await
.unwrap();
let backend = SuspensionRecordingBackend::new();
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, BTreeMap::new()),
backend: backend.clone(),
};
dispatch_after_revoke_action(
Some(&bridge),
&pool,
RevokeDispatchContext {
action_id,
subject_did: "did:plc:s",
revoke_reason: Some("manual lift"),
},
)
.await;
let snapshot = {
let restores = backend.restore_calls.lock().unwrap();
restores
.iter()
.map(|r| (r.did.clone(), r.prior_id.clone(), r.reason.clone()))
.collect::<Vec<_>>()
};
assert_eq!(snapshot.len(), 1);
assert_eq!(snapshot[0].0, "did:plc:s");
assert_eq!(snapshot[0].1, "ozone:did:plc:s:42");
assert_eq!(snapshot[0].2, "manual lift");
let rows = crate::pds_admin::list_pds_admin_audit_for_action(&pool, action_id)
.await
.unwrap();
assert_eq!(rows.len(), 2, "takedown + restore audit rows");
assert_eq!(rows[1].backend_method, BackendMethod::RestoreAccount);
assert_eq!(rows[1].outcome, crate::pds_admin::AuditOutcome::Success);
assert!(rows[1].backend_action_id.is_none());
}
#[tokio::test]
async fn revoke_dispatch_skips_when_prior_call_failed() {
let pool = fresh_pool().await;
let action_id = fixture_subject_action(&pool).await;
crate::pds_admin::record_pds_admin_call(
&pool,
action_id,
BackendMethod::TakedownAccount,
Err(BackendError::Network("dns".into())),
10,
20,
)
.await
.unwrap();
let backend = SuspensionRecordingBackend::new();
let bridge = PdsAdminBridge {
policy: policy_with_action_map(true, BTreeMap::new()),
backend: backend.clone(),
};
dispatch_after_revoke_action(
Some(&bridge),
&pool,
RevokeDispatchContext {
action_id,
subject_did: "did:plc:s",
revoke_reason: None,
},
)
.await;
assert!(
backend.restore_calls.lock().unwrap().is_empty(),
"no successful prior → no restore"
);
}
struct SuspensionRecordingBackend {
suspend_calls: Mutex<Vec<RecordedSuspend>>,
restore_calls: Mutex<Vec<RecordedRestore>>,
}
#[derive(Debug)]
struct RecordedSuspend {
did: String,
duration_days: Option<u32>,
}
#[derive(Debug)]
struct RecordedRestore {
did: String,
prior_id: String,
reason: String,
}
impl SuspensionRecordingBackend {
fn new() -> Arc<Self> {
Arc::new(Self {
suspend_calls: Mutex::new(Vec::new()),
restore_calls: Mutex::new(Vec::new()),
})
}
}
#[async_trait::async_trait]
impl PdsAdminBackend for SuspensionRecordingBackend {
async fn takedown_account(
&self,
_did: &str,
_reason: &str,
_notes: Option<&str>,
_id: i64,
) -> std::result::Result<BackendActionId, BackendError> {
unreachable!("test backend does not stub takedown")
}
async fn suspend_account(
&self,
did: &str,
_reason: &str,
duration_days: Option<u32>,
_notes: Option<&str>,
id: i64,
) -> std::result::Result<BackendActionId, BackendError> {
self.suspend_calls.lock().unwrap().push(RecordedSuspend {
did: did.to_string(),
duration_days,
});
Ok(BackendActionId::new(format!("ozone:{did}:{id}")))
}
async fn restore_account(
&self,
did: &str,
prior_action_id: &BackendActionId,
reason: &str,
) -> std::result::Result<(), BackendError> {
self.restore_calls.lock().unwrap().push(RecordedRestore {
did: did.to_string(),
prior_id: prior_action_id.as_str().to_string(),
reason: reason.to_string(),
});
Ok(())
}
async fn apply_label(
&self,
_subject: &Subject,
_val: &str,
_expires_days: Option<u32>,
) -> std::result::Result<(), BackendError> {
unreachable!()
}
async fn negate_label(
&self,
_subject: &Subject,
_val: &str,
) -> std::result::Result<(), BackendError> {
unreachable!()
}
async fn probe(
&self,
) -> std::result::Result<crate::pds_admin::backend::ProbeReport, BackendError> {
unreachable!("SuspensionRecordingBackend test stub: probe not exercised here")
}
}
}