use std::collections::HashSet;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use proto_blue_crypto::{K256Keypair, Keypair as _, format_multikey};
use sqlx::{Pool, Sqlite};
use time::format_description::FormatItem;
use time::macros::format_description;
use time::{OffsetDateTime, PrimitiveDateTime};
use tokio::sync::{broadcast, mpsc, oneshot, watch};
use tokio::time::{MissedTickBehavior, interval};
use uuid::Uuid;
use crate::error::{Error, Result};
use crate::label::Label;
use crate::labels::emission::{
ActionForEmission, LabelDraft, resolve_action_labels, resolve_reason_labels,
};
use crate::moderation::decay::calculate_strike_state;
use crate::moderation::policy::StrikePolicy;
use crate::moderation::reasons::ReasonVocabulary;
use crate::moderation::strike::{
StrikeApplication, calculate as strike_calculate, resolve_primary_reason,
};
use crate::moderation::types::{ActionRecord, ActionType};
use crate::moderation::window::compute_position_in_window;
use crate::server::RetentionConfig;
use crate::signing::sign_label;
use crate::signing_key::SigningKey;
pub(crate) const LEASE_STALE_MS: i64 = 60_000;
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
const SWEEP_CHECK_INTERVAL: Duration = Duration::from_secs(60);
const COMMAND_BUFFER: usize = 64;
const BROADCAST_BUFFER: usize = 1024;
#[doc(alias = "audit_log.reason")]
pub const AUDIT_REASON_SCHEMA: &str =
"label_applied / label_negated: { val, neg, moderator_reason }";
#[doc(alias = "audit_log.reason.report_resolved")]
pub const AUDIT_REASON_RESOLVE_REPORT: &str =
"report_resolved: { applied_label_val, resolution_reason }";
#[doc(alias = "audit_log.reason.flag_reporter")]
pub const AUDIT_REASON_FLAG_REPORTER: &str =
"reporter_flagged / reporter_unflagged: { did, suppressed, moderator_reason }";
#[doc(alias = "audit_log.reason.retention_sweep")]
pub const AUDIT_REASON_RETENTION_SWEEP: &str =
"retention_sweep: { rows_deleted, batches, duration_ms, retention_days_applied }";
pub const AUDIT_ACTION_VALUES: &[&str] = &[
"label_applied",
"label_negated",
"pending_policy_action_confirmed",
"pending_policy_action_dismissed",
"report_resolved",
"reporter_flagged",
"reporter_unflagged",
"retention_sweep",
"subject_action_recorded",
"subject_action_revoked",
];
pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
const CTS_FORMAT: &[FormatItem<'_>] =
format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
#[derive(Debug, Clone)]
pub struct ApplyLabelRequest {
pub actor_did: String,
pub uri: String,
pub cid: Option<String>,
pub val: String,
pub exp: Option<String>,
pub moderator_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct NegateLabelRequest {
pub actor_did: String,
pub uri: String,
pub val: String,
pub moderator_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct LabelEvent {
pub seq: i64,
pub label: Label,
}
enum WriteCommand {
Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
ResolveReport(
ResolveReportRequest,
oneshot::Sender<Result<ResolvedReport>>,
),
Sweep(SweepRequest, oneshot::Sender<Result<SweepBatchResult>>),
AppendAudit(
crate::audit::append::AuditRowForAppend,
oneshot::Sender<Result<i64>>,
),
RecordAction(RecordActionRequest, oneshot::Sender<Result<RecordedAction>>),
RevokeAction(RevokeActionRequest, oneshot::Sender<Result<RevokedAction>>),
ConfirmPendingAction(
ConfirmPendingActionRequest,
oneshot::Sender<Result<ConfirmedPendingAction>>,
),
DismissPendingAction(
DismissPendingActionRequest,
oneshot::Sender<Result<DismissedPendingAction>>,
),
Shutdown(oneshot::Sender<Result<()>>),
}
#[derive(Debug, Clone)]
pub struct ApplyLabelInline {
pub uri: String,
pub cid: Option<String>,
pub val: String,
pub exp: Option<String>,
}
#[derive(Debug, Clone)]
pub enum ResolutionAction {
Dismiss,
ApplyLabel(ApplyLabelInline),
}
impl ResolutionAction {
pub fn as_apply(&self) -> Option<&ApplyLabelInline> {
match self {
ResolutionAction::Dismiss => None,
ResolutionAction::ApplyLabel(a) => Some(a),
}
}
}
#[derive(Debug, Clone)]
pub struct ResolveReportRequest {
pub actor_did: String,
pub report_id: i64,
pub action: ResolutionAction,
pub resolution_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ResolvedReport {
pub report: crate::report::Report,
pub label_event: Option<LabelEvent>,
}
#[derive(Debug, Clone, Default)]
pub struct SweepRequest;
#[derive(Debug, Clone)]
pub struct SweepBatchResult {
pub rows_deleted: i64,
pub has_more: bool,
pub retention_days_applied: Option<u32>,
}
#[derive(Debug, Clone)]
pub struct RecordActionRequest {
pub subject: String,
pub actor_did: String,
pub action_type: ActionType,
pub reason_codes: Vec<String>,
pub duration_iso: Option<String>,
pub notes: Option<String>,
pub report_ids: Vec<i64>,
}
#[derive(Debug, Clone)]
pub struct RecordedAction {
pub action_id: i64,
pub strike_value_base: u32,
pub strike_value_applied: u32,
pub was_dampened: bool,
pub strikes_at_time_of_action: u32,
}
#[derive(Debug, Clone)]
pub struct RevokeActionRequest {
pub action_id: i64,
pub revoked_by_did: String,
pub revoked_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct RevokedAction {
pub action_id: i64,
pub revoked_at: String,
}
#[derive(Debug, Clone)]
pub struct ConfirmPendingActionRequest {
pub pending_id: i64,
pub moderator_did: String,
pub note: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ConfirmedPendingAction {
pub action_id: i64,
pub pending_id: i64,
pub resolved_at: String,
}
#[derive(Debug, Clone)]
pub struct DismissPendingActionRequest {
pub pending_id: i64,
pub moderator_did: String,
pub reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct DismissedPendingAction {
pub pending_id: i64,
pub resolved_at: String,
}
#[derive(Debug, Clone)]
pub struct SweepResult {
pub rows_deleted: i64,
pub batches: u64,
pub duration_ms: u64,
pub retention_days_applied: Option<u32>,
}
#[derive(Debug, Clone)]
pub struct WriterHandle {
tx: mpsc::Sender<WriteCommand>,
broadcast_tx: broadcast::Sender<LabelEvent>,
shutdown_rx: watch::Receiver<bool>,
}
impl WriterHandle {
pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Apply(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Negate(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::ResolveReport(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn sweep(&self, _req: SweepRequest) -> Result<SweepResult> {
let start = std::time::Instant::now();
let mut total_rows: i64 = 0;
let mut batches: u64 = 0;
loop {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Sweep(SweepRequest, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
let batch = reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))??;
total_rows += batch.rows_deleted;
batches += 1;
if !batch.has_more {
return Ok(SweepResult {
rows_deleted: total_rows,
batches,
duration_ms: start.elapsed().as_millis() as u64,
retention_days_applied: batch.retention_days_applied,
});
}
}
}
pub async fn record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::RecordAction(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::RevokeAction(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn confirm_pending_action(
&self,
req: ConfirmPendingActionRequest,
) -> Result<ConfirmedPendingAction> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::ConfirmPendingAction(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn dismiss_pending_action(
&self,
req: DismissPendingActionRequest,
) -> Result<DismissedPendingAction> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::DismissPendingAction(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn append_audit(&self, row: crate::audit::append::AuditRowForAppend) -> Result<i64> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::AppendAudit(row, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
self.broadcast_tx.subscribe()
}
pub fn receiver_count(&self) -> usize {
self.broadcast_tx.receiver_count()
}
pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
self.shutdown_rx.clone()
}
pub async fn shutdown(&self) -> Result<()> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Shutdown(reply_tx))
.await
.map_err(|_| Error::Signing("writer task is already shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
}
#[allow(clippy::too_many_arguments)]
pub async fn spawn(
pool: Pool<Sqlite>,
key: SigningKey,
service_did: String,
retention_days: Option<u32>,
retention: RetentionConfig,
reason_vocabulary: ReasonVocabulary,
strike_policy: StrikePolicy,
label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
) -> Result<WriterHandle> {
spawn_with_pds_admin(
pool,
key,
service_did,
retention_days,
retention,
reason_vocabulary,
strike_policy,
label_emission_policy,
policy_automation_policy,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn spawn_with_pds_admin(
pool: Pool<Sqlite>,
key: SigningKey,
service_did: String,
retention_days: Option<u32>,
retention: RetentionConfig,
reason_vocabulary: ReasonVocabulary,
strike_policy: StrikePolicy,
label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
pds_admin: Option<crate::pds_admin::dispatch::PdsAdminBridge>,
) -> Result<WriterHandle> {
let instance_id = acquire_lease(&pool).await?;
let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let writer = Writer {
pool,
key,
service_did,
signing_key_id,
instance_id,
rx,
broadcast_tx: broadcast_tx.clone(),
shutdown_tx,
retention_days,
retention,
reason_vocabulary,
strike_policy,
label_emission_policy,
policy_automation_policy,
pds_admin,
};
tokio::spawn(writer.run());
Ok(WriterHandle {
tx,
broadcast_tx,
shutdown_rx,
})
}
struct Writer {
pool: Pool<Sqlite>,
key: SigningKey,
service_did: String,
signing_key_id: i64,
instance_id: String,
rx: mpsc::Receiver<WriteCommand>,
broadcast_tx: broadcast::Sender<LabelEvent>,
shutdown_tx: watch::Sender<bool>,
retention_days: Option<u32>,
retention: RetentionConfig,
reason_vocabulary: ReasonVocabulary,
strike_policy: StrikePolicy,
label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
pds_admin: Option<crate::pds_admin::dispatch::PdsAdminBridge>,
}
struct SweepRunState {
started_at: Instant,
rows: i64,
batches: u64,
}
impl Writer {
fn compute_next_sweep_fire(&self) -> Option<Instant> {
if !self.retention.sweep_enabled || self.retention_days.is_none() {
return None;
}
let now_utc = OffsetDateTime::now_utc();
let target_hour = self.retention.sweep_run_at_utc_hour;
let target_today_time = time::Time::from_hms(target_hour, 0, 0)
.expect("hour validated < 24 by Config::validate");
let target_today = now_utc.replace_time(target_today_time);
let target_dt = if target_today <= now_utc {
target_today + time::Duration::days(1)
} else {
target_today
};
let wait = target_dt - now_utc;
let wait_secs = wait.whole_seconds().max(0) as u64;
Some(Instant::now() + Duration::from_secs(wait_secs))
}
async fn run(mut self) {
let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
heartbeat_timer.tick().await;
let mut sweep_check_timer = interval(SWEEP_CHECK_INTERVAL);
sweep_check_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
sweep_check_timer.tick().await;
let mut next_scheduled_fire: Option<Instant> = self.compute_next_sweep_fire();
let mut sweep_state: Option<SweepRunState> = None;
loop {
let sweep_in_progress = sweep_state.is_some();
tokio::select! {
biased;
cmd = self.rx.recv() => {
match cmd {
Some(WriteCommand::Apply(req, reply)) => {
let res = self.handle_apply(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::Negate(req, reply)) => {
let res = self.handle_negate(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::ResolveReport(req, reply)) => {
let res = self.handle_resolve_report(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::Sweep(req, reply)) => {
let res = self.handle_sweep(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::AppendAudit(row, reply)) => {
let res = self.handle_append_audit(row).await;
let _ = reply.send(res);
}
Some(WriteCommand::RecordAction(req, reply)) => {
let res = self.handle_record_action(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::RevokeAction(req, reply)) => {
let res = self.handle_revoke_action(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::ConfirmPendingAction(req, reply)) => {
let res = self.handle_confirm_pending_action(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::DismissPendingAction(req, reply)) => {
let res = self.handle_dismiss_pending_action(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::Shutdown(reply)) => {
let _ = self.shutdown_tx.send(true);
let res = self.release_lease().await;
let _ = reply.send(res);
return;
}
None => {
let _ = self.shutdown_tx.send(true);
if let Err(e) = self.release_lease().await {
tracing::error!("lease release on handle drop: {e}");
}
return;
}
}
}
_ = heartbeat_timer.tick() => {
if let Err(e) = self.heartbeat().await {
tracing::error!("lease heartbeat failed: {e}");
}
}
_ = sweep_check_timer.tick() => {
if sweep_state.is_none()
&& let Some(fire_at) = next_scheduled_fire
&& Instant::now() >= fire_at
{
sweep_state = Some(SweepRunState {
started_at: Instant::now(),
rows: 0,
batches: 0,
});
next_scheduled_fire = Some(fire_at + Duration::from_secs(86_400));
tracing::info!(
retention_days = ?self.retention_days,
sweep_batch_size = self.retention.sweep_batch_size,
"scheduled retention sweep starting"
);
}
}
_ = std::future::ready(()), if sweep_in_progress => {
match self.handle_sweep(SweepRequest).await {
Ok(batch) => {
let state = sweep_state
.as_mut()
.expect("sweep_in_progress => sweep_state Some");
state.rows += batch.rows_deleted;
state.batches += 1;
if !batch.has_more {
let final_state = sweep_state.take().expect("just set");
tracing::info!(
rows_deleted = final_state.rows,
batches = final_state.batches,
duration_ms = final_state.started_at.elapsed().as_millis() as u64,
retention_days_applied = ?batch.retention_days_applied,
"scheduled retention sweep complete"
);
}
}
Err(e) => {
let final_state = sweep_state.take();
tracing::error!(
error = %e,
batches_done = final_state.as_ref().map(|s| s.batches).unwrap_or(0),
rows_so_far = final_state.as_ref().map(|s| s.rows).unwrap_or(0),
"scheduled retention sweep batch failed; aborting run (will retry on next schedule)"
);
}
}
}
}
}
}
async fn handle_sweep(&self, _req: SweepRequest) -> Result<SweepBatchResult> {
let Some(days) = self.retention_days else {
return Ok(SweepBatchResult {
rows_deleted: 0,
has_more: false,
retention_days_applied: None,
});
};
let cutoff_ms = epoch_ms_now() - (days as i64) * 86_400_000;
let limit = self.retention.sweep_batch_size;
let mut tx = self.pool.begin().await?;
let result = sqlx::query!(
"DELETE FROM labels WHERE rowid IN (
SELECT rowid FROM labels WHERE created_at < ?1 LIMIT ?2
)",
cutoff_ms,
limit,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
let rows = result.rows_affected() as i64;
Ok(SweepBatchResult {
rows_deleted: rows,
has_more: rows >= limit,
retention_days_applied: Some(days),
})
}
async fn handle_append_audit(
&self,
row: crate::audit::append::AuditRowForAppend,
) -> Result<i64> {
let mut tx = self.pool.begin().await?;
let id = crate::audit::append::append_in_tx(&mut tx, &row).await?;
tx.commit().await?;
Ok(id)
}
async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
let mut tx = self.pool.begin().await?;
let created_at = epoch_ms_now();
let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
tx.commit().await?;
let _ = self.broadcast_tx.send(event.clone());
Ok(event)
}
async fn apply_label_inner(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
req: &ApplyLabelRequest,
created_at_ms: i64,
) -> Result<LabelEvent> {
let event = self
.sign_and_persist_label(
tx,
&req.val,
&req.uri,
req.cid.as_deref(),
false,
req.exp.as_deref(),
created_at_ms,
)
.await?;
let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
crate::audit::append::append_in_tx(
tx,
&crate::audit::append::AuditRowForAppend {
created_at: created_at_ms,
action: "label_applied".into(),
actor_did: req.actor_did.clone(),
target: Some(req.uri.clone()),
target_cid: req.cid.clone(),
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
Ok(event)
}
#[allow(clippy::too_many_arguments)]
async fn sign_and_persist_label(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
val: &str,
uri: &str,
cid: Option<&str>,
neg: bool,
exp: Option<&str>,
created_at_ms: i64,
) -> Result<LabelEvent> {
let seq = reserve_seq(tx).await?;
let prev_cts: Option<String> = sqlx::query_scalar!(
r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
self.service_did,
uri,
val,
)
.fetch_one(&mut **tx)
.await?;
let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
let cid_owned = cid.map(str::to_string);
let exp_owned = exp.map(str::to_string);
let val_owned = val.to_string();
let uri_owned = uri.to_string();
let mut label = Label {
ver: 1,
src: self.service_did.clone(),
uri: uri_owned,
cid: cid_owned,
val: val_owned,
neg,
cts,
exp: exp_owned,
sig: None,
};
label.sig = Some(sign_label(&self.key, &label)?);
let sig_bytes = label.sig.expect("just set").to_vec();
let neg_int: i64 = if neg { 1 } else { 0 };
sqlx::query!(
"INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
seq,
label.ver,
label.src,
label.uri,
label.cid,
label.val,
neg_int,
label.cts,
label.exp,
sig_bytes,
self.signing_key_id,
created_at_ms,
)
.execute(&mut **tx)
.await?;
Ok(LabelEvent { seq, label })
}
async fn sign_and_persist_label_from_draft(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
draft: &LabelDraft,
created_at_ms: i64,
) -> Result<LabelEvent> {
let exp_str = match draft.exp {
Some(st) => Some(rfc3339_from_epoch_ms(systemtime_to_epoch_ms(st)?)?),
None => None,
};
self.sign_and_persist_label(
tx,
&draft.val,
&draft.uri,
draft.cid.as_deref(),
draft.neg,
exp_str.as_deref(),
created_at_ms,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn insert_policy_auto_action(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
rule: &crate::policy::automation::PolicyRule,
subject_did: &str,
subject_uri: Option<&str>,
history_plus_precip: &[ActionRecord],
state_after_precip: &crate::moderation::decay::StrikeState,
now_systemtime: SystemTime,
created_at: i64,
predicted_auto_id: i64,
label_events: &mut Vec<LabelEvent>,
) -> Result<i64> {
let primary = resolve_primary_reason(&rule.reason_codes, &self.reason_vocabulary)?;
let strikes_at_time_of_auto = state_after_precip.current_count;
let position =
compute_position_in_window(history_plus_precip, &self.strike_policy, now_systemtime);
let calc: StrikeApplication = strike_calculate(
strikes_at_time_of_auto,
&primary,
&self.strike_policy,
position,
);
let (strike_base, strike_applied, was_dampened) = match rule.action_type {
ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
_ => (calc.base_weight, calc.applied, calc.was_dampened),
};
let auto_expires_at: Option<i64> = match (rule.action_type, rule.duration) {
(ActionType::TempSuspension, Some(d)) => Some(created_at + (d.as_secs() as i64) * 1000),
_ => None,
};
let auto_duration_iso: Option<String> = match (rule.action_type, rule.duration) {
(ActionType::TempSuspension, Some(d)) => {
Some(format!("PT{}S", d.as_secs()))
}
_ => None,
};
let action_for_emission = ActionForEmission {
action_type: rule.action_type,
expires_at: auto_expires_at.map(epoch_ms_to_systemtime),
subject_did: subject_did.to_string(),
subject_uri: subject_uri.map(str::to_string),
reason_codes: rule.reason_codes.clone(),
cid: None,
};
let action_drafts = resolve_action_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let reason_drafts = resolve_reason_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
.iter()
.chain(reason_drafts.iter())
.map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
.collect();
let next_id = predict_next_subject_action_id(tx).await?;
if next_id != predicted_auto_id {
return Err(Error::Signing(format!(
"policy auto-action: predicted id {predicted_auto_id} but next-id is now {next_id}; \
policy_consequence reservation diverged"
)));
}
let audit_reason = build_record_action_audit_reason(
next_id,
rule.action_type,
&primary.identifier,
&rule.reason_codes,
strike_base,
strike_applied,
was_dampened,
&emitted_labels_for_audit,
"policy",
Some(&rule.name),
None,
);
let audit_target = subject_uri
.map(str::to_string)
.unwrap_or_else(|| subject_did.to_string());
let auto_audit_log_id = crate::audit::append::append_in_tx(
tx,
&crate::audit::append::AuditRowForAppend {
created_at,
action: "subject_action_recorded".into(),
actor_did: crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID.to_string(),
target: Some(audit_target),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
let action_type_str = rule.action_type.as_db_str();
let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
let strikes_at_time_i64 = strikes_at_time_of_auto as i64;
let strike_base_i64 = strike_base as i64;
let strike_applied_i64 = strike_applied as i64;
let reason_codes_json = serde_json::to_string(&rule.reason_codes)
.map_err(|e| Error::Signing(format!("serialize policy reason_codes: {e}")))?;
let actor_kind_policy = "policy";
let synth_actor = crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID;
let auto_inserted_id = sqlx::query_scalar!(
"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 (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17)
RETURNING id",
subject_did,
subject_uri,
synth_actor,
action_type_str,
reason_codes_json,
auto_duration_iso,
created_at,
auto_expires_at,
None::<String>,
strike_base_i64,
strike_applied_i64,
was_dampened_int,
strikes_at_time_i64,
auto_audit_log_id,
created_at,
actor_kind_policy,
rule.name,
)
.fetch_one(&mut **tx)
.await?;
if auto_inserted_id != predicted_auto_id {
return Err(Error::Signing(format!(
"policy auto-action inserted at id {auto_inserted_id} but predicted \
{predicted_auto_id}; audit chain captured the predicted id (corrupted)"
)));
}
let mut auto_action_label_val: Option<String> = None;
for draft in &action_drafts {
let event = self
.sign_and_persist_label_from_draft(tx, draft, created_at)
.await?;
if auto_action_label_val.is_none() {
auto_action_label_val = Some(draft.val.clone());
}
label_events.push(event);
}
if let Some(ref val) = auto_action_label_val {
sqlx::query!(
"UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
val,
auto_inserted_id,
)
.execute(&mut **tx)
.await?;
}
for (draft, reason_code) in reason_drafts.iter().zip(rule.reason_codes.iter()) {
let event = self
.sign_and_persist_label_from_draft(tx, draft, created_at)
.await?;
sqlx::query!(
"INSERT INTO subject_action_reason_labels
(action_id, reason_code, emitted_label_uri, emitted_at)
VALUES (?1, ?2, ?3, ?4)",
auto_inserted_id,
reason_code,
draft.val,
created_at,
)
.execute(&mut **tx)
.await?;
label_events.push(event);
}
if rule.action_type == ActionType::Takedown {
auto_dismiss_pendings_on_takedown(
tx,
subject_did,
subject_uri,
auto_inserted_id,
crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID,
created_at,
)
.await?;
}
Ok(auto_inserted_id)
}
#[allow(clippy::too_many_arguments)]
async fn insert_pending_policy_action(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
rule: &crate::policy::automation::PolicyRule,
subject_did: &str,
subject_uri: Option<&str>,
triggering_action_id: i64,
created_at: i64,
predicted_pending_id: i64,
) -> Result<i64> {
let action_type_str = rule.action_type.as_db_str();
let duration_ms: Option<i64> = match (rule.action_type, rule.duration) {
(ActionType::TempSuspension, Some(d)) => Some((d.as_secs() as i64) * 1000),
_ => None,
};
let reason_codes_json = serde_json::to_string(&rule.reason_codes).map_err(|e| {
Error::Signing(format!("serialize policy reason_codes for pending: {e}"))
})?;
let inserted_id_opt = sqlx::query_scalar!(
r#"INSERT INTO pending_policy_actions (
subject_did, subject_uri, action_type, duration_ms,
reason_codes, triggered_by_policy_rule, triggered_at,
triggering_action_id
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
RETURNING id AS "id!: i64""#,
subject_did,
subject_uri,
action_type_str,
duration_ms,
reason_codes_json,
rule.name,
created_at,
triggering_action_id,
)
.fetch_one(&mut **tx)
.await?;
let inserted_id: i64 = inserted_id_opt;
if inserted_id != predicted_pending_id {
return Err(Error::Signing(format!(
"pending_policy_actions inserted at id {inserted_id} but predicted \
{predicted_pending_id}; policy_consequence reservation diverged"
)));
}
Ok(inserted_id)
}
async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
let mut tx = self.pool.begin().await?;
let latest = sqlx::query!(
"SELECT neg, cid FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3
ORDER BY seq DESC LIMIT 1",
self.service_did,
req.uri,
req.val,
)
.fetch_optional(&mut *tx)
.await?;
let cid = match latest {
Some(row) if row.neg == 0 => row.cid,
_ => {
return Err(Error::LabelNotFound {
src: self.service_did.clone(),
uri: req.uri,
val: req.val,
});
}
};
let seq = reserve_seq(&mut tx).await?;
let prev_cts: Option<String> = sqlx::query_scalar!(
r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
self.service_did,
req.uri,
req.val,
)
.fetch_one(&mut *tx)
.await?;
let wall_now_ms = epoch_ms_now();
let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
let mut label = Label {
ver: 1,
src: self.service_did.clone(),
uri: req.uri.clone(),
cid: cid.clone(),
val: req.val.clone(),
neg: true,
cts: cts.clone(),
exp: None, sig: None,
};
label.sig = Some(sign_label(&self.key, &label)?);
let sig_bytes = label.sig.expect("just set").to_vec();
let created_at = wall_now_ms;
let neg_int: i64 = 1;
sqlx::query!(
"INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
seq,
label.ver,
label.src,
label.uri,
label.cid,
label.val,
neg_int,
label.cts,
label.exp,
sig_bytes,
self.signing_key_id,
created_at,
)
.execute(&mut *tx)
.await?;
let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at,
action: "label_negated".into(),
actor_did: req.actor_did.clone(),
target: Some(req.uri.clone()),
target_cid: cid.clone(),
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
tx.commit().await?;
let event = LabelEvent { seq, label };
let _ = self.broadcast_tx.send(event.clone());
Ok(event)
}
async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
use crate::report::ReportStatus;
let mut tx = self.pool.begin().await?;
let current = sqlx::query_as!(
crate::report::Report,
r#"SELECT
id AS "id!: i64",
created_at AS "created_at!: String",
reported_by AS "reported_by!: String",
reason_type AS "reason_type!: String",
reason,
subject_type AS "subject_type!: String",
subject_did AS "subject_did!: String",
subject_uri,
subject_cid,
status AS "status!: ReportStatus",
resolved_at,
resolved_by,
resolution_label,
resolution_reason
FROM reports WHERE id = ?1"#,
req.report_id,
)
.fetch_optional(&mut *tx)
.await?;
let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
if report.status != ReportStatus::Pending {
return Err(Error::ReportAlreadyResolved { id: req.report_id });
}
let created_at = epoch_ms_now();
let label_event = if let Some(apply) = req.action.as_apply() {
let apply_req = ApplyLabelRequest {
actor_did: req.actor_did.clone(),
uri: apply.uri.clone(),
cid: apply.cid.clone(),
val: apply.val.clone(),
exp: apply.exp.clone(),
moderator_reason: None,
};
Some(
self.apply_label_inner(&mut tx, &apply_req, created_at)
.await?,
)
} else {
None
};
let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
let resolution_label = req.action.as_apply().map(|a| a.val.clone());
let resolved_status = ReportStatus::Resolved;
sqlx::query!(
"UPDATE reports SET
status = ?1,
resolved_at = ?2,
resolved_by = ?3,
resolution_label = ?4,
resolution_reason = ?5
WHERE id = ?6",
resolved_status,
resolved_at_rfc,
req.actor_did,
resolution_label,
req.resolution_reason,
req.report_id,
)
.execute(&mut *tx)
.await?;
let audit_reason = build_resolve_audit_reason(
req.action.as_apply().map(|a| a.val.as_str()),
req.resolution_reason.as_deref(),
);
crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at,
action: "report_resolved".into(),
actor_did: req.actor_did.clone(),
target: Some(req.report_id.to_string()),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
tx.commit().await?;
if let Some(event) = &label_event {
let _ = self.broadcast_tx.send(event.clone());
}
report.status = ReportStatus::Resolved;
report.resolved_at = Some(resolved_at_rfc);
report.resolved_by = Some(req.actor_did);
report.resolution_label = resolution_label;
report.resolution_reason = req.resolution_reason;
Ok(ResolvedReport {
report,
label_event,
})
}
async fn handle_record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
if req.reason_codes.is_empty() {
return Err(Error::Signing(
"recordAction: reason_codes must be non-empty".into(),
));
}
match req.action_type {
ActionType::TempSuspension => {
if req.duration_iso.as_deref().unwrap_or("").is_empty() {
return Err(Error::DurationRequiredForTempSuspension);
}
}
_ => {
if req.duration_iso.is_some() {
return Err(Error::DurationOnlyForTempSuspension);
}
}
}
let (subject_did, subject_uri) = route_subject(&req.subject)?;
let primary = resolve_primary_reason(&req.reason_codes, &self.reason_vocabulary)?;
let duration_secs = match (req.action_type, req.duration_iso.as_deref()) {
(ActionType::TempSuspension, Some(iso)) => Some(parse_iso8601_duration(iso)?),
_ => None,
};
let mut tx = self.pool.begin().await?;
let created_at = epoch_ms_now();
let effective_at = created_at;
let expires_at: Option<i64> = duration_secs.map(|s| effective_at + (s as i64) * 1000);
let history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
let now_systemtime = epoch_ms_to_systemtime(created_at);
let pre_state = calculate_strike_state(&history, &self.strike_policy, now_systemtime);
let strikes_at_time_of_action = pre_state.current_count;
let position = compute_position_in_window(&history, &self.strike_policy, now_systemtime);
let calc: StrikeApplication = strike_calculate(
strikes_at_time_of_action,
&primary,
&self.strike_policy,
position,
);
let (strike_base, strike_applied, was_dampened) = match req.action_type {
ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
_ => (calc.base_weight, calc.applied, calc.was_dampened),
};
let action_for_emission = ActionForEmission {
action_type: req.action_type,
expires_at: expires_at.map(epoch_ms_to_systemtime),
subject_did: subject_did.clone(),
subject_uri: subject_uri.clone(),
reason_codes: req.reason_codes.clone(),
cid: None,
};
let action_drafts = resolve_action_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let reason_drafts = resolve_reason_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
.iter()
.chain(reason_drafts.iter())
.map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
.collect();
let reason_codes_json = serde_json::to_string(&req.reason_codes)
.map_err(|e| Error::Signing(format!("serialize reason_codes: {e}")))?;
let report_ids_json = if req.report_ids.is_empty() {
None
} else {
Some(
serde_json::to_string(&req.report_ids)
.map_err(|e| Error::Signing(format!("serialize report_ids: {e}")))?,
)
};
let next_action_id = predict_next_subject_action_id(&mut tx).await?;
let policy_eval_history =
load_subject_actions_for_policy_eval(&mut tx, &subject_did).await?;
let pending_eval_actions = load_pending_for_policy_eval(&mut tx, &subject_did).await?;
let synthetic_precip = ActionRecord {
strike_value_applied: strike_applied,
effective_at: now_systemtime,
revoked_at: None,
action_type: req.action_type,
expires_at: expires_at.map(epoch_ms_to_systemtime),
was_dampened,
};
let mut history_plus_precip = history.clone();
history_plus_precip.push(synthetic_precip);
let state_after_precip =
calculate_strike_state(&history_plus_precip, &self.strike_policy, now_systemtime);
let mut policy_eval_history_plus_precip = policy_eval_history.clone();
policy_eval_history_plus_precip.push(crate::policy::evaluator::ActionForPolicyEval {
effective_at: now_systemtime,
action_type: req.action_type,
revoked_at: None,
triggered_by_policy_rule: None,
});
let firing_rule = crate::policy::evaluator::resolve_firing_rule(
&pre_state,
&state_after_precip,
&policy_eval_history_plus_precip,
&pending_eval_actions,
&self.policy_automation_policy,
);
let (policy_consequence, predicted_auto_action_id, predicted_pending_id, firing_rule_owned) =
if let Some(rule) = firing_rule {
match rule.mode {
crate::policy::automation::PolicyMode::Auto => {
let auto_id = next_action_id + 1;
let pc = serde_json::json!({
"rule_fired": rule.name,
"mode": "auto",
"auto_action_id": auto_id,
});
(Some(pc), Some(auto_id), None, Some(rule.clone()))
}
crate::policy::automation::PolicyMode::Flag => {
let pending_id = predict_next_pending_action_id(&mut tx).await?;
let pc = serde_json::json!({
"rule_fired": rule.name,
"mode": "flag",
"pending_action_id": pending_id,
});
(Some(pc), None, Some(pending_id), Some(rule.clone()))
}
}
} else {
(None, None, None, None)
};
let audit_reason = build_record_action_audit_reason(
next_action_id,
req.action_type,
&primary.identifier,
&req.reason_codes,
strike_base,
strike_applied,
was_dampened,
&emitted_labels_for_audit,
"moderator",
None,
policy_consequence,
);
let audit_target = subject_uri.clone().unwrap_or_else(|| subject_did.clone());
let audit_log_id = crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at,
action: "subject_action_recorded".into(),
actor_did: req.actor_did.clone(),
target: Some(audit_target),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
let action_type_str = req.action_type.as_db_str();
let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
let strikes_at_time_i64 = strikes_at_time_of_action as i64;
let strike_base_i64 = strike_base as i64;
let strike_applied_i64 = strike_applied as i64;
let actor_kind_moderator = "moderator";
let inserted_id = sqlx::query_scalar!(
"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 (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, NULL)
RETURNING id",
subject_did,
subject_uri,
req.actor_did,
action_type_str,
reason_codes_json,
req.duration_iso,
effective_at,
expires_at,
req.notes,
report_ids_json,
strike_base_i64,
strike_applied_i64,
was_dampened_int,
strikes_at_time_i64,
audit_log_id,
created_at,
actor_kind_moderator,
)
.fetch_one(&mut *tx)
.await?;
if inserted_id != next_action_id {
return Err(Error::Signing(format!(
"subject_actions inserted at id {inserted_id} but predicted {next_action_id}; audit chain captured the predicted id (corrupted)"
)));
}
let existing_action_label_val: Option<String> = sqlx::query_scalar!(
"SELECT emitted_label_uri FROM subject_actions WHERE id = ?1",
inserted_id,
)
.fetch_one(&mut *tx)
.await?;
let existing_reason_codes: Vec<String> = sqlx::query_scalar!(
"SELECT reason_code FROM subject_action_reason_labels WHERE action_id = ?1",
inserted_id,
)
.fetch_all(&mut *tx)
.await?;
let existing_reason_set: HashSet<&str> =
existing_reason_codes.iter().map(String::as_str).collect();
let mut label_events: Vec<LabelEvent> =
Vec::with_capacity(action_drafts.len() + reason_drafts.len());
let mut action_label_val: Option<String> = None;
if !should_skip_action_label_emission(existing_action_label_val.as_deref()) {
for draft in &action_drafts {
let event = self
.sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
.await?;
if action_label_val.is_none() {
action_label_val = Some(draft.val.clone());
}
label_events.push(event);
}
if let Some(ref val) = action_label_val {
sqlx::query!(
"UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
val,
inserted_id,
)
.execute(&mut *tx)
.await?;
}
}
for (draft, reason_code) in reason_drafts.iter().zip(req.reason_codes.iter()) {
if should_skip_reason_emission(reason_code, &existing_reason_set) {
continue;
}
let event = self
.sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
.await?;
sqlx::query!(
"INSERT INTO subject_action_reason_labels
(action_id, reason_code, emitted_label_uri, emitted_at)
VALUES (?1, ?2, ?3, ?4)",
inserted_id,
reason_code,
draft.val,
effective_at,
)
.execute(&mut *tx)
.await?;
label_events.push(event);
}
if let Some(rule) = firing_rule_owned.as_ref() {
match rule.mode {
crate::policy::automation::PolicyMode::Auto => {
let auto_id =
predicted_auto_action_id.expect("auto rule firing reserved an id above");
self.insert_policy_auto_action(
&mut tx,
rule,
&subject_did,
subject_uri.as_deref(),
&history_plus_precip,
&state_after_precip,
now_systemtime,
created_at,
auto_id,
&mut label_events,
)
.await?;
}
crate::policy::automation::PolicyMode::Flag => {
let pending_id =
predicted_pending_id.expect("flag rule firing reserved an id above");
self.insert_pending_policy_action(
&mut tx,
rule,
&subject_did,
subject_uri.as_deref(),
inserted_id,
created_at,
pending_id,
)
.await?;
}
}
}
if req.action_type == ActionType::Takedown {
auto_dismiss_pendings_on_takedown(
&mut tx,
&subject_did,
subject_uri.as_deref(),
inserted_id,
&req.actor_did,
created_at,
)
.await?;
}
let post_history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
let post_count_i64 = post_state.current_count as i64;
sqlx::query!(
"INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
VALUES (?1, ?2, ?3, ?3)
ON CONFLICT(subject_did) DO UPDATE SET
current_strike_count = excluded.current_strike_count,
last_action_at = excluded.last_action_at,
last_recompute_at = excluded.last_recompute_at",
subject_did,
post_count_i64,
created_at,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
for event in label_events {
let _ = self.broadcast_tx.send(event);
}
crate::pds_admin::dispatch::dispatch_after_record_action(
self.pds_admin.as_ref(),
&self.pool,
crate::pds_admin::dispatch::DispatchContext {
action_id: inserted_id,
action_type: req.action_type,
subject_did: &subject_did,
reason_codes: &req.reason_codes,
notes: req.notes.as_deref(),
duration_iso: req.duration_iso.as_deref(),
},
)
.await;
Ok(RecordedAction {
action_id: inserted_id,
strike_value_base: strike_base,
strike_value_applied: strike_applied,
was_dampened,
strikes_at_time_of_action,
})
}
async fn handle_revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
let mut tx = self.pool.begin().await?;
let row = sqlx::query!(
"SELECT subject_did, subject_uri, emitted_label_uri, revoked_at
FROM subject_actions WHERE id = ?1",
req.action_id,
)
.fetch_optional(&mut *tx)
.await?;
let row = row.ok_or(Error::ActionNotFound(req.action_id))?;
if row.revoked_at.is_some() {
return Err(Error::ActionAlreadyRevoked(req.action_id));
}
let revoked_at_ms = epoch_ms_now();
sqlx::query!(
"UPDATE subject_actions
SET revoked_at = ?1, revoked_by_did = ?2, revoked_reason = ?3
WHERE id = ?4",
revoked_at_ms,
req.revoked_by_did,
req.revoked_reason,
req.action_id,
)
.execute(&mut *tx)
.await?;
let label_uri = row
.subject_uri
.clone()
.unwrap_or_else(|| row.subject_did.clone());
let reason_linkage = sqlx::query!(
"SELECT reason_code, emitted_label_uri
FROM subject_action_reason_labels
WHERE action_id = ?1
ORDER BY reason_code ASC",
req.action_id,
)
.fetch_all(&mut *tx)
.await?;
let mut negated_for_audit: Vec<serde_json::Value> = Vec::new();
if let Some(ref val) = row.emitted_label_uri {
negated_for_audit.push(serde_json::json!({
"val": val,
"uri": label_uri,
}));
}
for r in &reason_linkage {
negated_for_audit.push(serde_json::json!({
"val": r.emitted_label_uri,
"uri": label_uri,
}));
}
let audit_reason = build_revoke_action_audit_reason(
req.action_id,
req.revoked_reason.as_deref(),
&negated_for_audit,
);
crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at: revoked_at_ms,
action: "subject_action_revoked".into(),
actor_did: req.revoked_by_did.clone(),
target: Some(req.action_id.to_string()),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
let mut negation_events: Vec<LabelEvent> = Vec::with_capacity(1 + reason_linkage.len());
if let Some(ref val) = row.emitted_label_uri {
let event = self
.sign_and_persist_label(
&mut tx,
val,
&label_uri,
None, true, None, revoked_at_ms,
)
.await?;
negation_events.push(event);
}
for r in &reason_linkage {
let event = self
.sign_and_persist_label(
&mut tx,
&r.emitted_label_uri,
&label_uri,
None,
true,
None,
revoked_at_ms,
)
.await?;
negation_events.push(event);
}
let post_history = load_subject_actions_for_calc(&mut tx, &row.subject_did).await?;
let now_systemtime = epoch_ms_to_systemtime(revoked_at_ms);
let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
let post_count_i64 = post_state.current_count as i64;
sqlx::query!(
"INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
VALUES (?1, ?2, ?3, ?3)
ON CONFLICT(subject_did) DO UPDATE SET
current_strike_count = excluded.current_strike_count,
last_recompute_at = excluded.last_recompute_at",
row.subject_did,
post_count_i64,
revoked_at_ms,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
for event in negation_events {
let _ = self.broadcast_tx.send(event);
}
crate::pds_admin::dispatch::dispatch_after_revoke_action(
self.pds_admin.as_ref(),
&self.pool,
crate::pds_admin::dispatch::RevokeDispatchContext {
action_id: req.action_id,
subject_did: &row.subject_did,
revoke_reason: req.revoked_reason.as_deref(),
},
)
.await;
Ok(RevokedAction {
action_id: req.action_id,
revoked_at: rfc3339_from_epoch_ms(revoked_at_ms)?,
})
}
async fn handle_confirm_pending_action(
&self,
req: ConfirmPendingActionRequest,
) -> Result<ConfirmedPendingAction> {
let mut tx = self.pool.begin().await?;
let pending = sqlx::query!(
r#"SELECT
subject_did AS "subject_did!: String",
subject_uri,
action_type AS "action_type!: String",
duration_ms,
reason_codes AS "reason_codes!: String",
triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
resolution
FROM pending_policy_actions
WHERE id = ?1"#,
req.pending_id,
)
.fetch_optional(&mut *tx)
.await?;
let pending = pending.ok_or(Error::PendingActionNotFound(req.pending_id))?;
if pending.resolution.is_some() {
return Err(Error::PendingAlreadyResolved(req.pending_id));
}
let action_type = ActionType::from_db_str(&pending.action_type).ok_or_else(|| {
Error::Signing(format!(
"pending_policy_actions row has invalid action_type {:?}",
pending.action_type
))
})?;
let history = load_subject_actions_for_calc(&mut tx, &pending.subject_did).await?;
let subject_is_takendown = history
.iter()
.any(|a| a.action_type == ActionType::Takedown && a.revoked_at.is_none());
if subject_is_takendown {
return Err(Error::SubjectTakendown(pending.subject_did.clone()));
}
let reason_codes: Vec<String> = serde_json::from_str(&pending.reason_codes)
.map_err(|e| Error::Signing(format!("parse pending reason_codes: {e}")))?;
if reason_codes.is_empty() {
return Err(Error::Signing(
"confirmPendingAction: pending row has empty reason_codes".into(),
));
}
let primary = resolve_primary_reason(&reason_codes, &self.reason_vocabulary)?;
let now_ms = epoch_ms_now();
let now_systemtime = epoch_ms_to_systemtime(now_ms);
let pre_state = calculate_strike_state(&history, &self.strike_policy, now_systemtime);
let strikes_at_time_of_action = pre_state.current_count;
let position = compute_position_in_window(&history, &self.strike_policy, now_systemtime);
let calc: StrikeApplication = strike_calculate(
strikes_at_time_of_action,
&primary,
&self.strike_policy,
position,
);
let (strike_base, strike_applied, was_dampened) = match action_type {
ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
_ => (calc.base_weight, calc.applied, calc.was_dampened),
};
let effective_at = now_ms;
let expires_at: Option<i64> = pending.duration_ms.map(|d| effective_at + d);
let duration_iso: Option<String> = pending.duration_ms.map(|d| format!("PT{}S", d / 1000));
let action_for_emission = ActionForEmission {
action_type,
expires_at: expires_at.map(epoch_ms_to_systemtime),
subject_did: pending.subject_did.clone(),
subject_uri: pending.subject_uri.clone(),
reason_codes: reason_codes.clone(),
cid: None,
};
let action_drafts = resolve_action_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let reason_drafts = resolve_reason_labels(
&action_for_emission,
&self.label_emission_policy,
now_systemtime,
);
let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
.iter()
.chain(reason_drafts.iter())
.map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
.collect();
let next_action_id = predict_next_subject_action_id(&mut tx).await?;
let audit_reason = build_pending_confirmed_audit_reason(
req.pending_id,
&pending.triggered_by_policy_rule,
next_action_id,
action_type,
&primary.identifier,
&reason_codes,
strike_base,
strike_applied,
was_dampened,
&emitted_labels_for_audit,
req.note.as_deref(),
);
let audit_target = pending
.subject_uri
.clone()
.unwrap_or_else(|| pending.subject_did.clone());
let audit_log_id = crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at: now_ms,
action: "pending_policy_action_confirmed".into(),
actor_did: req.moderator_did.clone(),
target: Some(audit_target),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
let action_type_str = action_type.as_db_str();
let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
let strikes_at_time_i64 = strikes_at_time_of_action as i64;
let strike_base_i64 = strike_base as i64;
let strike_applied_i64 = strike_applied as i64;
let reason_codes_json = serde_json::to_string(&reason_codes)
.map_err(|e| Error::Signing(format!("serialize reason_codes for confirm: {e}")))?;
let actor_kind_moderator = "moderator";
let inserted_id = sqlx::query_scalar!(
"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 (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17)
RETURNING id",
pending.subject_did,
pending.subject_uri,
req.moderator_did,
action_type_str,
reason_codes_json,
duration_iso,
effective_at,
expires_at,
req.note,
strike_base_i64,
strike_applied_i64,
was_dampened_int,
strikes_at_time_i64,
audit_log_id,
now_ms,
actor_kind_moderator,
pending.triggered_by_policy_rule,
)
.fetch_one(&mut *tx)
.await?;
if inserted_id != next_action_id {
return Err(Error::Signing(format!(
"confirm subject_actions inserted at id {inserted_id} but predicted \
{next_action_id}; audit chain captured the predicted id (corrupted)"
)));
}
let mut label_events: Vec<LabelEvent> =
Vec::with_capacity(action_drafts.len() + reason_drafts.len());
let mut action_label_val: Option<String> = None;
for draft in &action_drafts {
let event = self
.sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
.await?;
if action_label_val.is_none() {
action_label_val = Some(draft.val.clone());
}
label_events.push(event);
}
if let Some(ref val) = action_label_val {
sqlx::query!(
"UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
val,
inserted_id,
)
.execute(&mut *tx)
.await?;
}
for (draft, reason_code) in reason_drafts.iter().zip(reason_codes.iter()) {
let event = self
.sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
.await?;
sqlx::query!(
"INSERT INTO subject_action_reason_labels
(action_id, reason_code, emitted_label_uri, emitted_at)
VALUES (?1, ?2, ?3, ?4)",
inserted_id,
reason_code,
draft.val,
effective_at,
)
.execute(&mut *tx)
.await?;
label_events.push(event);
}
let confirmed_resolution = "confirmed";
sqlx::query!(
"UPDATE pending_policy_actions
SET resolution = ?1,
resolved_at = ?2,
resolved_by_did = ?3,
confirmed_action_id = ?4
WHERE id = ?5",
confirmed_resolution,
now_ms,
req.moderator_did,
inserted_id,
req.pending_id,
)
.execute(&mut *tx)
.await?;
if action_type == ActionType::Takedown {
auto_dismiss_pendings_on_takedown(
&mut tx,
&pending.subject_did,
pending.subject_uri.as_deref(),
inserted_id,
&req.moderator_did,
now_ms,
)
.await?;
}
let post_history = load_subject_actions_for_calc(&mut tx, &pending.subject_did).await?;
let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
let post_count_i64 = post_state.current_count as i64;
sqlx::query!(
"INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
VALUES (?1, ?2, ?3, ?3)
ON CONFLICT(subject_did) DO UPDATE SET
current_strike_count = excluded.current_strike_count,
last_action_at = excluded.last_action_at,
last_recompute_at = excluded.last_recompute_at",
pending.subject_did,
post_count_i64,
now_ms,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
for event in label_events {
let _ = self.broadcast_tx.send(event);
}
Ok(ConfirmedPendingAction {
action_id: inserted_id,
pending_id: req.pending_id,
resolved_at: rfc3339_from_epoch_ms(now_ms)?,
})
}
async fn handle_dismiss_pending_action(
&self,
req: DismissPendingActionRequest,
) -> Result<DismissedPendingAction> {
let mut tx = self.pool.begin().await?;
let pending = sqlx::query!(
r#"SELECT
subject_did AS "subject_did!: String",
subject_uri,
action_type AS "action_type!: String",
reason_codes AS "reason_codes!: String",
triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
resolution
FROM pending_policy_actions
WHERE id = ?1"#,
req.pending_id,
)
.fetch_optional(&mut *tx)
.await?;
let pending = pending.ok_or(Error::PendingActionNotFound(req.pending_id))?;
if pending.resolution.is_some() {
return Err(Error::PendingAlreadyResolved(req.pending_id));
}
let now_ms = epoch_ms_now();
let reason_codes_for_audit: serde_json::Value =
serde_json::from_str(&pending.reason_codes).unwrap_or(serde_json::Value::Null);
let audit_reason = build_pending_dismissed_audit_reason(
req.pending_id,
&pending.triggered_by_policy_rule,
&pending.action_type,
reason_codes_for_audit,
req.reason.as_deref(),
);
let audit_target = pending
.subject_uri
.clone()
.unwrap_or_else(|| pending.subject_did.clone());
crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at: now_ms,
action: "pending_policy_action_dismissed".into(),
actor_did: req.moderator_did.clone(),
target: Some(audit_target),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
let dismissed_resolution = "dismissed";
sqlx::query!(
"UPDATE pending_policy_actions
SET resolution = ?1,
resolved_at = ?2,
resolved_by_did = ?3
WHERE id = ?4",
dismissed_resolution,
now_ms,
req.moderator_did,
req.pending_id,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
Ok(DismissedPendingAction {
pending_id: req.pending_id,
resolved_at: rfc3339_from_epoch_ms(now_ms)?,
})
}
async fn heartbeat(&self) -> Result<()> {
let now_ms = epoch_ms_now();
sqlx::query!(
"UPDATE server_instance_lease SET last_heartbeat = ?1
WHERE id = 1 AND instance_id = ?2",
now_ms,
self.instance_id,
)
.execute(&self.pool)
.await?;
Ok(())
}
async fn release_lease(&self) -> Result<()> {
sqlx::query!(
"DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
self.instance_id,
)
.execute(&self.pool)
.await?;
Ok(())
}
}
pub(crate) async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
let now_ms = epoch_ms_now();
let existing =
sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
.fetch_optional(pool)
.await?;
if let Some(row) = existing {
let age_ms = (now_ms - row.last_heartbeat).max(0);
if age_ms < LEASE_STALE_MS {
return Err(Error::LeaseHeld {
instance_id: row.instance_id,
age_secs: (age_ms / 1000) as u64,
});
}
}
let new_id = Uuid::new_v4().to_string();
sqlx::query!(
"INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
VALUES (1, ?1, ?2, ?2)
ON CONFLICT(id) DO UPDATE SET
instance_id = excluded.instance_id,
acquired_at = excluded.acquired_at,
last_heartbeat = excluded.last_heartbeat",
new_id,
now_ms,
)
.execute(pool)
.await?;
Ok(new_id)
}
pub(crate) async fn release_lease_by_id(pool: &Pool<Sqlite>, instance_id: &str) -> Result<()> {
sqlx::query!(
"DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
instance_id,
)
.execute(pool)
.await?;
Ok(())
}
async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
let kp = K256Keypair::from_private_key(key.expose_secret())?;
let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
let existing =
sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
.fetch_optional(pool)
.await?;
if let Some(row) = existing {
if row.public_key_multibase != my_multibase {
return Err(Error::Signing(format!(
"signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
row.public_key_multibase, my_multibase
)));
}
return Ok(row.id);
}
let now_ms = epoch_ms_now();
let valid_from = rfc3339_from_epoch_ms(now_ms)?;
let id = sqlx::query_scalar!(
"INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
VALUES (?1, ?2, NULL, ?3)
RETURNING id",
my_multibase,
valid_from,
now_ms,
)
.fetch_one(pool)
.await?;
Ok(id)
}
async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
.fetch_one(tx)
.await?;
Ok(seq)
}
fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
let effective_ms = match prev_cts_str {
Some(s) => {
let prev_ms = parse_rfc3339_ms(s)?;
wall_now_ms.max(prev_ms + 1)
}
None => wall_now_ms,
};
rfc3339_from_epoch_ms(effective_ms)
}
pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
let nanos: i128 = (ms as i128) * 1_000_000;
let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
.map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
let formatted = dt
.format(&CTS_FORMAT)
.map_err(|e| Error::Signing(format!("format cts: {e}")))?;
Ok(format!("{formatted}Z"))
}
pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
let stripped = s
.strip_suffix('Z')
.ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
.map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
let nanos = pdt.assume_utc().unix_timestamp_nanos();
Ok((nanos / 1_000_000) as i64)
}
pub(crate) fn epoch_ms_now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system clock before unix epoch")
.as_millis() as i64
}
fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
let body = serde_json::json!({
"val": val,
"neg": neg,
"moderator_reason": moderator_reason,
});
body.to_string()
}
fn build_resolve_audit_reason(
applied_label_val: Option<&str>,
resolution_reason: Option<&str>,
) -> String {
serde_json::json!({
"applied_label_val": applied_label_val,
"resolution_reason": resolution_reason,
})
.to_string()
}
pub(crate) fn build_retention_sweep_audit_reason(result: &SweepResult) -> String {
serde_json::json!({
"rows_deleted": result.rows_deleted,
"batches": result.batches,
"duration_ms": result.duration_ms,
"retention_days_applied": result.retention_days_applied,
})
.to_string()
}
#[doc(alias = "audit_log.reason.subject_action_recorded")]
pub const AUDIT_REASON_RECORD_ACTION: &str = "subject_action_recorded: { action_id, action_type, primary_reason, reason_codes, strike_value_base, strike_value_applied, was_dampened }";
#[doc(alias = "audit_log.reason.subject_action_revoked")]
pub const AUDIT_REASON_REVOKE_ACTION: &str =
"subject_action_revoked: { action_id, revoked_reason }";
#[doc(alias = "audit_log.reason.pending_policy_action_confirmed")]
pub const AUDIT_REASON_CONFIRM_PENDING_ACTION: &str = "pending_policy_action_confirmed: { pending_id, action_id, triggered_by_policy_rule, action_type, primary_reason, reason_codes, strike_value_base, strike_value_applied, was_dampened, emitted_labels, moderator_note }";
#[doc(alias = "audit_log.reason.pending_policy_action_dismissed")]
pub const AUDIT_REASON_DISMISS_PENDING_ACTION: &str = "pending_policy_action_dismissed: { pending_id, triggered_by_policy_rule, action_type, reason_codes, triggered_by, moderator_reason | takedown_action_id }";
pub(crate) fn build_flag_reporter_audit_reason(
did: &str,
suppressed: bool,
moderator_reason: Option<&str>,
) -> String {
serde_json::json!({
"did": did,
"suppressed": suppressed,
"moderator_reason": moderator_reason,
})
.to_string()
}
fn should_skip_action_label_emission(existing_emitted_label_uri: Option<&str>) -> bool {
existing_emitted_label_uri.is_some()
}
fn should_skip_reason_emission(reason_code: &str, existing: &HashSet<&str>) -> bool {
existing.contains(reason_code)
}
#[allow(clippy::too_many_arguments)]
fn build_record_action_audit_reason(
action_id: i64,
action_type: ActionType,
primary_reason: &str,
reason_codes: &[String],
strike_value_base: u32,
strike_value_applied: u32,
was_dampened: bool,
emitted_labels: &[serde_json::Value],
actor_kind: &str,
triggered_by_policy_rule: Option<&str>,
policy_consequence: Option<serde_json::Value>,
) -> String {
let mut obj = serde_json::json!({
"action_id": action_id,
"action_type": action_type.as_db_str(),
"primary_reason": primary_reason,
"reason_codes": reason_codes,
"strike_value_base": strike_value_base,
"strike_value_applied": strike_value_applied,
"was_dampened": was_dampened,
"emitted_labels": emitted_labels,
"actor_kind": actor_kind,
});
if let Some(rule_name) = triggered_by_policy_rule {
obj["triggered_by_policy_rule"] = serde_json::Value::String(rule_name.to_string());
}
if let Some(pc) = policy_consequence {
obj["policy_consequence"] = pc;
}
obj.to_string()
}
fn build_revoke_action_audit_reason(
action_id: i64,
revoked_reason: Option<&str>,
negated_labels: &[serde_json::Value],
) -> String {
serde_json::json!({
"action_id": action_id,
"revoked_reason": revoked_reason,
"negated_labels": negated_labels,
})
.to_string()
}
#[allow(clippy::too_many_arguments)]
fn build_pending_confirmed_audit_reason(
pending_id: i64,
triggered_by_policy_rule: &str,
action_id: i64,
action_type: ActionType,
primary_reason: &str,
reason_codes: &[String],
strike_value_base: u32,
strike_value_applied: u32,
was_dampened: bool,
emitted_labels: &[serde_json::Value],
moderator_note: Option<&str>,
) -> String {
serde_json::json!({
"pending_id": pending_id,
"action_id": action_id,
"triggered_by_policy_rule": triggered_by_policy_rule,
"action_type": action_type.as_db_str(),
"primary_reason": primary_reason,
"reason_codes": reason_codes,
"strike_value_base": strike_value_base,
"strike_value_applied": strike_value_applied,
"was_dampened": was_dampened,
"emitted_labels": emitted_labels,
"moderator_note": moderator_note,
})
.to_string()
}
fn build_pending_dismissed_audit_reason(
pending_id: i64,
triggered_by_policy_rule: &str,
action_type: &str,
reason_codes_json: serde_json::Value,
moderator_reason: Option<&str>,
) -> String {
serde_json::json!({
"pending_id": pending_id,
"triggered_by_policy_rule": triggered_by_policy_rule,
"action_type": action_type,
"reason_codes": reason_codes_json,
"moderator_reason": moderator_reason,
"triggered_by": "moderator_dismissed",
})
.to_string()
}
fn build_pending_dismissed_on_takedown_audit_reason(
pending_id: i64,
triggered_by_policy_rule: &str,
action_type: &str,
reason_codes_json: serde_json::Value,
takedown_action_id: i64,
) -> String {
serde_json::json!({
"pending_id": pending_id,
"triggered_by_policy_rule": triggered_by_policy_rule,
"action_type": action_type,
"reason_codes": reason_codes_json,
"triggered_by": "takedown_terminal",
"takedown_action_id": takedown_action_id,
})
.to_string()
}
fn epoch_ms_to_systemtime(ms: i64) -> SystemTime {
if ms >= 0 {
UNIX_EPOCH + Duration::from_millis(ms as u64)
} else {
UNIX_EPOCH
}
}
fn systemtime_to_epoch_ms(st: SystemTime) -> Result<i64> {
st.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.map_err(|e| Error::Signing(format!("SystemTime before unix epoch: {e}")))
}
fn route_subject(subject: &str) -> Result<(String, Option<String>)> {
if let Some(rest) = subject.strip_prefix("at://") {
let repo = rest.split('/').next().unwrap_or("");
if repo.is_empty() || !repo.starts_with("did:") {
return Err(Error::Signing(format!(
"subject at://-URI must start at://did:.../...; got {subject:?}"
)));
}
Ok((repo.to_string(), Some(subject.to_string())))
} else if subject.starts_with("did:") && subject.len() > "did:".len() {
Ok((subject.to_string(), None))
} else {
Err(Error::Signing(format!(
"subject must be a DID (`did:...`) or AT-URI (`at://did:...`); got {subject:?}"
)))
}
}
pub(crate) fn parse_iso8601_duration(s: &str) -> Result<u64> {
let body = s.strip_prefix('P').ok_or_else(|| {
Error::Signing(format!(
"duration {s:?} must start with 'P' (ISO-8601 duration form, e.g. P7D)"
))
})?;
if body.is_empty() {
return Err(Error::Signing(format!(
"duration {s:?} has no components after 'P'"
)));
}
let (date_part, time_part) = match body.split_once('T') {
Some((d, t)) => (d, Some(t)),
None => (body, None),
};
let mut total_secs: u64 = 0;
let mut buf = String::new();
for c in date_part.chars() {
if c.is_ascii_digit() {
buf.push(c);
continue;
}
let n: u64 = buf
.parse()
.map_err(|_| Error::Signing(format!("duration {s:?}: malformed numeric run")))?;
buf.clear();
let mult = match c {
'D' => 86_400u64,
'W' => 7 * 86_400u64,
'Y' | 'M' => {
return Err(Error::Signing(format!(
"duration {s:?}: years (Y) and months (M, without T prefix) not supported in v1.4"
)));
}
other => {
return Err(Error::Signing(format!(
"duration {s:?}: unknown date-part unit {other:?}"
)));
}
};
total_secs = total_secs.saturating_add(n.saturating_mul(mult));
}
if !buf.is_empty() {
return Err(Error::Signing(format!(
"duration {s:?}: trailing digits without unit in date part"
)));
}
if let Some(time_body) = time_part {
if time_body.is_empty() {
return Err(Error::Signing(format!(
"duration {s:?}: 'T' separator with no time components"
)));
}
for c in time_body.chars() {
if c.is_ascii_digit() {
buf.push(c);
continue;
}
let n: u64 = buf.parse().map_err(|_| {
Error::Signing(format!(
"duration {s:?}: malformed numeric run in time part"
))
})?;
buf.clear();
let mult = match c {
'H' => 3600u64,
'M' => 60u64,
'S' => 1u64,
other => {
return Err(Error::Signing(format!(
"duration {s:?}: unknown time-part unit {other:?}"
)));
}
};
total_secs = total_secs.saturating_add(n.saturating_mul(mult));
}
if !buf.is_empty() {
return Err(Error::Signing(format!(
"duration {s:?}: trailing digits without unit in time part"
)));
}
}
if total_secs == 0 {
return Err(Error::Signing(format!(
"duration {s:?}: parsed to zero — must be at least one second"
)));
}
Ok(total_secs)
}
async fn predict_next_subject_action_id(tx: &mut sqlx::Transaction<'_, Sqlite>) -> Result<i64> {
let row = sqlx::query!(
r#"SELECT seq AS "seq!: i64" FROM sqlite_sequence WHERE name = 'subject_actions'"#
)
.fetch_optional(&mut **tx)
.await?;
Ok(row.map(|r| r.seq + 1).unwrap_or(1))
}
async fn predict_next_pending_action_id(tx: &mut sqlx::Transaction<'_, Sqlite>) -> Result<i64> {
let row = sqlx::query!(
r#"SELECT seq AS "seq!: i64" FROM sqlite_sequence WHERE name = 'pending_policy_actions'"#
)
.fetch_optional(&mut **tx)
.await?;
Ok(row.map(|r| r.seq + 1).unwrap_or(1))
}
async fn load_subject_actions_for_calc(
tx: &mut sqlx::Transaction<'_, Sqlite>,
subject_did: &str,
) -> Result<Vec<ActionRecord>> {
let rows = sqlx::query!(
"SELECT action_type, strike_value_applied, was_dampened,
effective_at, expires_at, revoked_at
FROM subject_actions
WHERE subject_did = ?1
ORDER BY id ASC",
subject_did,
)
.fetch_all(&mut **tx)
.await?;
let mut out = Vec::with_capacity(rows.len());
for r in rows {
let action_type = ActionType::from_db_str(&r.action_type).ok_or_else(|| {
Error::Signing(format!(
"subject_actions row has invalid action_type {:?}",
r.action_type
))
})?;
let strike_value_applied = u32::try_from(r.strike_value_applied).map_err(|_| {
Error::Signing(format!(
"subject_actions row strike_value_applied {} out of u32 range",
r.strike_value_applied
))
})?;
out.push(ActionRecord {
strike_value_applied,
effective_at: epoch_ms_to_systemtime(r.effective_at),
revoked_at: r.revoked_at.map(epoch_ms_to_systemtime),
action_type,
expires_at: r.expires_at.map(epoch_ms_to_systemtime),
was_dampened: r.was_dampened != 0,
});
}
Ok(out)
}
async fn load_subject_actions_for_policy_eval(
tx: &mut sqlx::Transaction<'_, Sqlite>,
subject_did: &str,
) -> Result<Vec<crate::policy::evaluator::ActionForPolicyEval>> {
let rows = sqlx::query!(
r#"SELECT
action_type AS "action_type!: String",
effective_at AS "effective_at!: i64",
revoked_at,
triggered_by_policy_rule
FROM subject_actions
WHERE subject_did = ?1
ORDER BY id ASC"#,
subject_did,
)
.fetch_all(&mut **tx)
.await?;
let mut out = Vec::with_capacity(rows.len());
for r in rows {
let action_type = ActionType::from_db_str(&r.action_type).ok_or_else(|| {
Error::Signing(format!(
"subject_actions row has invalid action_type {:?}",
r.action_type
))
})?;
out.push(crate::policy::evaluator::ActionForPolicyEval {
effective_at: epoch_ms_to_systemtime(r.effective_at),
action_type,
revoked_at: r.revoked_at.map(epoch_ms_to_systemtime),
triggered_by_policy_rule: r.triggered_by_policy_rule,
});
}
Ok(out)
}
async fn load_pending_for_policy_eval(
tx: &mut sqlx::Transaction<'_, Sqlite>,
subject_did: &str,
) -> Result<Vec<crate::policy::evaluator::PendingActionForPolicyEval>> {
let rows = sqlx::query!(
r#"SELECT
triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
resolution
FROM pending_policy_actions
WHERE subject_did = ?1
ORDER BY triggered_at ASC"#,
subject_did,
)
.fetch_all(&mut **tx)
.await?;
let mut out = Vec::with_capacity(rows.len());
for r in rows {
let resolution = match r.resolution.as_deref() {
None => None,
Some("confirmed") => Some(crate::policy::evaluator::PendingResolution::Confirmed),
Some("dismissed") => Some(crate::policy::evaluator::PendingResolution::Dismissed),
Some(other) => {
return Err(Error::Signing(format!(
"pending_policy_actions row has invalid resolution {:?}",
other
)));
}
};
out.push(crate::policy::evaluator::PendingActionForPolicyEval {
triggered_by_policy_rule: r.triggered_by_policy_rule,
resolution,
});
}
Ok(out)
}
async fn auto_dismiss_pendings_on_takedown(
tx: &mut sqlx::Transaction<'_, Sqlite>,
subject_did: &str,
subject_uri: Option<&str>,
triggering_takedown_id: i64,
actor_did: &str,
now_ms: i64,
) -> Result<Vec<i64>> {
let rows = sqlx::query!(
r#"SELECT
id AS "id!: i64",
triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
action_type AS "action_type!: String",
reason_codes AS "reason_codes!: String"
FROM pending_policy_actions
WHERE subject_did = ?1 AND resolution IS NULL
ORDER BY id ASC"#,
subject_did,
)
.fetch_all(&mut **tx)
.await?;
if rows.is_empty() {
return Ok(Vec::new());
}
let audit_target = subject_uri
.map(str::to_string)
.unwrap_or_else(|| subject_did.to_string());
let dismissed_resolution = "dismissed";
let mut dismissed_ids = Vec::with_capacity(rows.len());
for r in rows {
let reason_codes_json: serde_json::Value =
serde_json::from_str(&r.reason_codes).unwrap_or(serde_json::Value::Null);
let audit_reason = build_pending_dismissed_on_takedown_audit_reason(
r.id,
&r.triggered_by_policy_rule,
&r.action_type,
reason_codes_json,
triggering_takedown_id,
);
crate::audit::append::append_in_tx(
tx,
&crate::audit::append::AuditRowForAppend {
created_at: now_ms,
action: "pending_policy_action_dismissed".into(),
actor_did: actor_did.to_string(),
target: Some(audit_target.clone()),
target_cid: None,
outcome: "success".into(),
reason: Some(audit_reason),
},
)
.await?;
sqlx::query!(
"UPDATE pending_policy_actions
SET resolution = ?1,
resolved_at = ?2,
resolved_by_did = ?3
WHERE id = ?4",
dismissed_resolution,
now_ms,
actor_did,
r.id,
)
.execute(&mut **tx)
.await?;
dismissed_ids.push(r.id);
}
Ok(dismissed_ids)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn clamp_cts_uses_wall_clock_when_no_prior() {
let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.000Z");
}
#[test]
fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.501Z");
}
#[test]
fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:22.000Z");
}
#[test]
fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.501Z");
}
#[test]
fn build_audit_reason_shape_matches_documented_schema() {
let json = build_audit_reason("spam", false, Some("user reported"));
let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
assert_eq!(v["val"], "spam");
assert_eq!(v["neg"], false);
assert_eq!(v["moderator_reason"], "user reported");
}
#[test]
fn build_audit_reason_null_when_moderator_reason_absent() {
let json = build_audit_reason("spam", true, None);
let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
assert!(v["moderator_reason"].is_null());
}
#[test]
fn rfc3339_roundtrip() {
let s = "2026-04-22T12:00:00.938Z";
let ms = parse_rfc3339_ms(s).expect("parse");
let back = rfc3339_from_epoch_ms(ms).expect("format");
assert_eq!(back, s);
}
#[test]
fn route_subject_did_only_returns_did_no_uri() {
let (did, uri) = route_subject("did:plc:abc123").unwrap();
assert_eq!(did, "did:plc:abc123");
assert!(uri.is_none());
}
#[test]
fn route_subject_at_uri_extracts_repo_did_and_keeps_full_uri() {
let (did, uri) = route_subject("at://did:plc:abc/app.bsky.feed.post/3xx").unwrap();
assert_eq!(did, "did:plc:abc");
assert_eq!(
uri.as_deref(),
Some("at://did:plc:abc/app.bsky.feed.post/3xx")
);
}
#[test]
fn route_subject_rejects_at_uri_with_non_did_repo() {
let err = route_subject("at://handle.example/app.bsky.feed.post/3xx").unwrap_err();
assert!(matches!(err, Error::Signing(_)));
}
#[test]
fn route_subject_rejects_arbitrary_string() {
assert!(route_subject("not a subject").is_err());
assert!(route_subject("did:").is_err());
assert!(route_subject("").is_err());
}
#[test]
fn duration_p7d_is_seven_days() {
assert_eq!(parse_iso8601_duration("P7D").unwrap(), 7 * 86_400);
}
#[test]
fn duration_p1w_is_seven_days() {
assert_eq!(parse_iso8601_duration("P1W").unwrap(), 7 * 86_400);
}
#[test]
fn duration_pt12h_is_twelve_hours() {
assert_eq!(parse_iso8601_duration("PT12H").unwrap(), 12 * 3600);
}
#[test]
fn duration_pt30m_is_thirty_minutes() {
assert_eq!(parse_iso8601_duration("PT30M").unwrap(), 30 * 60);
}
#[test]
fn duration_pt45s_is_forty_five_seconds() {
assert_eq!(parse_iso8601_duration("PT45S").unwrap(), 45);
}
#[test]
fn duration_p1d_t12h_combines() {
assert_eq!(
parse_iso8601_duration("P1DT12H").unwrap(),
86_400 + 12 * 3600
);
}
#[test]
fn duration_no_p_prefix_rejected() {
assert!(parse_iso8601_duration("7D").is_err());
assert!(parse_iso8601_duration("").is_err());
}
#[test]
fn duration_p_alone_rejected() {
assert!(parse_iso8601_duration("P").is_err());
}
#[test]
fn duration_year_rejected_in_v14() {
let err = parse_iso8601_duration("P1Y").unwrap_err();
let msg = format!("{err}");
assert!(msg.contains("not supported"));
}
#[test]
fn duration_unknown_unit_rejected() {
assert!(parse_iso8601_duration("P5X").is_err());
assert!(parse_iso8601_duration("PT5X").is_err());
}
#[test]
fn duration_zero_rejected() {
assert!(parse_iso8601_duration("P0D").is_err());
}
#[test]
fn duration_trailing_digits_without_unit_rejected() {
assert!(parse_iso8601_duration("P5").is_err());
assert!(parse_iso8601_duration("PT12").is_err());
}
#[test]
fn duration_t_separator_with_no_time_rejected() {
assert!(parse_iso8601_duration("P1DT").is_err());
}
#[test]
fn record_action_audit_reason_shape() {
let emitted = vec![
serde_json::json!({"val": "!hide", "uri": "did:plc:subject0000000000000000"}),
serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
];
let json = build_record_action_audit_reason(
42,
ActionType::TempSuspension,
"spam",
&["spam".to_string(), "harassment".to_string()],
4,
2,
true,
&emitted,
"moderator",
None,
None,
);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["action_id"], 42);
assert_eq!(v["action_type"], "temp_suspension");
assert_eq!(v["primary_reason"], "spam");
assert_eq!(v["reason_codes"], serde_json::json!(["spam", "harassment"]));
assert_eq!(v["strike_value_base"], 4);
assert_eq!(v["strike_value_applied"], 2);
assert_eq!(v["was_dampened"], true);
assert_eq!(v["emitted_labels"][0]["val"], "!hide");
assert_eq!(v["emitted_labels"][1]["val"], "reason-spam");
assert_eq!(v["actor_kind"], "moderator");
assert!(v.get("triggered_by_policy_rule").is_none());
assert!(v.get("policy_consequence").is_none());
}
#[test]
fn record_action_audit_reason_empty_emitted_labels() {
let json = build_record_action_audit_reason(
7,
ActionType::Note,
"spam",
&["spam".to_string()],
0,
0,
false,
&[],
"moderator",
None,
None,
);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["emitted_labels"], serde_json::json!([]));
}
#[test]
fn record_action_audit_reason_with_policy_consequence_auto() {
let pc = serde_json::json!({
"rule_fired": "warn_at_5",
"mode": "auto",
"auto_action_id": 43,
});
let json = build_record_action_audit_reason(
42,
ActionType::Warning,
"spam",
&["spam".to_string()],
0,
0,
false,
&[],
"moderator",
None,
Some(pc),
);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["policy_consequence"]["rule_fired"], "warn_at_5");
assert_eq!(v["policy_consequence"]["mode"], "auto");
assert_eq!(v["policy_consequence"]["auto_action_id"], 43);
}
#[test]
fn record_action_audit_reason_for_policy_recorded_action() {
let json = build_record_action_audit_reason(
43,
ActionType::Warning,
"policy-threshold",
&["policy-threshold".to_string()],
0,
0,
false,
&[],
"policy",
Some("warn_at_5"),
None,
);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["actor_kind"], "policy");
assert_eq!(v["triggered_by_policy_rule"], "warn_at_5");
}
#[test]
fn revoke_action_audit_reason_with_reason() {
let json = build_revoke_action_audit_reason(7, Some("appeal granted"), &[]);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["action_id"], 7);
assert_eq!(v["revoked_reason"], "appeal granted");
assert_eq!(v["negated_labels"], serde_json::json!([]));
}
#[test]
fn revoke_action_audit_reason_without_reason_is_null() {
let json = build_revoke_action_audit_reason(7, None, &[]);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert!(v["revoked_reason"].is_null());
}
#[test]
fn revoke_action_audit_reason_with_negated_labels() {
let labels = vec![
serde_json::json!({"val": "!takedown", "uri": "did:plc:subject0000000000000000"}),
serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
];
let json = build_revoke_action_audit_reason(11, None, &labels);
let v: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(v["negated_labels"][0]["val"], "!takedown");
assert_eq!(v["negated_labels"][1]["val"], "reason-spam");
}
#[test]
fn skip_action_label_when_row_already_carries_an_emitted_val() {
assert!(should_skip_action_label_emission(Some("!takedown")));
}
#[test]
fn emit_action_label_when_row_has_no_emitted_val() {
assert!(!should_skip_action_label_emission(None));
}
#[test]
fn skip_reason_emission_when_already_linked() {
let existing: HashSet<&str> = ["spam", "harassment"].into_iter().collect();
assert!(should_skip_reason_emission("spam", &existing));
assert!(should_skip_reason_emission("harassment", &existing));
}
#[test]
fn emit_reason_when_not_already_linked() {
let existing: HashSet<&str> = ["spam"].into_iter().collect();
assert!(!should_skip_reason_emission("hate-speech", &existing));
let empty: HashSet<&str> = HashSet::new();
assert!(!should_skip_reason_emission("anything", &empty));
}
#[test]
fn reason_emission_skip_check_is_case_sensitive() {
let existing: HashSet<&str> = ["SPAM"].into_iter().collect();
assert!(!should_skip_reason_emission("spam", &existing));
}
}