cairn_mod/writer.rs
1//! Single-writer task (§F5 / §10 architecture).
2//!
3//! All mutations to `labels`, `label_sequence`, `audit_log`, and
4//! `server_instance_lease` flow through one task that owns the write pool.
5//! Other components obtain a cheaply-cloneable [`WriterHandle`] and submit
6//! [`ApplyLabelRequest`] / [`NegateLabelRequest`] over an mpsc channel;
7//! replies travel on per-request oneshot channels so the caller may cancel
8//! (drop the future) without corrupting writer state — the transaction
9//! still commits, the reply is discarded.
10//!
11//! Invariants this module enforces, any of which is load-bearing:
12//!
13//! - **Single instance.** On startup [`spawn`] either acquires the lease
14//! row in `server_instance_lease` (id = 1) or returns
15//! [`Error::LeaseHeld`]. A background heartbeat updates `last_heartbeat`
16//! every 10s; a second Cairn started against the same file sees a
17//! <60s-old heartbeat and refuses. The 60s threshold is 6× the heartbeat
18//! interval — tolerant of a single missed tick, strict enough that a
19//! legitimately-orphaned lease (hard crash, unclean shutdown) ages out
20//! before the operator manually restarts.
21//!
22//! - **Monotonic seq.** Sequence values come from
23//! `label_sequence.seq` (AUTOINCREMENT) reserved inside the same
24//! transaction as the `labels` INSERT. Under AUTOINCREMENT, SQLite
25//! never reuses a seq value even across rollbacks.
26//!
27//! - **Monotonic cts.** For each `(src, uri, val)` tuple the writer
28//! clamps to `max(wall_now, prev_cts + 1ms)` (§6.1) using the
29//! `labels_tuple_idx` composite index for the prev-cts lookup.
30//!
31//! - **Audit-per-write.** Every successful label INSERT produces exactly
32//! one `audit_log` row **in the same transaction**. The audit row's
33//! `reason` column carries a JSON object documenting the event
34//! ({"val", "neg", "moderator_reason"}). See [`AUDIT_REASON_SCHEMA`].
35//!
36//! - **Sign-then-insert.** `sign_label` runs between sequence reservation
37//! and the labels INSERT. The signed bytes include `ver` but not `seq`
38//! or `signing_key_id` (§6.2 step 1, §6.1: seq lives on the frame, not
39//! the label).
40
41use std::collections::HashSet;
42use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
43
44use proto_blue_crypto::{K256Keypair, Keypair as _, format_multikey};
45use sqlx::{Pool, Sqlite};
46use time::format_description::FormatItem;
47use time::macros::format_description;
48use time::{OffsetDateTime, PrimitiveDateTime};
49use tokio::sync::{broadcast, mpsc, oneshot, watch};
50use tokio::time::{MissedTickBehavior, interval};
51use uuid::Uuid;
52
53use crate::error::{Error, Result};
54use crate::label::Label;
55use crate::labels::emission::{
56 ActionForEmission, LabelDraft, resolve_action_labels, resolve_reason_labels,
57};
58use crate::moderation::decay::calculate_strike_state;
59use crate::moderation::policy::StrikePolicy;
60use crate::moderation::reasons::ReasonVocabulary;
61use crate::moderation::strike::{
62 StrikeApplication, calculate as strike_calculate, resolve_primary_reason,
63};
64use crate::moderation::types::{ActionRecord, ActionType};
65use crate::moderation::window::compute_position_in_window;
66use crate::server::RetentionConfig;
67use crate::signing::sign_label;
68use crate::signing_key::SigningKey;
69
70/// Lease freshness threshold (§F5). A lease younger than this is held by
71/// a live peer; younger than 10s would be flaky under a single missed
72/// heartbeat, older than ~2× the value risks a genuine zombie blocking a
73/// legitimate restart for too long.
74pub(crate) const LEASE_STALE_MS: i64 = 60_000;
75
76/// Heartbeat interval. 10s × 6 = 60s staleness budget per the threshold.
77const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
78
79/// Granularity at which the writer checks "is it time to start the
80/// scheduled sweep?". Daily fires bucketed by UTC hour means we don't
81/// need finer than this; 60s keeps the check cost negligible (one
82/// `Instant::now()` comparison + an `Option` peek per minute).
83const SWEEP_CHECK_INTERVAL: Duration = Duration::from_secs(60);
84
85/// Upper bound on queued write commands. Moderators rarely drive more
86/// than single-digit req/s even on busy operators; 64 is a comfortable
87/// overshoot that still makes backpressure visible quickly under a bug.
88const COMMAND_BUFFER: usize = 64;
89
90/// Broadcast buffer for [`LabelEvent`] fan-out. Slow subscribers hit
91/// `RecvError::Lagged` past this; the subscribeLabels endpoint (#7) turns
92/// that into a client-connection close.
93const BROADCAST_BUFFER: usize = 1024;
94
95/// Audit-log `reason` JSON schema for `label_applied` / `label_negated`.
96/// Centralized as a doc constant so future actions (signing_key_added,
97/// report_resolved, etc.) have an obvious place to register their shape.
98///
99/// ```json
100/// {
101/// "val": "<label value>",
102/// "neg": true | false,
103/// "moderator_reason": "<free text>" | null
104/// }
105/// ```
106#[doc(alias = "audit_log.reason")]
107pub const AUDIT_REASON_SCHEMA: &str =
108 "label_applied / label_negated: { val, neg, moderator_reason }";
109
110/// Audit-log `reason` JSON schema for `report_resolved` (§F12 resolveReport
111/// atomicity — two rows land in one transaction: the inner label_applied
112/// per [`AUDIT_REASON_SCHEMA`] plus this one).
113///
114/// ```json
115/// {
116/// "applied_label_val": "<val>" | null,
117/// "resolution_reason": "<free text>" | null
118/// }
119/// ```
120#[doc(alias = "audit_log.reason.report_resolved")]
121pub const AUDIT_REASON_RESOLVE_REPORT: &str =
122 "report_resolved: { applied_label_val, resolution_reason }";
123
124/// Audit-log `reason` JSON schema for `reporter_flagged` /
125/// `reporter_unflagged` (§F12 flagReporter).
126///
127/// ```json
128/// {
129/// "did": "<flagged DID>",
130/// "suppressed": true | false,
131/// "moderator_reason": "<free text>" | null
132/// }
133/// ```
134#[doc(alias = "audit_log.reason.flag_reporter")]
135pub const AUDIT_REASON_FLAG_REPORTER: &str =
136 "reporter_flagged / reporter_unflagged: { did, suppressed, moderator_reason }";
137
138/// Audit-log `reason` JSON schema for `retention_sweep` (§F4 — written
139/// only by the operator-initiated admin path; the scheduled-fire
140/// path does NOT audit per Q6/D2). Captures the result of the sweep
141/// run for reconstruction-friendly ops queries.
142///
143/// ```json
144/// {
145/// "rows_deleted": <i64>,
146/// "batches": <u64>,
147/// "duration_ms": <u64>,
148/// "retention_days_applied": <u32> | null
149/// }
150/// ```
151#[doc(alias = "audit_log.reason.retention_sweep")]
152pub const AUDIT_REASON_RETENTION_SWEEP: &str =
153 "retention_sweep: { rows_deleted, batches, duration_ms, retention_days_applied }";
154
155/// Closed set of `audit_log.action` values emitted by Cairn write paths.
156///
157/// §F10: audit rows only for moderation decisions, not for input/operational
158/// events — `createReport` (input) intentionally does NOT audit. The
159/// listAuditLog handler validates its `action` query param against this
160/// exact set; unknown values return `InvalidRequest` rather than silently
161/// matching zero rows.
162///
163/// Adding a new moderation action requires: (1) the handler writes an
164/// `INSERT INTO audit_log (action, ...)` row, (2) the value is added here,
165/// (3) the `reason` schema for the new action is documented as a const
166/// alongside [`AUDIT_REASON_SCHEMA`] and friends.
167pub const AUDIT_ACTION_VALUES: &[&str] = &[
168 "label_applied",
169 "label_negated",
170 "pending_policy_action_confirmed",
171 "pending_policy_action_dismissed",
172 "report_resolved",
173 "reporter_flagged",
174 "reporter_unflagged",
175 "retention_sweep",
176 "subject_action_recorded",
177 "subject_action_revoked",
178];
179
180/// Closed set of `audit_log.outcome` values matching the SQL `CHECK`
181/// constraint in `migrations/0001_init.sql`. listAuditLog validates the
182/// `outcome` query param against this set; the SQL constraint itself is
183/// the durable source of truth — this slice mirrors it so the handler
184/// can reject invalid values pre-query without round-tripping to SQLite.
185pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
186
187/// RFC-3339 with millisecond precision. `Z` is appended by
188/// [`rfc3339_from_epoch_ms`] and stripped by [`parse_rfc3339_ms`] — kept
189/// out of the format description because `time::OffsetDateTime::parse`
190/// can't infer a UTC offset from the literal character `Z`, and
191/// symmetric Z-handling at the string boundary is simpler than switching
192/// parsing-only to `well_known::Rfc3339`.
193const CTS_FORMAT: &[FormatItem<'_>] =
194 format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
195
196/// Moderator request to apply a label (positive event).
197#[derive(Debug, Clone)]
198pub struct ApplyLabelRequest {
199 /// The DID of the moderator issuing the request. Becomes
200 /// `audit_log.actor_did`. `src` on the label itself is the writer's
201 /// service DID, not this.
202 pub actor_did: String,
203 /// AT-URI or DID of the subject being labeled.
204 pub uri: String,
205 /// Optional record-version pin. Present means "this specific version";
206 /// absent means "all versions / the account" per §6.1.
207 pub cid: Option<String>,
208 /// Label value, ≤128 bytes (validated on insert by schema CHECK and
209 /// by caller-side input validation on the XRPC boundary).
210 pub val: String,
211 /// Optional expiration timestamp (RFC-3339 Z). Stored only; expiry
212 /// enforcement is v1.1 (§F7).
213 pub exp: Option<String>,
214 /// Free-text reason recorded to `audit_log` only. Not signed, not
215 /// included in the label record.
216 pub moderator_reason: Option<String>,
217}
218
219/// Moderator request to negate (withdraw) a previously-applied label.
220///
221/// Uniqueness is at the `(src, uri, val)` tuple. The negation copies the
222/// most-recent applied event's `cid` so the negation pins the same record
223/// version — callers don't supply it.
224#[derive(Debug, Clone)]
225pub struct NegateLabelRequest {
226 /// Moderator DID issuing the negation. Becomes
227 /// `audit_log.actor_did`.
228 pub actor_did: String,
229 /// AT-URI or DID of the subject whose label is being withdrawn.
230 /// The `cid` pinning (if any) is copied from the prior apply.
231 pub uri: String,
232 /// Label value being negated. Must match an existing applied
233 /// label on `(src, uri, val)` or the call errors with
234 /// `LabelNotFound`.
235 pub val: String,
236 /// Free-text reason recorded to `audit_log` only.
237 pub moderator_reason: Option<String>,
238}
239
240/// A committed label event. Returned by `apply_label` / `negate_label` and
241/// broadcast to subscribers (subscribeLabels consumers in #7). Wire
242/// serialization is that consumer's concern — the LabelEvent itself
243/// carries no serde derive because the canonical encoding for the wire is
244/// the DAG-CBOR path in `crate::signing`, not serde JSON.
245#[derive(Debug, Clone)]
246pub struct LabelEvent {
247 /// Frame sequence number from `label_sequence`. Strictly monotonic.
248 pub seq: i64,
249 /// The full signed label (`sig` populated).
250 pub label: Label,
251}
252
253/// Internal write command. One variant per public method.
254enum WriteCommand {
255 Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
256 Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
257 ResolveReport(
258 ResolveReportRequest,
259 oneshot::Sender<Result<ResolvedReport>>,
260 ),
261 Sweep(SweepRequest, oneshot::Sender<Result<SweepBatchResult>>),
262 /// Append a hash-chained audit row (#39). Used by in-process
263 /// callers that don't have an existing transaction (e.g.,
264 /// `retentionSweep` after the sweep itself completes). Callers
265 /// that already hold a transaction (writer-internal handlers,
266 /// `flag_reporter`) call [`crate::audit::append::append_in_tx`]
267 /// directly within that transaction instead.
268 AppendAudit(
269 crate::audit::append::AuditRowForAppend,
270 oneshot::Sender<Result<i64>>,
271 ),
272 /// Record a graduated-action subject_actions row (§F20 / #51).
273 /// Single-transaction: input validation, strike calculation,
274 /// subject_actions INSERT, strike_state UPSERT, audit_log
275 /// hash-chain append.
276 RecordAction(RecordActionRequest, oneshot::Sender<Result<RecordedAction>>),
277 /// Revoke a previously-recorded subject_actions row (§F20 /
278 /// #51). Single-transaction: row lookup + state-check,
279 /// revoked_at UPDATE, strike_state recompute, audit_log
280 /// hash-chain append.
281 RevokeAction(RevokeActionRequest, oneshot::Sender<Result<RevokedAction>>),
282 /// Confirm a pending policy action (§F22 / #74). Promotes the
283 /// pending row's proposed action to a real `subject_actions`
284 /// row (actor_kind='moderator', triggered_by_policy_rule
285 /// preserves the rule name as forensic provenance), emits
286 /// labels, UPDATEs the pending row's resolution columns, and
287 /// recomputes strike state — all hash-chained under one
288 /// `pending_policy_action_confirmed` audit row.
289 ConfirmPendingAction(
290 ConfirmPendingActionRequest,
291 oneshot::Sender<Result<ConfirmedPendingAction>>,
292 ),
293 /// Dismiss a pending policy action (§F22 / #75). Single-
294 /// transaction: pending row load + state-check, UPDATE to
295 /// resolution='dismissed', `pending_policy_action_dismissed`
296 /// audit row. No subject_actions row, no label emission, no
297 /// strike-state change — the moderator is explicitly closing
298 /// the loop on what the policy engine flagged.
299 DismissPendingAction(
300 DismissPendingActionRequest,
301 oneshot::Sender<Result<DismissedPendingAction>>,
302 ),
303 Shutdown(oneshot::Sender<Result<()>>),
304}
305
306/// Inline label-application sub-object for [`ResolveReportRequest`].
307/// Structurally aligned with [`ApplyLabelRequest`] minus the
308/// moderator-reason (that field on the outer request already captures
309/// the resolution rationale; the label's own audit reason is derived
310/// from the resolution flow).
311#[derive(Debug, Clone)]
312pub struct ApplyLabelInline {
313 /// Subject the label applies to (AT-URI or DID).
314 pub uri: String,
315 /// Optional record-version pin; `None` targets the account / all
316 /// versions of the record.
317 pub cid: Option<String>,
318 /// Label value to apply (≤128 bytes).
319 pub val: String,
320 /// Optional expiration timestamp (RFC-3339 Z). Stored only;
321 /// enforcement is v1.1.
322 pub exp: Option<String>,
323}
324
325/// Resolution outcome (#27): the implicit "with-label vs without-label"
326/// semantic, made explicit. The wire shape is `applyLabel: Option<…>`
327/// per the `resolveReport` lexicon; this enum is internal and the
328/// handler maps `None → Dismiss` / `Some(_) → ApplyLabel(_)` at the
329/// HTTP boundary.
330#[derive(Debug, Clone)]
331pub enum ResolutionAction {
332 /// Resolve without emitting a label (the operator-UX "dismiss"
333 /// flow). `resolution_label` on the row stays NULL.
334 Dismiss,
335 /// Resolve and emit a label in the same transaction. The label's
336 /// `val` is also recorded as `resolution_label` on the report row.
337 ApplyLabel(ApplyLabelInline),
338}
339
340impl ResolutionAction {
341 /// Borrow the inner `ApplyLabelInline` if this is an `ApplyLabel`
342 /// variant. Convenience for the writer's UPDATE / audit code that
343 /// needs both the optional label *and* its `val` projection.
344 pub fn as_apply(&self) -> Option<&ApplyLabelInline> {
345 match self {
346 ResolutionAction::Dismiss => None,
347 ResolutionAction::ApplyLabel(a) => Some(a),
348 }
349 }
350}
351
352/// Request to resolve a report (§F12 `resolveReport`). The optional
353/// label is applied **in the same transaction** as the report status
354/// update and both audit rows — §F5 single-writer invariant plus §F12
355/// atomicity requirement documented in the `resolveReport` lexicon.
356#[derive(Debug, Clone)]
357pub struct ResolveReportRequest {
358 /// Moderator DID issuing the resolution. Becomes
359 /// `audit_log.actor_did` on both the label-applied (if any) and
360 /// report_resolved rows.
361 pub actor_did: String,
362 /// Primary key of the report being resolved.
363 pub report_id: i64,
364 /// Whether the resolution emits a label or just closes the report.
365 /// Replaces the pre-#27 `apply_label: Option<ApplyLabelInline>`
366 /// representation; semantically identical, named for the operator
367 /// UX (dismiss vs apply-label).
368 pub action: ResolutionAction,
369 /// Free-text resolution rationale recorded to audit_log.
370 pub resolution_reason: Option<String>,
371}
372
373/// Result of a successful resolve. `label_event` is `Some(..)` iff
374/// the request carried `apply_label`; the broadcast has already
375/// happened inside the writer task post-commit.
376#[derive(Debug, Clone)]
377pub struct ResolvedReport {
378 /// The updated report row (status now `resolved`).
379 pub report: crate::report::Report,
380 /// The emitted label event, if the resolution included an
381 /// `apply_label`. Signed, broadcast, and committed as part of
382 /// the same transaction as the report UPDATE.
383 pub label_event: Option<LabelEvent>,
384}
385
386/// Trigger for the retention sweep (§F4). Carries no per-call
387/// parameters today — the cutoff comes from `retention_days` baked
388/// into the writer at spawn time, not from the request — but is
389/// kept as a typed unit so a future "sweep with explicit override
390/// for this one run" remains a non-breaking change to the variant.
391#[derive(Debug, Clone, Default)]
392pub struct SweepRequest;
393
394/// Per-batch outcome from one sweep dispatch through the writer's
395/// internal `WriteCommand::Sweep` channel. The writer task processes
396/// ONE batch per command so its main `select!` can interleave
397/// incoming label writes between batches (§F4 + §F5 — single-writer
398/// invariant + bounded latency). Callers who want a full sweep loop
399/// until [`Self::has_more`] is `false`; the [`WriterHandle::sweep`]
400/// convenience wrapper does this internally.
401#[derive(Debug, Clone)]
402pub struct SweepBatchResult {
403 /// Rows deleted in this batch. Zero means "no more old rows
404 /// match the cutoff" — caller stops looping.
405 pub rows_deleted: i64,
406 /// `true` when the batch hit the configured `sweep_batch_size`
407 /// limit and a follow-up batch may find more rows. `false`
408 /// indicates the batch was partial (last batch) and the sweep
409 /// is complete.
410 pub has_more: bool,
411 /// Cutoff days actually applied. `None` when the writer was
412 /// spawned with `retention_days = None` — the sweep is a no-op
413 /// in that configuration and `rows_deleted` is always 0.
414 pub retention_days_applied: Option<u32>,
415}
416
417/// Request to record a new subject_actions row (§F20 / #51 graduated-
418/// action moderation). The handler validates the inputs, computes
419/// the strike value at action time via the v1.4 calculators
420/// (#48/#49/#50/#51), inserts the row, updates the strike-state
421/// cache, and writes a hash-chained audit_log row — all in a
422/// single transaction per §F5 atomicity.
423#[derive(Debug, Clone)]
424pub struct RecordActionRequest {
425 /// Raw subject — DID (`did:plc:...`, `did:web:...`) or AT-URI
426 /// (`at://did:.../col/r`). The handler routes to subject_did vs
427 /// subject_uri based on prefix; for AT-URIs the parent repo DID
428 /// is extracted as subject_did and strike accounting rolls up
429 /// to the account.
430 pub subject: String,
431 /// DID of the moderator/admin recording the action (JWT iss
432 /// from XRPC; CLI caller DID otherwise). Becomes
433 /// subject_actions.actor_did and audit_log.actor_did.
434 pub actor_did: String,
435 /// Graduated-action category. The handler enforces:
436 /// `temp_suspension` requires `duration_iso`; everything else
437 /// rejects it.
438 pub action_type: ActionType,
439 /// Operator-vocabulary identifiers from `[moderation_reasons]`.
440 /// Must be non-empty. Multi-reason resolution: severe wins; else
441 /// highest base_weight wins; ties → first-listed.
442 pub reason_codes: Vec<String>,
443 /// ISO-8601 duration string (e.g. `P7D`). Required for
444 /// `temp_suspension`; rejected for other types. Stored verbatim
445 /// on the row for display; `expires_at` is the canonical
446 /// "when does it end" surface.
447 pub duration_iso: Option<String>,
448 /// Optional moderator-facing rationale.
449 pub notes: Option<String>,
450 /// Optional report row ids that motivated this action.
451 pub report_ids: Vec<i64>,
452}
453
454/// Result of a successful [`WriterHandle::record_action`]. The
455/// admin handler echoes these fields verbatim in the
456/// `tools.cairn.admin.recordAction` response; the CLI surfaces them
457/// in human and JSON output.
458#[derive(Debug, Clone)]
459pub struct RecordedAction {
460 /// Inserted subject_actions row id.
461 pub action_id: i64,
462 /// Reason's `base_weight` before dampening. `0` for warning/note.
463 pub strike_value_base: u32,
464 /// Strike weight actually applied after dampening.
465 pub strike_value_applied: u32,
466 /// `true` iff the dampening curve was consulted (see #49).
467 pub was_dampened: bool,
468 /// Subject's `current_strike_count` BEFORE this action — frozen
469 /// for forensic history.
470 pub strikes_at_time_of_action: u32,
471}
472
473/// Request to revoke a previously-recorded subject_actions row
474/// (§F20 / #51). Sets the row's revoked_at/revoked_by_did/
475/// revoked_reason columns (the schema's no-update-except-revoke
476/// trigger permits exactly this NULL→non-NULL transition);
477/// recomputes the strike-state cache; writes a hash-chained
478/// audit_log row. All in a single transaction.
479#[derive(Debug, Clone)]
480pub struct RevokeActionRequest {
481 /// subject_actions.id to revoke. Errors with
482 /// [`Error::ActionNotFound`] when the row doesn't exist;
483 /// [`Error::ActionAlreadyRevoked`] when revoked_at is already
484 /// non-NULL.
485 pub action_id: i64,
486 /// DID of the moderator/admin performing the revocation.
487 pub revoked_by_did: String,
488 /// Optional rationale stored on the row's revoked_reason column.
489 pub revoked_reason: Option<String>,
490}
491
492/// Result of a successful [`WriterHandle::revoke_action`].
493#[derive(Debug, Clone)]
494pub struct RevokedAction {
495 /// The revoked row's id (echoed for confirmation).
496 pub action_id: i64,
497 /// Wall-clock the revocation took effect, as RFC-3339 Z
498 /// (matches the [`AUDIT_REASON_RECORD_ACTION`] surface for
499 /// consumer convenience).
500 pub revoked_at: String,
501}
502
503/// Request to confirm a pending policy action (§F22 / #74). The
504/// pending row's proposed action is "promoted" to a real
505/// `subject_actions` row (actor_kind='moderator'; the moderator
506/// takes responsibility by confirming) with `triggered_by_policy_rule`
507/// preserved as forensic provenance. Single-transaction: pending
508/// load + state-check, subject_actions INSERT, label emission,
509/// pending UPDATE (resolution='confirmed'), strike-state recompute,
510/// audit_log hash-chain append.
511#[derive(Debug, Clone)]
512pub struct ConfirmPendingActionRequest {
513 /// `pending_policy_actions.id` to confirm. Errors with
514 /// [`Error::PendingActionNotFound`] when the row doesn't
515 /// exist; [`Error::PendingAlreadyResolved`] when the
516 /// resolution column is already non-NULL.
517 pub pending_id: i64,
518 /// DID of the moderator confirming the pending. Becomes the
519 /// new subject_actions row's `actor_did` and the audit row's
520 /// `actor_did`, and lands on `pending_policy_actions.resolved_by_did`.
521 pub moderator_did: String,
522 /// Optional moderator-facing rationale. Stored on the new
523 /// subject_actions row's `notes` column and echoed in the audit
524 /// row's reason JSON as `moderator_note` for forensic
525 /// reconstruction.
526 pub note: Option<String>,
527}
528
529/// Result of a successful [`WriterHandle::confirm_pending_action`].
530#[derive(Debug, Clone)]
531pub struct ConfirmedPendingAction {
532 /// Inserted subject_actions row id (the materialized action).
533 pub action_id: i64,
534 /// The pending row that was just resolved (echoed for
535 /// confirmation).
536 pub pending_id: i64,
537 /// Wall-clock the confirmation took effect, as RFC-3339 Z.
538 pub resolved_at: String,
539}
540
541/// Request to dismiss a pending policy action (§F22 / #75). The
542/// pending row stays in the table as forensic record with
543/// `resolution = 'dismissed'`; no subject_actions row, no label
544/// emission, no strike-state change. Single-transaction: pending
545/// load + state-check, UPDATE, audit_log hash-chain append.
546///
547/// Note that there is no `SubjectTakendown` defensive check here
548/// (unlike confirm in #74): explicit dismissal is meaningful
549/// regardless of takedown state, and is in fact part of the
550/// cleanup path #76 will automate when a takedown lands.
551#[derive(Debug, Clone)]
552pub struct DismissPendingActionRequest {
553 /// `pending_policy_actions.id` to dismiss. Errors with
554 /// [`Error::PendingActionNotFound`] when the row doesn't
555 /// exist; [`Error::PendingAlreadyResolved`] when the
556 /// resolution column is already non-NULL (already confirmed
557 /// or dismissed).
558 pub pending_id: i64,
559 /// DID of the moderator dismissing the pending. Lands on
560 /// `pending_policy_actions.resolved_by_did` and the audit
561 /// row's `actor_did`.
562 pub moderator_did: String,
563 /// Optional moderator-facing rationale. Captured in the
564 /// audit row's reason JSON as `moderator_reason` (option
565 /// (b) per the design — the pending table tracks resolution
566 /// state; rationale lives in audit). The pending row itself
567 /// has no `resolved_reason` column.
568 pub reason: Option<String>,
569}
570
571/// Result of a successful [`WriterHandle::dismiss_pending_action`].
572#[derive(Debug, Clone)]
573pub struct DismissedPendingAction {
574 /// The pending row that was just resolved (echoed for
575 /// confirmation).
576 pub pending_id: i64,
577 /// Wall-clock the dismissal took effect, as RFC-3339 Z.
578 pub resolved_at: String,
579}
580
581/// Aggregate result of a full sweep run (returned by
582/// [`WriterHandle::sweep`] after looping over batches).
583#[derive(Debug, Clone)]
584pub struct SweepResult {
585 /// Total rows deleted across all batches.
586 pub rows_deleted: i64,
587 /// Number of batches issued.
588 pub batches: u64,
589 /// Wall-clock duration of the full sweep, in milliseconds.
590 pub duration_ms: u64,
591 /// Cutoff days actually applied, or `None` when the writer's
592 /// `retention_days` is `None` (sweep is a no-op).
593 pub retention_days_applied: Option<u32>,
594}
595
596/// Cheap handle to the writer task. Clones share the same underlying
597/// mpsc channel; a drop of the last clone is a silent shutdown signal
598/// to the writer task (receiver closes). For a clean shutdown that
599/// releases the lease row, call [`WriterHandle::shutdown`].
600#[derive(Debug, Clone)]
601pub struct WriterHandle {
602 tx: mpsc::Sender<WriteCommand>,
603 broadcast_tx: broadcast::Sender<LabelEvent>,
604 /// Flipped to `true` when the writer task starts its shutdown path.
605 /// Exposed via [`WriterHandle::shutdown_signal`] so peer components
606 /// (the subscribeLabels handler, future maintenance tasks) can close
607 /// their own resources cleanly. Using a watch channel rather than the
608 /// broadcast-channel close signal because clones of `WriterHandle`
609 /// keep `broadcast_tx` alive — receivers would otherwise never see
610 /// `RecvError::Closed`.
611 shutdown_rx: watch::Receiver<bool>,
612}
613
614impl WriterHandle {
615 /// Submit an apply-label request. Resolves when the writer has
616 /// committed the transaction and broadcast the event, or returns an
617 /// `Err` describing why the write was rejected.
618 pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
619 let (reply_tx, reply_rx) = oneshot::channel();
620 self.tx
621 .send(WriteCommand::Apply(req, reply_tx))
622 .await
623 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
624 reply_rx
625 .await
626 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
627 }
628
629 /// Submit a negate-label request. Returns [`Error::LabelNotFound`] if
630 /// no applied label currently exists for `(service_did, uri, val)`.
631 pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
632 let (reply_tx, reply_rx) = oneshot::channel();
633 self.tx
634 .send(WriteCommand::Negate(req, reply_tx))
635 .await
636 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
637 reply_rx
638 .await
639 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
640 }
641
642 /// Resolve a report, optionally applying a label in the same
643 /// atomic transaction. The label INSERT + audit row + report
644 /// UPDATE + resolution audit row all commit together or not at
645 /// all; the label broadcast fires post-commit inside the writer
646 /// task (§F5 + §F12 atomicity contract).
647 ///
648 /// Errors:
649 /// - [`Error::ReportNotFound`] — `report_id` doesn't exist.
650 /// - [`Error::ReportAlreadyResolved`] — report is not in the
651 /// `pending` state. Handler maps to a generic `InvalidRequest`.
652 pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
653 let (reply_tx, reply_rx) = oneshot::channel();
654 self.tx
655 .send(WriteCommand::ResolveReport(req, reply_tx))
656 .await
657 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
658 reply_rx
659 .await
660 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
661 }
662
663 /// Trigger a §F4 retention sweep through the writer task. Loops
664 /// over single-batch dispatches until a batch returns
665 /// `has_more = false`, aggregating rows + duration.
666 /// Other writer commands interleave between batches (single-writer
667 /// invariant + bounded latency).
668 ///
669 /// Returns `rows_deleted = 0, batches = 0` when the writer was
670 /// spawned with `retention_days = None` (sweep configured off):
671 /// the first batch returns immediately and the loop exits with
672 /// the no-op result.
673 pub async fn sweep(&self, _req: SweepRequest) -> Result<SweepResult> {
674 let start = std::time::Instant::now();
675 let mut total_rows: i64 = 0;
676 let mut batches: u64 = 0;
677
678 loop {
679 let (reply_tx, reply_rx) = oneshot::channel();
680 self.tx
681 .send(WriteCommand::Sweep(SweepRequest, reply_tx))
682 .await
683 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
684 let batch = reply_rx
685 .await
686 .map_err(|_| Error::Signing("writer dropped reply channel".into()))??;
687
688 total_rows += batch.rows_deleted;
689 // Count the no-op early-exit batch too — it represents
690 // the round-trip the caller paid for. Useful for tracing.
691 batches += 1;
692
693 if !batch.has_more {
694 return Ok(SweepResult {
695 rows_deleted: total_rows,
696 batches,
697 duration_ms: start.elapsed().as_millis() as u64,
698 retention_days_applied: batch.retention_days_applied,
699 });
700 }
701 }
702 }
703
704 /// Record a graduated-action moderation event (§F20 / #51).
705 /// The writer validates inputs, computes the strike value at
706 /// action time via the v1.4 calculators (#48/#49/#50/#51),
707 /// inserts the subject_actions row, updates the strike-state
708 /// cache, and writes a hash-chained audit_log row — all in a
709 /// single transaction. Errors:
710 ///
711 /// - [`Error::ReasonNotFound`] when a reason identifier is not
712 /// declared in `[moderation_reasons]`.
713 /// - [`Error::DurationRequiredForTempSuspension`] for
714 /// `temp_suspension` without `duration_iso`.
715 /// - [`Error::DurationOnlyForTempSuspension`] for non-temp
716 /// types with `duration_iso` set.
717 /// - [`Error::Signing`] for malformed subject / duration / etc.
718 pub async fn record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
719 let (reply_tx, reply_rx) = oneshot::channel();
720 self.tx
721 .send(WriteCommand::RecordAction(req, reply_tx))
722 .await
723 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
724 reply_rx
725 .await
726 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
727 }
728
729 /// Revoke a previously-recorded action (§F20 / #51). Errors:
730 ///
731 /// - [`Error::ActionNotFound`] when no row matches `action_id`.
732 /// - [`Error::ActionAlreadyRevoked`] when the row's revoked_at
733 /// is already non-NULL — the schema trigger forbids
734 /// re-revocation; the handler catches it before the UPDATE.
735 pub async fn revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
736 let (reply_tx, reply_rx) = oneshot::channel();
737 self.tx
738 .send(WriteCommand::RevokeAction(req, reply_tx))
739 .await
740 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
741 reply_rx
742 .await
743 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
744 }
745
746 /// Confirm a pending policy action (§F22 / #74). The proposed
747 /// action carried on the pending row is materialized as a real
748 /// `subject_actions` row (actor_kind='moderator',
749 /// triggered_by_policy_rule preserves the rule that proposed
750 /// it), labels emit per [`crate::labels::policy::LabelEmissionPolicy`],
751 /// the pending row's resolution columns transition NULL →
752 /// 'confirmed', and the strike-state cache is recomputed — all
753 /// in one transaction under a single
754 /// `pending_policy_action_confirmed` audit row. Errors:
755 ///
756 /// - [`Error::PendingActionNotFound`] when no row matches
757 /// `pending_id`.
758 /// - [`Error::PendingAlreadyResolved`] when the row's
759 /// resolution column is already non-NULL.
760 /// - [`Error::SubjectTakendown`] when the subject already
761 /// carries an unrevoked Takedown — defensive guard against
762 /// the auto-dismissal-on-takedown race (#76).
763 pub async fn confirm_pending_action(
764 &self,
765 req: ConfirmPendingActionRequest,
766 ) -> Result<ConfirmedPendingAction> {
767 let (reply_tx, reply_rx) = oneshot::channel();
768 self.tx
769 .send(WriteCommand::ConfirmPendingAction(req, reply_tx))
770 .await
771 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
772 reply_rx
773 .await
774 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
775 }
776
777 /// Dismiss a pending policy action (§F22 / #75). The pending
778 /// row's resolution columns transition NULL → 'dismissed' and
779 /// a single hash-chained `pending_policy_action_dismissed`
780 /// audit row commits with the UPDATE — no subject_actions
781 /// row, no label emission, no strike-state change. The
782 /// pending row stays in the table as forensic record
783 /// (confirmed_action_id stays NULL since no action was
784 /// created).
785 ///
786 /// Errors:
787 ///
788 /// - [`Error::PendingActionNotFound`] when no row matches
789 /// `pending_id`.
790 /// - [`Error::PendingAlreadyResolved`] when the row's
791 /// resolution column is already non-NULL.
792 pub async fn dismiss_pending_action(
793 &self,
794 req: DismissPendingActionRequest,
795 ) -> Result<DismissedPendingAction> {
796 let (reply_tx, reply_rx) = oneshot::channel();
797 self.tx
798 .send(WriteCommand::DismissPendingAction(req, reply_tx))
799 .await
800 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
801 reply_rx
802 .await
803 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
804 }
805
806 /// Append a hash-chained audit row through the writer task (#39).
807 /// Used by in-process callers that don't already hold a transaction
808 /// — e.g., `retentionSweep`'s post-sweep audit row. Callers that
809 /// have an open transaction (writer-internal handlers,
810 /// `flag_reporter`) use [`crate::audit::append::append_in_tx`]
811 /// directly so the audit row commits atomically with the rest of
812 /// their work. Cross-process CLIs (publish/unpublish-service-
813 /// record) use [`crate::audit::append::append_via_pool`].
814 ///
815 /// Returns the inserted `audit_log.id`.
816 pub async fn append_audit(&self, row: crate::audit::append::AuditRowForAppend) -> Result<i64> {
817 let (reply_tx, reply_rx) = oneshot::channel();
818 self.tx
819 .send(WriteCommand::AppendAudit(row, reply_tx))
820 .await
821 .map_err(|_| Error::Signing("writer task is shut down".into()))?;
822 reply_rx
823 .await
824 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
825 }
826
827 /// Subscribe to committed events. The returned receiver lags past
828 /// the internal broadcast buffer; the consumer (subscribeLabels, #7)
829 /// turns `RecvError::Lagged` into a connection close.
830 pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
831 self.broadcast_tx.subscribe()
832 }
833
834 /// Count of live broadcast receivers. Exposed for test sync points
835 /// (the WS handler subscribes asynchronously after upgrade; tests
836 /// that want to emit an event "after the subscriber is ready" poll
837 /// this until it reaches the expected count). Not a production hot
838 /// path — `broadcast::Sender::receiver_count` walks an atomic chain.
839 pub fn receiver_count(&self) -> usize {
840 self.broadcast_tx.receiver_count()
841 }
842
843 /// Observe writer lifecycle. The returned receiver starts at `false`
844 /// and flips to `true` once exactly, when the writer has accepted a
845 /// shutdown command and is about to release the lease. Holders close
846 /// downstream connections gracefully on the `true` transition.
847 pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
848 self.shutdown_rx.clone()
849 }
850
851 /// Explicit shutdown. Drains in-flight writes, releases the lease
852 /// row, stops the heartbeat, and returns. Idempotent-ish: calling on
853 /// an already-shut-down writer returns a channel-closed error, not a
854 /// panic.
855 pub async fn shutdown(&self) -> Result<()> {
856 let (reply_tx, reply_rx) = oneshot::channel();
857 self.tx
858 .send(WriteCommand::Shutdown(reply_tx))
859 .await
860 .map_err(|_| Error::Signing("writer task is already shut down".into()))?;
861 reply_rx
862 .await
863 .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
864 }
865}
866
867/// Start the writer task. Acquires the instance lease, bootstraps the
868/// `signing_keys` row if empty, and spawns the event loop + heartbeat.
869///
870/// `service_did` stamps `labels.src` on every emitted event. `key` must
871/// be the private key whose public form is recorded (or will be recorded)
872/// in `signing_keys`.
873///
874/// `retention_days` is the §F4 retention cutoff in days. `None` disables
875/// the retention sweep entirely (`WriteCommand::Sweep` becomes a no-op
876/// returning `rows_deleted = 0`); `Some(N)` lets the sweep delete labels
877/// whose `created_at` is older than `now - N days`. Source of truth is
878/// [`crate::SubscribeConfig::retention_days`] — pass it through verbatim.
879///
880/// `retention` is the sweep execution policy (schedule + batching);
881/// distinct from `retention_days` per the [§F4 split](crate::RetentionConfig).
882///
883/// `reason_vocabulary` and `strike_policy` are the resolved v1.4
884/// moderation surface (#47 / #48). The recorder (#51 RecordAction
885/// handler) consults them at action time to validate reason codes
886/// and compute strike values. They're held by the writer task so
887/// the recorder doesn't have to thread them through every call.
888/// Operators restart `cairn serve` to change either — same posture
889/// as `[labeler]` config.
890///
891/// `label_emission_policy` is the resolved `[label_emission]`
892/// surface (#58, v1.5). The recorder (#60) consults it post-INSERT
893/// to translate the freshly-recorded action into ATProto labels
894/// emitted in the same transaction. Held by the writer task for
895/// the same reason as the v1.4 surfaces above.
896#[allow(clippy::too_many_arguments)]
897pub async fn spawn(
898 pool: Pool<Sqlite>,
899 key: SigningKey,
900 service_did: String,
901 retention_days: Option<u32>,
902 retention: RetentionConfig,
903 reason_vocabulary: ReasonVocabulary,
904 strike_policy: StrikePolicy,
905 label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
906 policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
907) -> Result<WriterHandle> {
908 spawn_with_pds_admin(
909 pool,
910 key,
911 service_did,
912 retention_days,
913 retention,
914 reason_vocabulary,
915 strike_policy,
916 label_emission_policy,
917 policy_automation_policy,
918 None,
919 )
920 .await
921}
922
923/// Spawn the writer task with an optional PDS-admin bridge
924/// attached (#87 / §F23).
925///
926/// Behaves identically to [`spawn`] when `pds_admin` is `None`.
927/// When `Some(bridge)`, every successful recordAction's
928/// post-commit hook consults
929/// [`crate::pds_admin::dispatch::dispatch_after_record_action`]
930/// to dispatch a backend call against the configured PDS and
931/// audit-record the result.
932///
933/// Lives as a sibling function rather than extending [`spawn`]'s
934/// positional signature so the 18+ existing test files calling
935/// `spawn_writer` (the lib.rs re-export of [`spawn`]) don't have
936/// to thread a `None` argument through. `serve::run` calls this
937/// directly with a constructed bridge when `[pds_admin].enabled`
938/// is true.
939#[allow(clippy::too_many_arguments)]
940pub async fn spawn_with_pds_admin(
941 pool: Pool<Sqlite>,
942 key: SigningKey,
943 service_did: String,
944 retention_days: Option<u32>,
945 retention: RetentionConfig,
946 reason_vocabulary: ReasonVocabulary,
947 strike_policy: StrikePolicy,
948 label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
949 policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
950 pds_admin: Option<crate::pds_admin::dispatch::PdsAdminBridge>,
951) -> Result<WriterHandle> {
952 let instance_id = acquire_lease(&pool).await?;
953 let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
954
955 let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
956 let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
957 let (shutdown_tx, shutdown_rx) = watch::channel(false);
958
959 let writer = Writer {
960 pool,
961 key,
962 service_did,
963 signing_key_id,
964 instance_id,
965 rx,
966 broadcast_tx: broadcast_tx.clone(),
967 shutdown_tx,
968 retention_days,
969 retention,
970 reason_vocabulary,
971 strike_policy,
972 label_emission_policy,
973 policy_automation_policy,
974 pds_admin,
975 };
976
977 tokio::spawn(writer.run());
978
979 Ok(WriterHandle {
980 tx,
981 broadcast_tx,
982 shutdown_rx,
983 })
984}
985
986// ---------- Writer task ----------
987
988struct Writer {
989 pool: Pool<Sqlite>,
990 key: SigningKey,
991 service_did: String,
992 signing_key_id: i64,
993 instance_id: String,
994 rx: mpsc::Receiver<WriteCommand>,
995 broadcast_tx: broadcast::Sender<LabelEvent>,
996 shutdown_tx: watch::Sender<bool>,
997 /// Retention cutoff in days. `None` makes the sweep a no-op.
998 /// Source of truth is [`crate::SubscribeConfig::retention_days`].
999 retention_days: Option<u32>,
1000 /// Sweep execution policy (schedule + batching).
1001 retention: RetentionConfig,
1002 /// Resolved [moderation_reasons] vocabulary (#47). Consulted by
1003 /// the RecordAction handler to validate reason codes and look
1004 /// up base_weight / severe.
1005 reason_vocabulary: ReasonVocabulary,
1006 /// Resolved [strike_policy] (#48). Consulted by the RecordAction
1007 /// handler for the dampening curve, threshold, and decay
1008 /// parameters.
1009 strike_policy: StrikePolicy,
1010 /// Resolved [label_emission] (#58, v1.5). Consulted by the
1011 /// RecordAction handler post-INSERT to translate the action into
1012 /// ATProto labels, emitted in the same tx.
1013 label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
1014 /// Resolved [policy_automation] (#71, v1.6). Consulted by the
1015 /// RecordAction handler post-emission to evaluate threshold-
1016 /// crossing rules and either auto-record a consequent action
1017 /// or queue a `pending_policy_actions` row for moderator
1018 /// review (#73).
1019 policy_automation_policy: crate::policy::automation::PolicyAutomationPolicy,
1020 /// Resolved §F23 PDS-admin bridge (#83 policy + #86
1021 /// `OzoneBackend` instance). `None` when `[pds_admin]` is
1022 /// disabled — the post-commit dispatch in
1023 /// [`Self::handle_record_action`] short-circuits without
1024 /// touching the trait. `Some(_)` when enabled — every
1025 /// successfully-committed recordAction passes through
1026 /// [`crate::pds_admin::dispatch::dispatch_after_record_action`]
1027 /// per §A13.
1028 pds_admin: Option<crate::pds_admin::dispatch::PdsAdminBridge>,
1029}
1030
1031/// Internal accumulator for an in-flight scheduled sweep. Lives only
1032/// while [`Writer::run`] is mid-sweep; absent between sweeps.
1033struct SweepRunState {
1034 started_at: Instant,
1035 rows: i64,
1036 batches: u64,
1037}
1038
1039impl Writer {
1040 /// Compute the next absolute [`Instant`] at which the scheduled
1041 /// sweep should fire, or `None` when the sweep is disabled.
1042 ///
1043 /// "Disabled" = `sweep_enabled = false` OR `retention_days = None`.
1044 /// In the second case the sweep would be a no-op anyway; skipping
1045 /// the timer entirely keeps the run loop quiet.
1046 ///
1047 /// Targets the next occurrence of `sweep_run_at_utc_hour:00:00`
1048 /// UTC. If we're already past that hour today, target tomorrow at
1049 /// the same hour.
1050 fn compute_next_sweep_fire(&self) -> Option<Instant> {
1051 if !self.retention.sweep_enabled || self.retention_days.is_none() {
1052 return None;
1053 }
1054 let now_utc = OffsetDateTime::now_utc();
1055 let target_hour = self.retention.sweep_run_at_utc_hour;
1056 let target_today_time = time::Time::from_hms(target_hour, 0, 0)
1057 .expect("hour validated < 24 by Config::validate");
1058 let target_today = now_utc.replace_time(target_today_time);
1059 let target_dt = if target_today <= now_utc {
1060 target_today + time::Duration::days(1)
1061 } else {
1062 target_today
1063 };
1064 let wait = target_dt - now_utc;
1065 let wait_secs = wait.whole_seconds().max(0) as u64;
1066 Some(Instant::now() + Duration::from_secs(wait_secs))
1067 }
1068
1069 async fn run(mut self) {
1070 let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
1071 heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
1072 // First tick fires immediately; skip it — we just acquired the lease
1073 // with a fresh timestamp, so touching it again is redundant.
1074 heartbeat_timer.tick().await;
1075
1076 // Scheduled-sweep wiring (§F4). The check timer wakes the run
1077 // loop once a minute to evaluate "is it time to start a sweep?";
1078 // when a sweep starts, `sweep_state` is populated and the
1079 // immediate-batch arm runs one batch per loop iteration until
1080 // `has_more = false`. Inter-batch yielding to incoming commands
1081 // is automatic — the biased select prefers `rx.recv()` and
1082 // `heartbeat_timer.tick()` over running another batch.
1083 let mut sweep_check_timer = interval(SWEEP_CHECK_INTERVAL);
1084 sweep_check_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
1085 sweep_check_timer.tick().await;
1086
1087 let mut next_scheduled_fire: Option<Instant> = self.compute_next_sweep_fire();
1088 let mut sweep_state: Option<SweepRunState> = None;
1089
1090 loop {
1091 let sweep_in_progress = sweep_state.is_some();
1092 tokio::select! {
1093 biased;
1094 // Prefer draining write commands over heartbeats so a
1095 // steady-state of inbound work does not starve itself
1096 // behind a background housekeeping task.
1097 cmd = self.rx.recv() => {
1098 match cmd {
1099 Some(WriteCommand::Apply(req, reply)) => {
1100 let res = self.handle_apply(req).await;
1101 // Caller may have cancelled; dropping the reply is fine.
1102 let _ = reply.send(res);
1103 }
1104 Some(WriteCommand::Negate(req, reply)) => {
1105 let res = self.handle_negate(req).await;
1106 let _ = reply.send(res);
1107 }
1108 Some(WriteCommand::ResolveReport(req, reply)) => {
1109 let res = self.handle_resolve_report(req).await;
1110 let _ = reply.send(res);
1111 }
1112 Some(WriteCommand::Sweep(req, reply)) => {
1113 // Single-batch dispatch — the caller loops
1114 // until has_more=false. Per-batch return
1115 // lets the writer's biased select pick up
1116 // pending Apply / Negate / ResolveReport
1117 // commands between batches (§F4 + §F5).
1118 let res = self.handle_sweep(req).await;
1119 let _ = reply.send(res);
1120 }
1121 Some(WriteCommand::AppendAudit(row, reply)) => {
1122 let res = self.handle_append_audit(row).await;
1123 let _ = reply.send(res);
1124 }
1125 Some(WriteCommand::RecordAction(req, reply)) => {
1126 let res = self.handle_record_action(req).await;
1127 let _ = reply.send(res);
1128 }
1129 Some(WriteCommand::RevokeAction(req, reply)) => {
1130 let res = self.handle_revoke_action(req).await;
1131 let _ = reply.send(res);
1132 }
1133 Some(WriteCommand::ConfirmPendingAction(req, reply)) => {
1134 let res = self.handle_confirm_pending_action(req).await;
1135 let _ = reply.send(res);
1136 }
1137 Some(WriteCommand::DismissPendingAction(req, reply)) => {
1138 let res = self.handle_dismiss_pending_action(req).await;
1139 let _ = reply.send(res);
1140 }
1141 Some(WriteCommand::Shutdown(reply)) => {
1142 // Flip the shutdown watch *before* releasing the
1143 // lease so subscriber tasks see the signal while
1144 // the DB is still accessible for their close-
1145 // frame sends.
1146 let _ = self.shutdown_tx.send(true);
1147 let res = self.release_lease().await;
1148 let _ = reply.send(res);
1149 return;
1150 }
1151 None => {
1152 // All handles dropped without explicit shutdown.
1153 // Best-effort lease release so a same-process
1154 // restart doesn't trip the 60s wait.
1155 let _ = self.shutdown_tx.send(true);
1156 if let Err(e) = self.release_lease().await {
1157 tracing::error!("lease release on handle drop: {e}");
1158 }
1159 return;
1160 }
1161 }
1162 }
1163 _ = heartbeat_timer.tick() => {
1164 if let Err(e) = self.heartbeat().await {
1165 // Transient DB errors are logged, not fatal. A
1166 // sustained failure is caught at the next-instance
1167 // startup check — not at this writer's expense.
1168 tracing::error!("lease heartbeat failed: {e}");
1169 }
1170 }
1171 _ = sweep_check_timer.tick() => {
1172 if sweep_state.is_none()
1173 && let Some(fire_at) = next_scheduled_fire
1174 && Instant::now() >= fire_at
1175 {
1176 sweep_state = Some(SweepRunState {
1177 started_at: Instant::now(),
1178 rows: 0,
1179 batches: 0,
1180 });
1181 next_scheduled_fire = Some(fire_at + Duration::from_secs(86_400));
1182 tracing::info!(
1183 retention_days = ?self.retention_days,
1184 sweep_batch_size = self.retention.sweep_batch_size,
1185 "scheduled retention sweep starting"
1186 );
1187 }
1188 }
1189 // Always-ready arm gated on sweep_in_progress. With biased
1190 // order this ranks below rx.recv() + the two timers, so
1191 // normal commands and heartbeats interleave naturally
1192 // between batches (§F4 inter-batch yield, §F5 single-
1193 // writer invariant).
1194 _ = std::future::ready(()), if sweep_in_progress => {
1195 match self.handle_sweep(SweepRequest).await {
1196 Ok(batch) => {
1197 let state = sweep_state
1198 .as_mut()
1199 .expect("sweep_in_progress => sweep_state Some");
1200 state.rows += batch.rows_deleted;
1201 state.batches += 1;
1202 if !batch.has_more {
1203 let final_state = sweep_state.take().expect("just set");
1204 tracing::info!(
1205 rows_deleted = final_state.rows,
1206 batches = final_state.batches,
1207 duration_ms = final_state.started_at.elapsed().as_millis() as u64,
1208 retention_days_applied = ?batch.retention_days_applied,
1209 "scheduled retention sweep complete"
1210 );
1211 }
1212 }
1213 Err(e) => {
1214 let final_state = sweep_state.take();
1215 tracing::error!(
1216 error = %e,
1217 batches_done = final_state.as_ref().map(|s| s.batches).unwrap_or(0),
1218 rows_so_far = final_state.as_ref().map(|s| s.rows).unwrap_or(0),
1219 "scheduled retention sweep batch failed; aborting run (will retry on next schedule)"
1220 );
1221 }
1222 }
1223 }
1224 }
1225 }
1226 }
1227
1228 /// Process one batch of the §F4 retention sweep. Single-batch
1229 /// dispatch — see [`WriteCommand::Sweep`] for the loop ordering.
1230 ///
1231 /// Returns immediately with `rows_deleted = 0, has_more = false`
1232 /// when `retention_days` is `None`. Otherwise issues one
1233 /// `DELETE FROM labels WHERE created_at < cutoff LIMIT N` in its
1234 /// own transaction; sets `has_more = true` iff the batch hit the
1235 /// `sweep_batch_size` limit (suggesting more rows may match).
1236 ///
1237 /// Errors are propagated as `Err(_)` rather than swallowed: a
1238 /// transient DB failure should surface to the operator on a
1239 /// manual sweep, and to the schedule-loop logger on a scheduled
1240 /// sweep. Idempotency (Q5) means the next sweep retries cleanly.
1241 async fn handle_sweep(&self, _req: SweepRequest) -> Result<SweepBatchResult> {
1242 let Some(days) = self.retention_days else {
1243 return Ok(SweepBatchResult {
1244 rows_deleted: 0,
1245 has_more: false,
1246 retention_days_applied: None,
1247 });
1248 };
1249
1250 let cutoff_ms = epoch_ms_now() - (days as i64) * 86_400_000;
1251 let limit = self.retention.sweep_batch_size;
1252
1253 let mut tx = self.pool.begin().await?;
1254 // SQLite's DELETE doesn't support LIMIT without the
1255 // SQLITE_ENABLE_UPDATE_DELETE_LIMIT compile flag (off in
1256 // bundled builds). The rowid IN (SELECT ... LIMIT) trick is
1257 // the canonical portable workaround.
1258 let result = sqlx::query!(
1259 "DELETE FROM labels WHERE rowid IN (
1260 SELECT rowid FROM labels WHERE created_at < ?1 LIMIT ?2
1261 )",
1262 cutoff_ms,
1263 limit,
1264 )
1265 .execute(&mut *tx)
1266 .await?;
1267 tx.commit().await?;
1268
1269 let rows = result.rows_affected() as i64;
1270 Ok(SweepBatchResult {
1271 rows_deleted: rows,
1272 has_more: rows >= limit,
1273 retention_days_applied: Some(days),
1274 })
1275 }
1276
1277 /// Handler for [`WriteCommand::AppendAudit`]. Opens its own
1278 /// transaction (BEGIN DEFERRED — fine because the writer task is
1279 /// the only in-process audit appender during a `cairn serve`
1280 /// session, and cross-process appenders use BEGIN IMMEDIATE on
1281 /// their side), inserts the audit row with a freshly-computed
1282 /// hash, commits.
1283 async fn handle_append_audit(
1284 &self,
1285 row: crate::audit::append::AuditRowForAppend,
1286 ) -> Result<i64> {
1287 let mut tx = self.pool.begin().await?;
1288 let id = crate::audit::append::append_in_tx(&mut tx, &row).await?;
1289 tx.commit().await?;
1290 Ok(id)
1291 }
1292
1293 async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
1294 let mut tx = self.pool.begin().await?;
1295 let created_at = epoch_ms_now();
1296 let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
1297 tx.commit().await?;
1298 // No-receivers is not a write failure (§plan point G).
1299 let _ = self.broadcast_tx.send(event.clone());
1300 Ok(event)
1301 }
1302
1303 /// Inner label-application pipeline: reserve seq → clamp cts →
1304 /// sign → INSERT label → INSERT audit. Does NOT commit the
1305 /// transaction and does NOT broadcast — caller handles both.
1306 ///
1307 /// Extracted from `handle_apply` so `handle_resolve_report` can
1308 /// reuse it inside the same transaction as the report update,
1309 /// preserving §F5's single-writer-owns-label-emission invariant
1310 /// (this helper only runs on the writer task; no other code path
1311 /// can access seq allocation or cts clamping).
1312 async fn apply_label_inner(
1313 &self,
1314 tx: &mut sqlx::Transaction<'_, Sqlite>,
1315 req: &ApplyLabelRequest,
1316 created_at_ms: i64,
1317 ) -> Result<LabelEvent> {
1318 let event = self
1319 .sign_and_persist_label(
1320 tx,
1321 &req.val,
1322 &req.uri,
1323 req.cid.as_deref(),
1324 false,
1325 req.exp.as_deref(),
1326 created_at_ms,
1327 )
1328 .await?;
1329
1330 let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
1331 crate::audit::append::append_in_tx(
1332 tx,
1333 &crate::audit::append::AuditRowForAppend {
1334 created_at: created_at_ms,
1335 action: "label_applied".into(),
1336 actor_did: req.actor_did.clone(),
1337 target: Some(req.uri.clone()),
1338 target_cid: req.cid.clone(),
1339 outcome: "success".into(),
1340 reason: Some(audit_reason),
1341 },
1342 )
1343 .await?;
1344
1345 Ok(event)
1346 }
1347
1348 /// Tx-scoped sign-and-persist core for one label row: reserve
1349 /// seq → fetch prev cts → clamp → build wire-level [`Label`] →
1350 /// sign → INSERT into `labels`. Returns the seq + signed label.
1351 /// Does NOT write audit rows and does NOT broadcast — callers
1352 /// own both (the post-tx broadcast and any audit row layered
1353 /// over the bare label INSERT).
1354 ///
1355 /// Reused by [`Self::apply_label_inner`] (which adds a
1356 /// `label_applied` audit row) and by
1357 /// [`Self::handle_record_action`] (which writes one consolidated
1358 /// `subject_action_recorded` audit row capturing the action plus
1359 /// every emitted label, rather than per-label rows). Both reach
1360 /// it only from the single writer task, preserving §F5.
1361 #[allow(clippy::too_many_arguments)]
1362 async fn sign_and_persist_label(
1363 &self,
1364 tx: &mut sqlx::Transaction<'_, Sqlite>,
1365 val: &str,
1366 uri: &str,
1367 cid: Option<&str>,
1368 neg: bool,
1369 exp: Option<&str>,
1370 created_at_ms: i64,
1371 ) -> Result<LabelEvent> {
1372 let seq = reserve_seq(tx).await?;
1373 let prev_cts: Option<String> = sqlx::query_scalar!(
1374 // sqlx type override — MAX() strips column origin metadata so
1375 // sqlx can't infer the TEXT+nullable shape on its own. `?:`
1376 // forces the nullable wrapper (Option<String>).
1377 r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
1378 WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
1379 self.service_did,
1380 uri,
1381 val,
1382 )
1383 .fetch_one(&mut **tx)
1384 .await?;
1385
1386 let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
1387
1388 let cid_owned = cid.map(str::to_string);
1389 let exp_owned = exp.map(str::to_string);
1390 let val_owned = val.to_string();
1391 let uri_owned = uri.to_string();
1392
1393 let mut label = Label {
1394 ver: 1,
1395 src: self.service_did.clone(),
1396 uri: uri_owned,
1397 cid: cid_owned,
1398 val: val_owned,
1399 neg,
1400 cts,
1401 exp: exp_owned,
1402 sig: None,
1403 };
1404 label.sig = Some(sign_label(&self.key, &label)?);
1405 let sig_bytes = label.sig.expect("just set").to_vec();
1406
1407 let neg_int: i64 = if neg { 1 } else { 0 };
1408 sqlx::query!(
1409 "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
1410 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1411 seq,
1412 label.ver,
1413 label.src,
1414 label.uri,
1415 label.cid,
1416 label.val,
1417 neg_int,
1418 label.cts,
1419 label.exp,
1420 sig_bytes,
1421 self.signing_key_id,
1422 created_at_ms,
1423 )
1424 .execute(&mut **tx)
1425 .await?;
1426
1427 Ok(LabelEvent { seq, label })
1428 }
1429
1430 /// Tx-scoped wrapper around [`Self::sign_and_persist_label`] that
1431 /// converts a [`LabelDraft`] (the emission core's purpose-shaped
1432 /// output, with `cts/exp` as `SystemTime`) into the wire-level
1433 /// arguments. Used by [`Self::handle_record_action`] to sign +
1434 /// persist every draft produced by `resolve_action_labels` /
1435 /// `resolve_reason_labels`. The draft's `cts` is informational
1436 /// here — the actual cts on the labels row comes from
1437 /// `clamp_cts(created_at_ms, prev_cts)` inside
1438 /// [`Self::sign_and_persist_label`]. In v1.5 the recorder calls
1439 /// the resolvers with `now = effective_at`, so the two are
1440 /// equal pre-clamp.
1441 async fn sign_and_persist_label_from_draft(
1442 &self,
1443 tx: &mut sqlx::Transaction<'_, Sqlite>,
1444 draft: &LabelDraft,
1445 created_at_ms: i64,
1446 ) -> Result<LabelEvent> {
1447 let exp_str = match draft.exp {
1448 Some(st) => Some(rfc3339_from_epoch_ms(systemtime_to_epoch_ms(st)?)?),
1449 None => None,
1450 };
1451 self.sign_and_persist_label(
1452 tx,
1453 &draft.val,
1454 &draft.uri,
1455 draft.cid.as_deref(),
1456 draft.neg,
1457 exp_str.as_deref(),
1458 created_at_ms,
1459 )
1460 .await
1461 }
1462
1463 /// Insert the auto-recorded `subject_actions` row produced by a
1464 /// `[policy_automation]` rule firing (#73, v1.6,
1465 /// `mode=auto`). Runs the same shape as the moderator-recorded
1466 /// path: predict id, write audit row, INSERT, verify, emit
1467 /// labels (action label + reason labels), append linkage rows.
1468 /// All inside the existing transaction.
1469 ///
1470 /// Differs from the moderator path in three places:
1471 /// - `actor_kind = 'policy'` and `actor_did =
1472 /// SYNTHETIC_POLICY_ACTOR_DID`.
1473 /// - `triggered_by_policy_rule = rule.name`.
1474 /// - The audit reason JSON has NO `policy_consequence` field
1475 /// (the auto-action is the consequence; the precipitating
1476 /// action's audit row carries the cross-reference).
1477 #[allow(clippy::too_many_arguments)]
1478 async fn insert_policy_auto_action(
1479 &self,
1480 tx: &mut sqlx::Transaction<'_, Sqlite>,
1481 rule: &crate::policy::automation::PolicyRule,
1482 subject_did: &str,
1483 subject_uri: Option<&str>,
1484 history_plus_precip: &[ActionRecord],
1485 state_after_precip: &crate::moderation::decay::StrikeState,
1486 now_systemtime: SystemTime,
1487 created_at: i64,
1488 predicted_auto_id: i64,
1489 label_events: &mut Vec<LabelEvent>,
1490 ) -> Result<i64> {
1491 // The rule's reason_codes resolve via the operator's
1492 // [moderation_reasons] vocabulary — config validation in
1493 // #71 guarantees every code is declared, so the resolver
1494 // always succeeds here.
1495 let primary = resolve_primary_reason(&rule.reason_codes, &self.reason_vocabulary)?;
1496
1497 // Strike state for the auto-action: starts from
1498 // state_after_precip's count. Position-in-window is
1499 // computed from history+precipitating (the auto-action
1500 // hasn't been INSERTed yet).
1501 let strikes_at_time_of_auto = state_after_precip.current_count;
1502 let position =
1503 compute_position_in_window(history_plus_precip, &self.strike_policy, now_systemtime);
1504 let calc: StrikeApplication = strike_calculate(
1505 strikes_at_time_of_auto,
1506 &primary,
1507 &self.strike_policy,
1508 position,
1509 );
1510 // Note/Warning don't carry strikes regardless of reason.
1511 // Defense-in-depth (mirrors the moderator path).
1512 let (strike_base, strike_applied, was_dampened) = match rule.action_type {
1513 ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
1514 _ => (calc.base_weight, calc.applied, calc.was_dampened),
1515 };
1516
1517 // Auto-action expiry: derive from rule.duration when
1518 // action_type == TempSuspension. #71 already enforced the
1519 // pairing at config-load time.
1520 let auto_expires_at: Option<i64> = match (rule.action_type, rule.duration) {
1521 (ActionType::TempSuspension, Some(d)) => Some(created_at + (d.as_secs() as i64) * 1000),
1522 _ => None,
1523 };
1524 let auto_duration_iso: Option<String> = match (rule.action_type, rule.duration) {
1525 (ActionType::TempSuspension, Some(d)) => {
1526 // Round-trip: rule.duration was parsed from the
1527 // operator's config. We store the original
1528 // `P{n}D` / `PT{n}H` shape verbatim so audit
1529 // readers see what the operator declared.
1530 // Reconstruction would require remembering the
1531 // original string; v1.6 just regenerates a
1532 // canonical form (PT<seconds>S) for the row's
1533 // `duration` column. Operators reading the row
1534 // can compute the wall-clock end via
1535 // `expires_at`; the duration column is display-
1536 // only.
1537 Some(format!("PT{}S", d.as_secs()))
1538 }
1539 _ => None,
1540 };
1541
1542 // Resolve labels for the auto-action.
1543 let action_for_emission = ActionForEmission {
1544 action_type: rule.action_type,
1545 expires_at: auto_expires_at.map(epoch_ms_to_systemtime),
1546 subject_did: subject_did.to_string(),
1547 subject_uri: subject_uri.map(str::to_string),
1548 reason_codes: rule.reason_codes.clone(),
1549 cid: None,
1550 };
1551 let action_drafts = resolve_action_labels(
1552 &action_for_emission,
1553 &self.label_emission_policy,
1554 now_systemtime,
1555 );
1556 let reason_drafts = resolve_reason_labels(
1557 &action_for_emission,
1558 &self.label_emission_policy,
1559 now_systemtime,
1560 );
1561 let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
1562 .iter()
1563 .chain(reason_drafts.iter())
1564 .map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
1565 .collect();
1566
1567 // Predict + verify the auto-action's id matches what we
1568 // reserved when building the precipitating audit row's
1569 // policy_consequence pointer.
1570 let next_id = predict_next_subject_action_id(tx).await?;
1571 if next_id != predicted_auto_id {
1572 return Err(Error::Signing(format!(
1573 "policy auto-action: predicted id {predicted_auto_id} but next-id is now {next_id}; \
1574 policy_consequence reservation diverged"
1575 )));
1576 }
1577
1578 let audit_reason = build_record_action_audit_reason(
1579 next_id,
1580 rule.action_type,
1581 &primary.identifier,
1582 &rule.reason_codes,
1583 strike_base,
1584 strike_applied,
1585 was_dampened,
1586 &emitted_labels_for_audit,
1587 "policy",
1588 Some(&rule.name),
1589 None,
1590 );
1591 let audit_target = subject_uri
1592 .map(str::to_string)
1593 .unwrap_or_else(|| subject_did.to_string());
1594 let auto_audit_log_id = crate::audit::append::append_in_tx(
1595 tx,
1596 &crate::audit::append::AuditRowForAppend {
1597 created_at,
1598 action: "subject_action_recorded".into(),
1599 actor_did: crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID.to_string(),
1600 target: Some(audit_target),
1601 target_cid: None,
1602 outcome: "success".into(),
1603 reason: Some(audit_reason),
1604 },
1605 )
1606 .await?;
1607
1608 // INSERT the auto-action row. actor_kind = 'policy'
1609 // discriminates from moderator-recorded rows; actor_did =
1610 // SYNTHETIC_POLICY_ACTOR_DID gives downstream filtering
1611 // a stable identifier; triggered_by_policy_rule preserves
1612 // provenance for audit + revocation.
1613 let action_type_str = rule.action_type.as_db_str();
1614 let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
1615 let strikes_at_time_i64 = strikes_at_time_of_auto as i64;
1616 let strike_base_i64 = strike_base as i64;
1617 let strike_applied_i64 = strike_applied as i64;
1618 let reason_codes_json = serde_json::to_string(&rule.reason_codes)
1619 .map_err(|e| Error::Signing(format!("serialize policy reason_codes: {e}")))?;
1620 let actor_kind_policy = "policy";
1621 let synth_actor = crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID;
1622 let auto_inserted_id = sqlx::query_scalar!(
1623 "INSERT INTO subject_actions (
1624 subject_did, subject_uri, actor_did, action_type, reason_codes,
1625 duration, effective_at, expires_at, notes, report_ids,
1626 strike_value_base, strike_value_applied, was_dampened,
1627 strikes_at_time_of_action, audit_log_id, created_at,
1628 actor_kind, triggered_by_policy_rule
1629 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17)
1630 RETURNING id",
1631 subject_did,
1632 subject_uri,
1633 synth_actor,
1634 action_type_str,
1635 reason_codes_json,
1636 auto_duration_iso,
1637 created_at,
1638 auto_expires_at,
1639 None::<String>,
1640 strike_base_i64,
1641 strike_applied_i64,
1642 was_dampened_int,
1643 strikes_at_time_i64,
1644 auto_audit_log_id,
1645 created_at,
1646 actor_kind_policy,
1647 rule.name,
1648 )
1649 .fetch_one(&mut **tx)
1650 .await?;
1651
1652 if auto_inserted_id != predicted_auto_id {
1653 return Err(Error::Signing(format!(
1654 "policy auto-action inserted at id {auto_inserted_id} but predicted \
1655 {predicted_auto_id}; audit chain captured the predicted id (corrupted)"
1656 )));
1657 }
1658
1659 // Emit labels for the auto-action — same path as the
1660 // moderator-recorded action's emission. Idempotency
1661 // guards from #64 are no-ops on a fresh INSERT.
1662 let mut auto_action_label_val: Option<String> = None;
1663 for draft in &action_drafts {
1664 let event = self
1665 .sign_and_persist_label_from_draft(tx, draft, created_at)
1666 .await?;
1667 if auto_action_label_val.is_none() {
1668 auto_action_label_val = Some(draft.val.clone());
1669 }
1670 label_events.push(event);
1671 }
1672 if let Some(ref val) = auto_action_label_val {
1673 sqlx::query!(
1674 "UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
1675 val,
1676 auto_inserted_id,
1677 )
1678 .execute(&mut **tx)
1679 .await?;
1680 }
1681 for (draft, reason_code) in reason_drafts.iter().zip(rule.reason_codes.iter()) {
1682 let event = self
1683 .sign_and_persist_label_from_draft(tx, draft, created_at)
1684 .await?;
1685 sqlx::query!(
1686 "INSERT INTO subject_action_reason_labels
1687 (action_id, reason_code, emitted_label_uri, emitted_at)
1688 VALUES (?1, ?2, ?3, ?4)",
1689 auto_inserted_id,
1690 reason_code,
1691 draft.val,
1692 created_at,
1693 )
1694 .execute(&mut **tx)
1695 .await?;
1696 label_events.push(event);
1697 }
1698
1699 // ---------- takedown cascade (#76, v1.6) ----------
1700 //
1701 // A policy-auto-recorded Takedown auto-dismisses every
1702 // unresolved pending_policy_actions row for the subject —
1703 // including any pending whose own rule fired in mode=flag
1704 // earlier in the subject's history. The cascade's
1705 // resolved_by_did is the synthetic policy DID, marking
1706 // the dismissal as system-driven rather than moderator-
1707 // driven for downstream filtering.
1708 //
1709 // The pending row that the rule itself proposed (when
1710 // this auto-action's own rule was originally mode=flag)
1711 // can't appear here: a rule firing in mode=auto inserts a
1712 // subject_actions row directly via this function and does
1713 // NOT create a pending_policy_actions row. The cascade
1714 // sees only OTHER subject pendings.
1715 if rule.action_type == ActionType::Takedown {
1716 auto_dismiss_pendings_on_takedown(
1717 tx,
1718 subject_did,
1719 subject_uri,
1720 auto_inserted_id,
1721 crate::policy::automation::SYNTHETIC_POLICY_ACTOR_DID,
1722 created_at,
1723 )
1724 .await?;
1725 }
1726
1727 Ok(auto_inserted_id)
1728 }
1729
1730 /// Insert a `pending_policy_actions` row for a `mode=flag`
1731 /// rule firing (#73, v1.6). The pending row holds the
1732 /// proposed action's shape until a moderator confirms (#74)
1733 /// or dismisses (#75); no `subject_actions` row, no label
1734 /// emission. The precipitating action's audit row already
1735 /// cross-references this pending row's id via the
1736 /// `policy_consequence` field.
1737 #[allow(clippy::too_many_arguments)]
1738 async fn insert_pending_policy_action(
1739 &self,
1740 tx: &mut sqlx::Transaction<'_, Sqlite>,
1741 rule: &crate::policy::automation::PolicyRule,
1742 subject_did: &str,
1743 subject_uri: Option<&str>,
1744 triggering_action_id: i64,
1745 created_at: i64,
1746 predicted_pending_id: i64,
1747 ) -> Result<i64> {
1748 let action_type_str = rule.action_type.as_db_str();
1749 let duration_ms: Option<i64> = match (rule.action_type, rule.duration) {
1750 (ActionType::TempSuspension, Some(d)) => Some((d.as_secs() as i64) * 1000),
1751 _ => None,
1752 };
1753 let reason_codes_json = serde_json::to_string(&rule.reason_codes).map_err(|e| {
1754 Error::Signing(format!("serialize policy reason_codes for pending: {e}"))
1755 })?;
1756
1757 let inserted_id_opt = sqlx::query_scalar!(
1758 r#"INSERT INTO pending_policy_actions (
1759 subject_did, subject_uri, action_type, duration_ms,
1760 reason_codes, triggered_by_policy_rule, triggered_at,
1761 triggering_action_id
1762 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
1763 RETURNING id AS "id!: i64""#,
1764 subject_did,
1765 subject_uri,
1766 action_type_str,
1767 duration_ms,
1768 reason_codes_json,
1769 rule.name,
1770 created_at,
1771 triggering_action_id,
1772 )
1773 .fetch_one(&mut **tx)
1774 .await?;
1775 let inserted_id: i64 = inserted_id_opt;
1776
1777 if inserted_id != predicted_pending_id {
1778 return Err(Error::Signing(format!(
1779 "pending_policy_actions inserted at id {inserted_id} but predicted \
1780 {predicted_pending_id}; policy_consequence reservation diverged"
1781 )));
1782 }
1783 Ok(inserted_id)
1784 }
1785
1786 async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
1787 let mut tx = self.pool.begin().await?;
1788
1789 // Most-recent event for the tuple. Not-found OR latest-is-already-
1790 // a-negation both surface as LabelNotFound (§F6).
1791 let latest = sqlx::query!(
1792 "SELECT neg, cid FROM labels
1793 WHERE src = ?1 AND uri = ?2 AND val = ?3
1794 ORDER BY seq DESC LIMIT 1",
1795 self.service_did,
1796 req.uri,
1797 req.val,
1798 )
1799 .fetch_optional(&mut *tx)
1800 .await?;
1801
1802 let cid = match latest {
1803 Some(row) if row.neg == 0 => row.cid,
1804 _ => {
1805 return Err(Error::LabelNotFound {
1806 src: self.service_did.clone(),
1807 uri: req.uri,
1808 val: req.val,
1809 });
1810 }
1811 };
1812
1813 let seq = reserve_seq(&mut tx).await?;
1814 // Prev cts lookup is the same query as apply; tuple uniqueness is
1815 // on (src, uri, val) regardless of neg.
1816 let prev_cts: Option<String> = sqlx::query_scalar!(
1817 // sqlx type override — MAX() strips column origin metadata so
1818 // sqlx can't infer the TEXT+nullable shape on its own. `?:`
1819 // forces the nullable wrapper (Option<String>).
1820 r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
1821 WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
1822 self.service_did,
1823 req.uri,
1824 req.val,
1825 )
1826 .fetch_one(&mut *tx)
1827 .await?;
1828
1829 let wall_now_ms = epoch_ms_now();
1830 let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
1831
1832 let mut label = Label {
1833 ver: 1,
1834 src: self.service_did.clone(),
1835 uri: req.uri.clone(),
1836 cid: cid.clone(),
1837 val: req.val.clone(),
1838 neg: true,
1839 cts: cts.clone(),
1840 exp: None, // Negations never carry an expiry.
1841 sig: None,
1842 };
1843 label.sig = Some(sign_label(&self.key, &label)?);
1844 let sig_bytes = label.sig.expect("just set").to_vec();
1845
1846 let created_at = wall_now_ms;
1847 let neg_int: i64 = 1;
1848 sqlx::query!(
1849 "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
1850 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1851 seq,
1852 label.ver,
1853 label.src,
1854 label.uri,
1855 label.cid,
1856 label.val,
1857 neg_int,
1858 label.cts,
1859 label.exp,
1860 sig_bytes,
1861 self.signing_key_id,
1862 created_at,
1863 )
1864 .execute(&mut *tx)
1865 .await?;
1866
1867 let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
1868 crate::audit::append::append_in_tx(
1869 &mut tx,
1870 &crate::audit::append::AuditRowForAppend {
1871 created_at,
1872 action: "label_negated".into(),
1873 actor_did: req.actor_did.clone(),
1874 target: Some(req.uri.clone()),
1875 target_cid: cid.clone(),
1876 outcome: "success".into(),
1877 reason: Some(audit_reason),
1878 },
1879 )
1880 .await?;
1881
1882 tx.commit().await?;
1883
1884 let event = LabelEvent { seq, label };
1885 let _ = self.broadcast_tx.send(event.clone());
1886 Ok(event)
1887 }
1888
1889 /// Atomic resolveReport flow — §F12 requires label INSERT +
1890 /// audit(label_applied) + report UPDATE + audit(report_resolved)
1891 /// all commit together or not at all. Everything happens inside
1892 /// one SQLite transaction on the writer task, preserving §F5's
1893 /// single-writer invariant for the optional label emission.
1894 ///
1895 /// Returns [`Error::ReportNotFound`] if the report doesn't
1896 /// exist, [`Error::ReportAlreadyResolved`] if it is not in
1897 /// `status='pending'`. Label-value validation is a handler
1898 /// concern (the handler consults `AdminConfig.label_values`
1899 /// before sending the command here — saves a writer round-trip
1900 /// on invalid input and keeps the anti-leak message local).
1901 async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
1902 use crate::report::ReportStatus;
1903
1904 let mut tx = self.pool.begin().await?;
1905
1906 // 1. Load the report; confirm pending status.
1907 let current = sqlx::query_as!(
1908 crate::report::Report,
1909 r#"SELECT
1910 id AS "id!: i64",
1911 created_at AS "created_at!: String",
1912 reported_by AS "reported_by!: String",
1913 reason_type AS "reason_type!: String",
1914 reason,
1915 subject_type AS "subject_type!: String",
1916 subject_did AS "subject_did!: String",
1917 subject_uri,
1918 subject_cid,
1919 status AS "status!: ReportStatus",
1920 resolved_at,
1921 resolved_by,
1922 resolution_label,
1923 resolution_reason
1924 FROM reports WHERE id = ?1"#,
1925 req.report_id,
1926 )
1927 .fetch_optional(&mut *tx)
1928 .await?;
1929
1930 let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
1931 if report.status != ReportStatus::Pending {
1932 return Err(Error::ReportAlreadyResolved { id: req.report_id });
1933 }
1934
1935 let created_at = epoch_ms_now();
1936
1937 // 2. Optional label-apply BEFORE the report UPDATE so audit
1938 // row insertion order reflects logical sequence
1939 // (label_applied, then report_resolved — §G per #15 criteria).
1940 let label_event = if let Some(apply) = req.action.as_apply() {
1941 let apply_req = ApplyLabelRequest {
1942 actor_did: req.actor_did.clone(),
1943 uri: apply.uri.clone(),
1944 cid: apply.cid.clone(),
1945 val: apply.val.clone(),
1946 exp: apply.exp.clone(),
1947 // The label's own audit_reason JSON captures {val,
1948 // neg, moderator_reason}; the resolution-level reason
1949 // goes on the report_resolved audit row below.
1950 moderator_reason: None,
1951 };
1952 Some(
1953 self.apply_label_inner(&mut tx, &apply_req, created_at)
1954 .await?,
1955 )
1956 } else {
1957 None
1958 };
1959
1960 // 3. UPDATE reports.
1961 let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
1962 let resolution_label = req.action.as_apply().map(|a| a.val.clone());
1963 let resolved_status = ReportStatus::Resolved;
1964 sqlx::query!(
1965 "UPDATE reports SET
1966 status = ?1,
1967 resolved_at = ?2,
1968 resolved_by = ?3,
1969 resolution_label = ?4,
1970 resolution_reason = ?5
1971 WHERE id = ?6",
1972 resolved_status,
1973 resolved_at_rfc,
1974 req.actor_did,
1975 resolution_label,
1976 req.resolution_reason,
1977 req.report_id,
1978 )
1979 .execute(&mut *tx)
1980 .await?;
1981
1982 // 4. Audit: report_resolved.
1983 let audit_reason = build_resolve_audit_reason(
1984 req.action.as_apply().map(|a| a.val.as_str()),
1985 req.resolution_reason.as_deref(),
1986 );
1987 crate::audit::append::append_in_tx(
1988 &mut tx,
1989 &crate::audit::append::AuditRowForAppend {
1990 created_at,
1991 action: "report_resolved".into(),
1992 actor_did: req.actor_did.clone(),
1993 target: Some(req.report_id.to_string()),
1994 target_cid: None,
1995 outcome: "success".into(),
1996 reason: Some(audit_reason),
1997 },
1998 )
1999 .await?;
2000
2001 tx.commit().await?;
2002
2003 // 5. Broadcast (post-commit, per §F5 broadcast-after-commit rule).
2004 if let Some(event) = &label_event {
2005 let _ = self.broadcast_tx.send(event.clone());
2006 }
2007
2008 // 6. Mutate the loaded report struct to reflect the committed
2009 // state and return it. Avoids a re-fetch round-trip since we
2010 // know exactly what changed.
2011 report.status = ReportStatus::Resolved;
2012 report.resolved_at = Some(resolved_at_rfc);
2013 report.resolved_by = Some(req.actor_did);
2014 report.resolution_label = resolution_label;
2015 report.resolution_reason = req.resolution_reason;
2016
2017 Ok(ResolvedReport {
2018 report,
2019 label_event,
2020 })
2021 }
2022
2023 /// Atomic recordAction flow (§F20 / #51). Validates inputs,
2024 /// resolves the primary reason, computes strike values via the
2025 /// v1.4 calculators (#48/#49/#50/#51), inserts the audit_log
2026 /// row first (so subject_actions can carry audit_log_id at
2027 /// INSERT — the schema's no-update-except-revoke trigger
2028 /// forbids backfilling it later), inserts the subject_actions
2029 /// row, then UPSERTs subject_strike_state. All in one
2030 /// transaction.
2031 async fn handle_record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
2032 // ---------- pre-flight validation ----------
2033
2034 if req.reason_codes.is_empty() {
2035 return Err(Error::Signing(
2036 "recordAction: reason_codes must be non-empty".into(),
2037 ));
2038 }
2039 // Note/Warning don't take reasons in the §F20 design intent,
2040 // but the issue body says "warning requires --reason" for
2041 // the audit trail. Accept reasons for all types; only
2042 // strike-bearing types actually use them for strike
2043 // calculation.
2044 match req.action_type {
2045 ActionType::TempSuspension => {
2046 if req.duration_iso.as_deref().unwrap_or("").is_empty() {
2047 return Err(Error::DurationRequiredForTempSuspension);
2048 }
2049 }
2050 _ => {
2051 if req.duration_iso.is_some() {
2052 return Err(Error::DurationOnlyForTempSuspension);
2053 }
2054 }
2055 }
2056
2057 let (subject_did, subject_uri) = route_subject(&req.subject)?;
2058
2059 // Resolve primary reason against the writer's vocabulary.
2060 // Vocabulary lookups are cheap; doing it before opening the
2061 // transaction keeps a bad-reason failure from acquiring the
2062 // SQLite write lock unnecessarily.
2063 let primary = resolve_primary_reason(&req.reason_codes, &self.reason_vocabulary)?;
2064
2065 // Parse duration (cheap; pre-tx for the same reason).
2066 let duration_secs = match (req.action_type, req.duration_iso.as_deref()) {
2067 (ActionType::TempSuspension, Some(iso)) => Some(parse_iso8601_duration(iso)?),
2068 _ => None,
2069 };
2070
2071 // ---------- transaction ----------
2072
2073 let mut tx = self.pool.begin().await?;
2074 let created_at = epoch_ms_now();
2075 let effective_at = created_at;
2076 let expires_at: Option<i64> = duration_secs.map(|s| effective_at + (s as i64) * 1000);
2077
2078 // Load history for strike + position calculation. id-ascending
2079 // (oldest-first) per the calculators' contract.
2080 let history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
2081
2082 // Compute current count via the decay calculator.
2083 let now_systemtime = epoch_ms_to_systemtime(created_at);
2084 let pre_state = calculate_strike_state(&history, &self.strike_policy, now_systemtime);
2085 let strikes_at_time_of_action = pre_state.current_count;
2086
2087 // Compute position and strike value.
2088 let position = compute_position_in_window(&history, &self.strike_policy, now_systemtime);
2089 let calc: StrikeApplication = strike_calculate(
2090 strikes_at_time_of_action,
2091 &primary,
2092 &self.strike_policy,
2093 position,
2094 );
2095
2096 // Note/Warning don't carry strikes regardless of reason
2097 // (§F20 design + #48 semantics). Override with zero values
2098 // and was_dampened=false; defense-in-depth for an operator
2099 // who attaches a strike-bearing reason to a Note row.
2100 let (strike_base, strike_applied, was_dampened) = match req.action_type {
2101 ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
2102 _ => (calc.base_weight, calc.applied, calc.was_dampened),
2103 };
2104
2105 // Compute label-emission drafts before the audit row so the
2106 // audit reason JSON captures every (val, uri) tuple this
2107 // action will produce. The drafts are pure data — no I/O,
2108 // no signing — so committing them to the audit row before
2109 // signing is safe; the actual label INSERTs land below in
2110 // the same tx, and any failure rolls everything back.
2111 //
2112 // LabelDraft.cid is None for v1.5: subject CIDs aren't yet
2113 // plumbed through RecordActionRequest. The
2114 // ActionForEmission.cid field is reserved for the wiring a
2115 // future ticket adds via admin XRPC + CLI.
2116 let action_for_emission = ActionForEmission {
2117 action_type: req.action_type,
2118 expires_at: expires_at.map(epoch_ms_to_systemtime),
2119 subject_did: subject_did.clone(),
2120 subject_uri: subject_uri.clone(),
2121 reason_codes: req.reason_codes.clone(),
2122 cid: None,
2123 };
2124 let action_drafts = resolve_action_labels(
2125 &action_for_emission,
2126 &self.label_emission_policy,
2127 now_systemtime,
2128 );
2129 let reason_drafts = resolve_reason_labels(
2130 &action_for_emission,
2131 &self.label_emission_policy,
2132 now_systemtime,
2133 );
2134 let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
2135 .iter()
2136 .chain(reason_drafts.iter())
2137 .map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
2138 .collect();
2139
2140 // Build audit_log reason JSON before the INSERT so the
2141 // audit row can carry the resolved values for forensic
2142 // reconstruction.
2143 // action_id is unknown until subject_actions INSERT, so we
2144 // patch it in after that step.
2145 let reason_codes_json = serde_json::to_string(&req.reason_codes)
2146 .map_err(|e| Error::Signing(format!("serialize reason_codes: {e}")))?;
2147 let report_ids_json = if req.report_ids.is_empty() {
2148 None
2149 } else {
2150 Some(
2151 serde_json::to_string(&req.report_ids)
2152 .map_err(|e| Error::Signing(format!("serialize report_ids: {e}")))?,
2153 )
2154 };
2155
2156 // Append audit_log row first. target = subject_did (the
2157 // account being moderated); target_cid = None. We don't
2158 // know action_id yet — pass 0 as a placeholder; the actual
2159 // id is patched into the reason JSON after the
2160 // subject_actions INSERT, but only the reason JSON gets
2161 // the patch. The audit row's hash chain locks at this
2162 // INSERT, so the action_id is captured at write time.
2163 //
2164 // Order rationale: audit_log first so subject_actions
2165 // can carry audit_log_id at INSERT (the trigger forbids
2166 // UPDATE of audit_log_id). The audit row's reason JSON
2167 // captures the action_id IT will reference, computed
2168 // post-INSERT via the next-id query.
2169 let next_action_id = predict_next_subject_action_id(&mut tx).await?;
2170
2171 // ---------- policy evaluation (#73, v1.6) ----------
2172 //
2173 // Evaluate operator-declared rules BEFORE writing the
2174 // precipitating audit row so the audit reason JSON's
2175 // `policy_consequence` field can cross-reference the
2176 // predicted auto-action id (mode=auto) or the predicted
2177 // pending row id (mode=flag). Hash chain locks the
2178 // bundle (action + policy consequence) atomically.
2179 //
2180 // The evaluator is pure — it consumes already-computed
2181 // strike states and projections. We construct
2182 // `state_after_precip` synthetically (history + the
2183 // about-to-INSERT precipitating row's effect) so we can
2184 // call into #72 before the row actually lands. This
2185 // matches the v1.5 #60 pattern of computing label
2186 // drafts pre-INSERT and verifying via the
2187 // predict-then-verify dance.
2188 let policy_eval_history =
2189 load_subject_actions_for_policy_eval(&mut tx, &subject_did).await?;
2190 let pending_eval_actions = load_pending_for_policy_eval(&mut tx, &subject_did).await?;
2191 let synthetic_precip = ActionRecord {
2192 strike_value_applied: strike_applied,
2193 effective_at: now_systemtime,
2194 revoked_at: None,
2195 action_type: req.action_type,
2196 expires_at: expires_at.map(epoch_ms_to_systemtime),
2197 was_dampened,
2198 };
2199 let mut history_plus_precip = history.clone();
2200 history_plus_precip.push(synthetic_precip);
2201 let state_after_precip =
2202 calculate_strike_state(&history_plus_precip, &self.strike_policy, now_systemtime);
2203
2204 // Project the synthetic precipitating action onto
2205 // ActionForPolicyEval shape and append to the eval
2206 // history; the evaluator's takedown-detection +
2207 // idempotency checks include the row that's about to
2208 // INSERT (matters when this very recordAction is itself
2209 // a takedown — the evaluator returns None immediately).
2210 let mut policy_eval_history_plus_precip = policy_eval_history.clone();
2211 policy_eval_history_plus_precip.push(crate::policy::evaluator::ActionForPolicyEval {
2212 effective_at: now_systemtime,
2213 action_type: req.action_type,
2214 revoked_at: None,
2215 triggered_by_policy_rule: None,
2216 });
2217
2218 let firing_rule = crate::policy::evaluator::resolve_firing_rule(
2219 &pre_state,
2220 &state_after_precip,
2221 &policy_eval_history_plus_precip,
2222 &pending_eval_actions,
2223 &self.policy_automation_policy,
2224 );
2225
2226 // Reserve IDs + build policy_consequence pointer for the
2227 // precipitating audit row. For mode=auto the auto-action
2228 // id is `next_action_id + 1` (the writer task is the only
2229 // in-process appender, so consecutive INSERTs assign
2230 // consecutive ids; the predict-then-verify check below
2231 // catches any drift). For mode=flag we predict the next
2232 // pending_policy_actions id.
2233 let (policy_consequence, predicted_auto_action_id, predicted_pending_id, firing_rule_owned) =
2234 if let Some(rule) = firing_rule {
2235 match rule.mode {
2236 crate::policy::automation::PolicyMode::Auto => {
2237 let auto_id = next_action_id + 1;
2238 let pc = serde_json::json!({
2239 "rule_fired": rule.name,
2240 "mode": "auto",
2241 "auto_action_id": auto_id,
2242 });
2243 (Some(pc), Some(auto_id), None, Some(rule.clone()))
2244 }
2245 crate::policy::automation::PolicyMode::Flag => {
2246 let pending_id = predict_next_pending_action_id(&mut tx).await?;
2247 let pc = serde_json::json!({
2248 "rule_fired": rule.name,
2249 "mode": "flag",
2250 "pending_action_id": pending_id,
2251 });
2252 (Some(pc), None, Some(pending_id), Some(rule.clone()))
2253 }
2254 }
2255 } else {
2256 (None, None, None, None)
2257 };
2258
2259 let audit_reason = build_record_action_audit_reason(
2260 next_action_id,
2261 req.action_type,
2262 &primary.identifier,
2263 &req.reason_codes,
2264 strike_base,
2265 strike_applied,
2266 was_dampened,
2267 &emitted_labels_for_audit,
2268 "moderator",
2269 None,
2270 policy_consequence,
2271 );
2272 let audit_target = subject_uri.clone().unwrap_or_else(|| subject_did.clone());
2273 let audit_log_id = crate::audit::append::append_in_tx(
2274 &mut tx,
2275 &crate::audit::append::AuditRowForAppend {
2276 created_at,
2277 action: "subject_action_recorded".into(),
2278 actor_did: req.actor_did.clone(),
2279 target: Some(audit_target),
2280 target_cid: None,
2281 outcome: "success".into(),
2282 reason: Some(audit_reason),
2283 },
2284 )
2285 .await?;
2286
2287 // INSERT subject_actions. action_type goes through
2288 // ActionType::as_db_str so the SQL CHECK constraint binds
2289 // the same set as the Rust enum.
2290 // actor_kind = 'moderator' here — this is the moderator-
2291 // recorded path; the policy-recorded path lower in this
2292 // function sets 'policy' explicitly.
2293 let action_type_str = req.action_type.as_db_str();
2294 let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
2295 let strikes_at_time_i64 = strikes_at_time_of_action as i64;
2296 let strike_base_i64 = strike_base as i64;
2297 let strike_applied_i64 = strike_applied as i64;
2298 let actor_kind_moderator = "moderator";
2299 let inserted_id = sqlx::query_scalar!(
2300 "INSERT INTO subject_actions (
2301 subject_did, subject_uri, actor_did, action_type, reason_codes,
2302 duration, effective_at, expires_at, notes, report_ids,
2303 strike_value_base, strike_value_applied, was_dampened,
2304 strikes_at_time_of_action, audit_log_id, created_at,
2305 actor_kind, triggered_by_policy_rule
2306 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, NULL)
2307 RETURNING id",
2308 subject_did,
2309 subject_uri,
2310 req.actor_did,
2311 action_type_str,
2312 reason_codes_json,
2313 req.duration_iso,
2314 effective_at,
2315 expires_at,
2316 req.notes,
2317 report_ids_json,
2318 strike_base_i64,
2319 strike_applied_i64,
2320 was_dampened_int,
2321 strikes_at_time_i64,
2322 audit_log_id,
2323 created_at,
2324 actor_kind_moderator,
2325 )
2326 .fetch_one(&mut *tx)
2327 .await?;
2328
2329 if inserted_id != next_action_id {
2330 // Audit row already committed; we'd have to invalidate
2331 // it. With BEGIN DEFERRED + the writer task being the
2332 // only in-process appender, the predict-vs-insert
2333 // race is closed — but if it ever fires, surface as
2334 // an internal error rather than a silent mismatch.
2335 return Err(Error::Signing(format!(
2336 "subject_actions inserted at id {inserted_id} but predicted {next_action_id}; audit chain captured the predicted id (corrupted)"
2337 )));
2338 }
2339
2340 // ---------- idempotency guard (#64, v1.5) ----------
2341 //
2342 // Defense-in-depth: skip emission for any state already
2343 // present on the row. v1.5's recordAction flow always
2344 // finds these queries returning the empty/NULL case (the
2345 // INSERT just landed and nothing else has touched the row
2346 // yet), so the guard is structurally a no-op in
2347 // production. It exists to protect against future paths
2348 // that might invoke emission against a row already
2349 // carrying linkage state — backfill migrations, retry
2350 // helpers, alternate code paths — and against bugs that
2351 // would otherwise silently produce duplicate (src, uri,
2352 // val) records on the wire.
2353 //
2354 // The audit row's `emitted_labels` was already written
2355 // above (intent, derived from the drafts). In v1.5 normal
2356 // flow intent matches reality. If a future defensive
2357 // scenario fires the guard and skips emissions, the audit
2358 // row will claim more emissions than the labels table
2359 // holds. That divergence is the future-implementer's
2360 // responsibility to handle (e.g., by adding a re-emission
2361 // entry point that constructs its own audit row); v1.5
2362 // accepts the limitation because the divergence is
2363 // structurally unreachable from this code path.
2364 //
2365 // The guard's "fire" branch is structurally unreachable
2366 // from the public XRPC + writer-task API in v1.5 (every
2367 // recordAction goes through this same INSERT). That makes
2368 // it untestable end-to-end without a re-emission entry
2369 // point we don't ship; the pure decision logic is
2370 // exercised via unit tests on
2371 // `should_skip_action_label_emission` /
2372 // `should_skip_reason_emission` instead.
2373 let existing_action_label_val: Option<String> = sqlx::query_scalar!(
2374 "SELECT emitted_label_uri FROM subject_actions WHERE id = ?1",
2375 inserted_id,
2376 )
2377 .fetch_one(&mut *tx)
2378 .await?;
2379 let existing_reason_codes: Vec<String> = sqlx::query_scalar!(
2380 "SELECT reason_code FROM subject_action_reason_labels WHERE action_id = ?1",
2381 inserted_id,
2382 )
2383 .fetch_all(&mut *tx)
2384 .await?;
2385 let existing_reason_set: HashSet<&str> =
2386 existing_reason_codes.iter().map(String::as_str).collect();
2387
2388 // ---------- emission (#60, v1.5) ----------
2389 //
2390 // Sign + persist each draft (action label first so its
2391 // seq < reason labels'), then UPDATE
2392 // subject_actions.emitted_label_uri, then write per-reason
2393 // linkage rows. All in the same tx as the action INSERT
2394 // and the audit row. Failure here rolls back everything
2395 // up to and including the audit row, so the audit chain
2396 // never claims emission that didn't happen.
2397 //
2398 // Empty action_drafts (note / suppressed warning / policy
2399 // disabled) cleanly short-circuits both loops, leaves
2400 // emitted_label_uri NULL, and writes no reason-label
2401 // linkage rows — the action still records, just without
2402 // an ATProto label tail.
2403 let mut label_events: Vec<LabelEvent> =
2404 Vec::with_capacity(action_drafts.len() + reason_drafts.len());
2405 let mut action_label_val: Option<String> = None;
2406 if !should_skip_action_label_emission(existing_action_label_val.as_deref()) {
2407 for draft in &action_drafts {
2408 let event = self
2409 .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
2410 .await?;
2411 if action_label_val.is_none() {
2412 action_label_val = Some(draft.val.clone());
2413 }
2414 label_events.push(event);
2415 }
2416
2417 if let Some(ref val) = action_label_val {
2418 // emitted_label_uri stores the action label's `val`,
2419 // not a URI: ATProto labels have no canonical URI; the
2420 // discriminator within (src=service_did, uri=subject,
2421 // val) is just the val. Column name predates the
2422 // realization (#57); locked once shipped.
2423 sqlx::query!(
2424 "UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
2425 val,
2426 inserted_id,
2427 )
2428 .execute(&mut *tx)
2429 .await?;
2430 }
2431 }
2432
2433 // Per-(action, reason) linkage rows. resolve_reason_labels
2434 // emits drafts in `reason_codes` order, so a parallel
2435 // iteration recovers the source reason_code without
2436 // re-parsing the label val. Reasons whose linkage row
2437 // already exists are skipped per the idempotency guard
2438 // above.
2439 for (draft, reason_code) in reason_drafts.iter().zip(req.reason_codes.iter()) {
2440 if should_skip_reason_emission(reason_code, &existing_reason_set) {
2441 continue;
2442 }
2443 let event = self
2444 .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
2445 .await?;
2446 sqlx::query!(
2447 "INSERT INTO subject_action_reason_labels
2448 (action_id, reason_code, emitted_label_uri, emitted_at)
2449 VALUES (?1, ?2, ?3, ?4)",
2450 inserted_id,
2451 reason_code,
2452 draft.val,
2453 effective_at,
2454 )
2455 .execute(&mut *tx)
2456 .await?;
2457 label_events.push(event);
2458 }
2459
2460 // ---------- policy consequence (#73, v1.6) ----------
2461 //
2462 // If a rule fires, this is where the consequence lands —
2463 // either a second `subject_actions` row (mode=auto) with
2464 // its own audit row + label emission, or a
2465 // `pending_policy_actions` row (mode=flag) awaiting
2466 // moderator review. All in the same transaction as the
2467 // precipitating action; failure rolls everything back.
2468 if let Some(rule) = firing_rule_owned.as_ref() {
2469 match rule.mode {
2470 crate::policy::automation::PolicyMode::Auto => {
2471 let auto_id =
2472 predicted_auto_action_id.expect("auto rule firing reserved an id above");
2473 self.insert_policy_auto_action(
2474 &mut tx,
2475 rule,
2476 &subject_did,
2477 subject_uri.as_deref(),
2478 &history_plus_precip,
2479 &state_after_precip,
2480 now_systemtime,
2481 created_at,
2482 auto_id,
2483 &mut label_events,
2484 )
2485 .await?;
2486 }
2487 crate::policy::automation::PolicyMode::Flag => {
2488 let pending_id =
2489 predicted_pending_id.expect("flag rule firing reserved an id above");
2490 self.insert_pending_policy_action(
2491 &mut tx,
2492 rule,
2493 &subject_did,
2494 subject_uri.as_deref(),
2495 inserted_id,
2496 created_at,
2497 pending_id,
2498 )
2499 .await?;
2500 }
2501 }
2502 }
2503
2504 // ---------- takedown cascade (#76, v1.6) ----------
2505 //
2506 // Moderator-recorded Takedown auto-dismisses every
2507 // unresolved pending_policy_actions row for the subject.
2508 // The cascade audit rows hash-chain after the
2509 // precipitating action's audit row (and any policy-
2510 // consequence audit row from the auto-action branch
2511 // above), so a forensic reader sees: takedown → cascade
2512 // dismissals, in commit order.
2513 //
2514 // Note: the policy evaluator (#72) returns None when the
2515 // precipitating action itself is a Takedown (its
2516 // subject_is_takendown gate fires post-precip-projection),
2517 // so `firing_rule_owned` is always None on this path —
2518 // the cascade is the only consequence of a moderator
2519 // takedown.
2520 if req.action_type == ActionType::Takedown {
2521 auto_dismiss_pendings_on_takedown(
2522 &mut tx,
2523 &subject_did,
2524 subject_uri.as_deref(),
2525 inserted_id,
2526 &req.actor_did,
2527 created_at,
2528 )
2529 .await?;
2530 }
2531
2532 // Recompute strike state for the cache. Reload history with
2533 // the new row so the cache reflects post-insert reality.
2534 // Includes the auto-recorded action when one fired.
2535 let post_history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
2536 let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
2537
2538 let post_count_i64 = post_state.current_count as i64;
2539 sqlx::query!(
2540 "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
2541 VALUES (?1, ?2, ?3, ?3)
2542 ON CONFLICT(subject_did) DO UPDATE SET
2543 current_strike_count = excluded.current_strike_count,
2544 last_action_at = excluded.last_action_at,
2545 last_recompute_at = excluded.last_recompute_at",
2546 subject_did,
2547 post_count_i64,
2548 created_at,
2549 )
2550 .execute(&mut *tx)
2551 .await?;
2552
2553 tx.commit().await?;
2554
2555 // Broadcast each emitted label to subscribeLabels consumers.
2556 // Order matches the persistence order (action label first,
2557 // then reasons). No-receivers is not a write failure
2558 // (§plan point G).
2559 for event in label_events {
2560 let _ = self.broadcast_tx.send(event);
2561 }
2562
2563 // §F23 / §A13 post-commit PDS-admin dispatch (#87). Fires
2564 // AFTER the recordAction transaction has committed and
2565 // labels have been broadcast — holding the SQLite write
2566 // lock across an HTTP call would deadlock the writer task.
2567 // Failures here log + audit-record but do NOT propagate to
2568 // the caller; the cairn-mod-side action stays committed
2569 // regardless of PDS-side outcome (per §A13's "fail loud,
2570 // let the operator decide" posture).
2571 //
2572 // duration_iso is plumbed through #89 so SuspendAccount's
2573 // dispatch path can encode duration_days in the bsky-PDS
2574 // ref field. Other action types ignore it.
2575 crate::pds_admin::dispatch::dispatch_after_record_action(
2576 self.pds_admin.as_ref(),
2577 &self.pool,
2578 crate::pds_admin::dispatch::DispatchContext {
2579 action_id: inserted_id,
2580 action_type: req.action_type,
2581 subject_did: &subject_did,
2582 reason_codes: &req.reason_codes,
2583 notes: req.notes.as_deref(),
2584 duration_iso: req.duration_iso.as_deref(),
2585 },
2586 )
2587 .await;
2588
2589 Ok(RecordedAction {
2590 action_id: inserted_id,
2591 strike_value_base: strike_base,
2592 strike_value_applied: strike_applied,
2593 was_dampened,
2594 strikes_at_time_of_action,
2595 })
2596 }
2597
2598 /// Atomic revokeAction flow (§F20 / #51). Looks up the row,
2599 /// verifies state, sets the revoked_* columns (the schema's
2600 /// no-update-except-revoke trigger permits this NULL→non-NULL
2601 /// transition), recomputes strike_state for the subject, and
2602 /// appends an audit_log row. All in one transaction.
2603 async fn handle_revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
2604 let mut tx = self.pool.begin().await?;
2605
2606 // Look up the row + check current state. Pull the fields
2607 // needed for both v1.4 revocation (subject_did, revoked_at)
2608 // and v1.5 #62 negation (subject_uri, emitted_label_uri).
2609 let row = sqlx::query!(
2610 "SELECT subject_did, subject_uri, emitted_label_uri, revoked_at
2611 FROM subject_actions WHERE id = ?1",
2612 req.action_id,
2613 )
2614 .fetch_optional(&mut *tx)
2615 .await?;
2616 let row = row.ok_or(Error::ActionNotFound(req.action_id))?;
2617 if row.revoked_at.is_some() {
2618 return Err(Error::ActionAlreadyRevoked(req.action_id));
2619 }
2620
2621 let revoked_at_ms = epoch_ms_now();
2622
2623 // UPDATE the revocation columns. The schema trigger permits
2624 // NULL→non-NULL on these three columns; any other column
2625 // change in this UPDATE would abort.
2626 sqlx::query!(
2627 "UPDATE subject_actions
2628 SET revoked_at = ?1, revoked_by_did = ?2, revoked_reason = ?3
2629 WHERE id = ?4",
2630 revoked_at_ms,
2631 req.revoked_by_did,
2632 req.revoked_reason,
2633 req.action_id,
2634 )
2635 .execute(&mut *tx)
2636 .await?;
2637
2638 // ---------- negation (#62, v1.5) ----------
2639 //
2640 // Read the action's emitted-label linkage from storage
2641 // (NOT from current policy — see module docs for the
2642 // val-from-storage rule). For each emitted val, sign and
2643 // persist a fresh label row with neg=true targeting the
2644 // same (src, uri, val) tuple. The original rows stay in
2645 // place; ATProto consumers honor the latest record per
2646 // tuple and so see the negation. Reason linkage rows
2647 // (subject_action_reason_labels) are PRESERVED — they're
2648 // forensic record of "at this point these labels were
2649 // emitted," not a cache of what's currently in force.
2650 //
2651 // exp = None on every negation. The original temp_suspension
2652 // label may have carried an expiry (its exp said "stop
2653 // honoring me at this wall-clock"), but the negation itself
2654 // is a permanent statement that supersedes the original;
2655 // expiring the negation would resurrect the original label
2656 // in consumer caches.
2657 //
2658 // Action that was never emitted (note, suppressed warning,
2659 // emission disabled at recording time) → both the action
2660 // label fetch and the reason-label fetch yield zero rows
2661 // and the negation step is a no-op. The revocation
2662 // audit row still lands with negated_labels = [].
2663 let label_uri = row
2664 .subject_uri
2665 .clone()
2666 .unwrap_or_else(|| row.subject_did.clone());
2667
2668 // Reason-label linkage rows ordered by reason_code for
2669 // deterministic audit shape. Original emission order isn't
2670 // recoverable (storage doesn't preserve req.reason_codes
2671 // ordering), and ATProto consumers don't care about
2672 // negation order; alphabetical is the cheapest stable
2673 // ordering for forensic readers.
2674 let reason_linkage = sqlx::query!(
2675 "SELECT reason_code, emitted_label_uri
2676 FROM subject_action_reason_labels
2677 WHERE action_id = ?1
2678 ORDER BY reason_code ASC",
2679 req.action_id,
2680 )
2681 .fetch_all(&mut *tx)
2682 .await?;
2683
2684 let mut negated_for_audit: Vec<serde_json::Value> = Vec::new();
2685 if let Some(ref val) = row.emitted_label_uri {
2686 negated_for_audit.push(serde_json::json!({
2687 "val": val,
2688 "uri": label_uri,
2689 }));
2690 }
2691 for r in &reason_linkage {
2692 negated_for_audit.push(serde_json::json!({
2693 "val": r.emitted_label_uri,
2694 "uri": label_uri,
2695 }));
2696 }
2697
2698 // Audit row first — captures the negation set in the hash
2699 // chain before the label INSERTs land. Mirrors #60's
2700 // emission-side ordering: a forensic reader who sees the
2701 // audit row knows what labels SHOULD have been written;
2702 // the labels table is the witness. tx atomicity guarantees
2703 // they agree.
2704 let audit_reason = build_revoke_action_audit_reason(
2705 req.action_id,
2706 req.revoked_reason.as_deref(),
2707 &negated_for_audit,
2708 );
2709 crate::audit::append::append_in_tx(
2710 &mut tx,
2711 &crate::audit::append::AuditRowForAppend {
2712 created_at: revoked_at_ms,
2713 action: "subject_action_revoked".into(),
2714 actor_did: req.revoked_by_did.clone(),
2715 target: Some(req.action_id.to_string()),
2716 target_cid: None,
2717 outcome: "success".into(),
2718 reason: Some(audit_reason),
2719 },
2720 )
2721 .await?;
2722
2723 // Sign + persist each negation label. Action label first so
2724 // its seq < reason negations'.
2725 let mut negation_events: Vec<LabelEvent> = Vec::with_capacity(1 + reason_linkage.len());
2726 if let Some(ref val) = row.emitted_label_uri {
2727 let event = self
2728 .sign_and_persist_label(
2729 &mut tx,
2730 val,
2731 &label_uri,
2732 None, // cid: not plumbed in v1.5
2733 true, // neg
2734 None, // exp: negations don't expire
2735 revoked_at_ms,
2736 )
2737 .await?;
2738 negation_events.push(event);
2739 }
2740 for r in &reason_linkage {
2741 let event = self
2742 .sign_and_persist_label(
2743 &mut tx,
2744 &r.emitted_label_uri,
2745 &label_uri,
2746 None,
2747 true,
2748 None,
2749 revoked_at_ms,
2750 )
2751 .await?;
2752 negation_events.push(event);
2753 }
2754
2755 // Recompute strike state for the subject. The decay
2756 // calculator excludes revoked rows from current_count.
2757 let post_history = load_subject_actions_for_calc(&mut tx, &row.subject_did).await?;
2758 let now_systemtime = epoch_ms_to_systemtime(revoked_at_ms);
2759 let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
2760
2761 let post_count_i64 = post_state.current_count as i64;
2762 sqlx::query!(
2763 "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
2764 VALUES (?1, ?2, ?3, ?3)
2765 ON CONFLICT(subject_did) DO UPDATE SET
2766 current_strike_count = excluded.current_strike_count,
2767 last_recompute_at = excluded.last_recompute_at",
2768 row.subject_did,
2769 post_count_i64,
2770 revoked_at_ms,
2771 )
2772 .execute(&mut *tx)
2773 .await?;
2774
2775 tx.commit().await?;
2776
2777 // Broadcast each negation to subscribeLabels consumers
2778 // post-commit. Same ordering and no-receivers-is-fine
2779 // posture as the emission path (§plan point G).
2780 for event in negation_events {
2781 let _ = self.broadcast_tx.send(event);
2782 }
2783
2784 // §F23 / §A13 post-revoke PDS-admin dispatch (#89). Looks
2785 // up the action's prior pds_admin_audit row (the takedown
2786 // / suspension call); if a successful one exists, fires
2787 // restore_account. Same fail-loud posture as the
2788 // recordAction dispatch — the cairn-mod-side revocation
2789 // stays committed regardless of PDS-side outcome.
2790 crate::pds_admin::dispatch::dispatch_after_revoke_action(
2791 self.pds_admin.as_ref(),
2792 &self.pool,
2793 crate::pds_admin::dispatch::RevokeDispatchContext {
2794 action_id: req.action_id,
2795 subject_did: &row.subject_did,
2796 revoke_reason: req.revoked_reason.as_deref(),
2797 },
2798 )
2799 .await;
2800
2801 Ok(RevokedAction {
2802 action_id: req.action_id,
2803 revoked_at: rfc3339_from_epoch_ms(revoked_at_ms)?,
2804 })
2805 }
2806
2807 /// Atomic confirmPendingAction flow (§F22 / #74). Loads the
2808 /// pending row, validates state, materializes the proposed
2809 /// action as a real `subject_actions` row (actor_kind =
2810 /// 'moderator', triggered_by_policy_rule preserved as
2811 /// provenance), emits labels, UPDATEs the pending row's
2812 /// resolution columns, recomputes strike state, and writes a
2813 /// hash-chained `pending_policy_action_confirmed` audit row —
2814 /// all in one transaction.
2815 ///
2816 /// Strike values are computed at confirmation time, not at
2817 /// proposal time: the moderator is the one taking
2818 /// responsibility, and the subject's strike state may have
2819 /// shifted since the rule fired (e.g., decay, intervening
2820 /// revocations). The resulting subject_actions row's
2821 /// `effective_at` and `expires_at` (for temp_suspension)
2822 /// likewise anchor on `now`, not on the original triggered_at.
2823 ///
2824 /// No policy-evaluation re-runs against the new row: the
2825 /// originating rule is already-fired (per #72's idempotency,
2826 /// which gates on triggered_by_policy_rule), and rule-fan-out
2827 /// from a confirmation would compound moderator-tier
2828 /// authority. The single audit row reflects this single
2829 /// decision.
2830 async fn handle_confirm_pending_action(
2831 &self,
2832 req: ConfirmPendingActionRequest,
2833 ) -> Result<ConfirmedPendingAction> {
2834 let mut tx = self.pool.begin().await?;
2835
2836 // ---------- load + validate pending ----------
2837 let pending = sqlx::query!(
2838 r#"SELECT
2839 subject_did AS "subject_did!: String",
2840 subject_uri,
2841 action_type AS "action_type!: String",
2842 duration_ms,
2843 reason_codes AS "reason_codes!: String",
2844 triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
2845 resolution
2846 FROM pending_policy_actions
2847 WHERE id = ?1"#,
2848 req.pending_id,
2849 )
2850 .fetch_optional(&mut *tx)
2851 .await?;
2852 let pending = pending.ok_or(Error::PendingActionNotFound(req.pending_id))?;
2853 if pending.resolution.is_some() {
2854 return Err(Error::PendingAlreadyResolved(req.pending_id));
2855 }
2856
2857 let action_type = ActionType::from_db_str(&pending.action_type).ok_or_else(|| {
2858 Error::Signing(format!(
2859 "pending_policy_actions row has invalid action_type {:?}",
2860 pending.action_type
2861 ))
2862 })?;
2863
2864 // ---------- defensive: subject_takendown ----------
2865 //
2866 // #76 (auto-dismissal-on-takedown) is meant to resolve all
2867 // pendings the moment a takedown lands — so reaching this
2868 // path with an active takedown means a race the auto-dismiss
2869 // hasn't closed yet. Reject defensively rather than
2870 // materialize a redundant action under terminal-severity
2871 // semantics.
2872 let history = load_subject_actions_for_calc(&mut tx, &pending.subject_did).await?;
2873 let subject_is_takendown = history
2874 .iter()
2875 .any(|a| a.action_type == ActionType::Takedown && a.revoked_at.is_none());
2876 if subject_is_takendown {
2877 return Err(Error::SubjectTakendown(pending.subject_did.clone()));
2878 }
2879
2880 // ---------- decode pending fields ----------
2881 let reason_codes: Vec<String> = serde_json::from_str(&pending.reason_codes)
2882 .map_err(|e| Error::Signing(format!("parse pending reason_codes: {e}")))?;
2883 if reason_codes.is_empty() {
2884 return Err(Error::Signing(
2885 "confirmPendingAction: pending row has empty reason_codes".into(),
2886 ));
2887 }
2888
2889 // Vocabulary lookup. The pending was originally created
2890 // with rule.reason_codes (config-validated against
2891 // [moderation_reasons] at startup per #71), so this should
2892 // always resolve — but the vocabulary may have shifted
2893 // since the pending was queued, so route the failure
2894 // through the same ReasonNotFound surface the moderator
2895 // path uses.
2896 let primary = resolve_primary_reason(&reason_codes, &self.reason_vocabulary)?;
2897
2898 // ---------- strike calc at confirmation time ----------
2899 let now_ms = epoch_ms_now();
2900 let now_systemtime = epoch_ms_to_systemtime(now_ms);
2901 let pre_state = calculate_strike_state(&history, &self.strike_policy, now_systemtime);
2902 let strikes_at_time_of_action = pre_state.current_count;
2903 let position = compute_position_in_window(&history, &self.strike_policy, now_systemtime);
2904 let calc: StrikeApplication = strike_calculate(
2905 strikes_at_time_of_action,
2906 &primary,
2907 &self.strike_policy,
2908 position,
2909 );
2910 let (strike_base, strike_applied, was_dampened) = match action_type {
2911 ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
2912 _ => (calc.base_weight, calc.applied, calc.was_dampened),
2913 };
2914
2915 // ---------- effective_at / expires_at ----------
2916 //
2917 // Confirmation is when the action takes effect, so
2918 // expires_at re-anchors on `now`. duration_ms was frozen
2919 // on the pending row at proposal time (rule.duration via
2920 // #71 / #73); the canonical wire shape is `PT<seconds>S`
2921 // matching the policy auto-action path.
2922 let effective_at = now_ms;
2923 let expires_at: Option<i64> = pending.duration_ms.map(|d| effective_at + d);
2924 let duration_iso: Option<String> = pending.duration_ms.map(|d| format!("PT{}S", d / 1000));
2925
2926 // ---------- emission drafts ----------
2927 let action_for_emission = ActionForEmission {
2928 action_type,
2929 expires_at: expires_at.map(epoch_ms_to_systemtime),
2930 subject_did: pending.subject_did.clone(),
2931 subject_uri: pending.subject_uri.clone(),
2932 reason_codes: reason_codes.clone(),
2933 cid: None,
2934 };
2935 let action_drafts = resolve_action_labels(
2936 &action_for_emission,
2937 &self.label_emission_policy,
2938 now_systemtime,
2939 );
2940 let reason_drafts = resolve_reason_labels(
2941 &action_for_emission,
2942 &self.label_emission_policy,
2943 now_systemtime,
2944 );
2945 let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
2946 .iter()
2947 .chain(reason_drafts.iter())
2948 .map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
2949 .collect();
2950
2951 // ---------- predict + audit ----------
2952 let next_action_id = predict_next_subject_action_id(&mut tx).await?;
2953
2954 let audit_reason = build_pending_confirmed_audit_reason(
2955 req.pending_id,
2956 &pending.triggered_by_policy_rule,
2957 next_action_id,
2958 action_type,
2959 &primary.identifier,
2960 &reason_codes,
2961 strike_base,
2962 strike_applied,
2963 was_dampened,
2964 &emitted_labels_for_audit,
2965 req.note.as_deref(),
2966 );
2967 let audit_target = pending
2968 .subject_uri
2969 .clone()
2970 .unwrap_or_else(|| pending.subject_did.clone());
2971 let audit_log_id = crate::audit::append::append_in_tx(
2972 &mut tx,
2973 &crate::audit::append::AuditRowForAppend {
2974 created_at: now_ms,
2975 action: "pending_policy_action_confirmed".into(),
2976 actor_did: req.moderator_did.clone(),
2977 target: Some(audit_target),
2978 target_cid: None,
2979 outcome: "success".into(),
2980 reason: Some(audit_reason),
2981 },
2982 )
2983 .await?;
2984
2985 // ---------- INSERT subject_actions ----------
2986 //
2987 // actor_kind = 'moderator' (the moderator takes
2988 // responsibility by confirming); triggered_by_policy_rule
2989 // preserves the rule that proposed the action — forensic
2990 // provenance plus the idempotency-gate input for #72's
2991 // already-fired check.
2992 let action_type_str = action_type.as_db_str();
2993 let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
2994 let strikes_at_time_i64 = strikes_at_time_of_action as i64;
2995 let strike_base_i64 = strike_base as i64;
2996 let strike_applied_i64 = strike_applied as i64;
2997 let reason_codes_json = serde_json::to_string(&reason_codes)
2998 .map_err(|e| Error::Signing(format!("serialize reason_codes for confirm: {e}")))?;
2999 let actor_kind_moderator = "moderator";
3000 let inserted_id = sqlx::query_scalar!(
3001 "INSERT INTO subject_actions (
3002 subject_did, subject_uri, actor_did, action_type, reason_codes,
3003 duration, effective_at, expires_at, notes, report_ids,
3004 strike_value_base, strike_value_applied, was_dampened,
3005 strikes_at_time_of_action, audit_log_id, created_at,
3006 actor_kind, triggered_by_policy_rule
3007 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17)
3008 RETURNING id",
3009 pending.subject_did,
3010 pending.subject_uri,
3011 req.moderator_did,
3012 action_type_str,
3013 reason_codes_json,
3014 duration_iso,
3015 effective_at,
3016 expires_at,
3017 req.note,
3018 strike_base_i64,
3019 strike_applied_i64,
3020 was_dampened_int,
3021 strikes_at_time_i64,
3022 audit_log_id,
3023 now_ms,
3024 actor_kind_moderator,
3025 pending.triggered_by_policy_rule,
3026 )
3027 .fetch_one(&mut *tx)
3028 .await?;
3029 if inserted_id != next_action_id {
3030 return Err(Error::Signing(format!(
3031 "confirm subject_actions inserted at id {inserted_id} but predicted \
3032 {next_action_id}; audit chain captured the predicted id (corrupted)"
3033 )));
3034 }
3035
3036 // ---------- emit labels ----------
3037 //
3038 // Same pattern as handle_record_action and
3039 // insert_policy_auto_action: action label first (its seq
3040 // < reason labels'), then per-(action, reason) linkage
3041 // rows. Fresh INSERT so the #64 idempotency guard is a
3042 // no-op; skipped here for clarity.
3043 let mut label_events: Vec<LabelEvent> =
3044 Vec::with_capacity(action_drafts.len() + reason_drafts.len());
3045 let mut action_label_val: Option<String> = None;
3046 for draft in &action_drafts {
3047 let event = self
3048 .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
3049 .await?;
3050 if action_label_val.is_none() {
3051 action_label_val = Some(draft.val.clone());
3052 }
3053 label_events.push(event);
3054 }
3055 if let Some(ref val) = action_label_val {
3056 sqlx::query!(
3057 "UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
3058 val,
3059 inserted_id,
3060 )
3061 .execute(&mut *tx)
3062 .await?;
3063 }
3064 for (draft, reason_code) in reason_drafts.iter().zip(reason_codes.iter()) {
3065 let event = self
3066 .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
3067 .await?;
3068 sqlx::query!(
3069 "INSERT INTO subject_action_reason_labels
3070 (action_id, reason_code, emitted_label_uri, emitted_at)
3071 VALUES (?1, ?2, ?3, ?4)",
3072 inserted_id,
3073 reason_code,
3074 draft.val,
3075 effective_at,
3076 )
3077 .execute(&mut *tx)
3078 .await?;
3079 label_events.push(event);
3080 }
3081
3082 // ---------- UPDATE pending: NULL → 'confirmed' ----------
3083 //
3084 // The schema trigger (migration 0005) permits this single
3085 // NULL → non-NULL transition on resolution / resolved_at /
3086 // resolved_by_did / confirmed_action_id; any other change
3087 // to the row would abort. Setting all four columns
3088 // together so the row's invariant
3089 // ("confirmed implies confirmed_action_id non-NULL")
3090 // holds at every queryable instant.
3091 let confirmed_resolution = "confirmed";
3092 sqlx::query!(
3093 "UPDATE pending_policy_actions
3094 SET resolution = ?1,
3095 resolved_at = ?2,
3096 resolved_by_did = ?3,
3097 confirmed_action_id = ?4
3098 WHERE id = ?5",
3099 confirmed_resolution,
3100 now_ms,
3101 req.moderator_did,
3102 inserted_id,
3103 req.pending_id,
3104 )
3105 .execute(&mut *tx)
3106 .await?;
3107
3108 // ---------- takedown cascade (#76, v1.6) ----------
3109 //
3110 // Confirming a pending whose proposed action is a Takedown
3111 // makes the subject takendown the moment the new
3112 // subject_actions row INSERTed; the cascade then dismisses
3113 // every OTHER unresolved pending for the subject. The
3114 // pending currently being confirmed is naturally excluded
3115 // because its UPDATE to resolution='confirmed' has already
3116 // landed above — the cascade's `WHERE resolution IS NULL`
3117 // filter sees it as 'confirmed', not NULL.
3118 //
3119 // resolved_by_did on cascaded rows is the confirming
3120 // moderator's DID — not the synthetic policy DID — because
3121 // the moderator chose to confirm the takedown and is the
3122 // ultimate authority for the cascade.
3123 if action_type == ActionType::Takedown {
3124 auto_dismiss_pendings_on_takedown(
3125 &mut tx,
3126 &pending.subject_did,
3127 pending.subject_uri.as_deref(),
3128 inserted_id,
3129 &req.moderator_did,
3130 now_ms,
3131 )
3132 .await?;
3133 }
3134
3135 // ---------- recompute strike state ----------
3136 let post_history = load_subject_actions_for_calc(&mut tx, &pending.subject_did).await?;
3137 let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
3138 let post_count_i64 = post_state.current_count as i64;
3139 sqlx::query!(
3140 "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
3141 VALUES (?1, ?2, ?3, ?3)
3142 ON CONFLICT(subject_did) DO UPDATE SET
3143 current_strike_count = excluded.current_strike_count,
3144 last_action_at = excluded.last_action_at,
3145 last_recompute_at = excluded.last_recompute_at",
3146 pending.subject_did,
3147 post_count_i64,
3148 now_ms,
3149 )
3150 .execute(&mut *tx)
3151 .await?;
3152
3153 tx.commit().await?;
3154
3155 // ---------- broadcast ----------
3156 for event in label_events {
3157 let _ = self.broadcast_tx.send(event);
3158 }
3159
3160 Ok(ConfirmedPendingAction {
3161 action_id: inserted_id,
3162 pending_id: req.pending_id,
3163 resolved_at: rfc3339_from_epoch_ms(now_ms)?,
3164 })
3165 }
3166
3167 /// Atomic dismissPendingAction flow (§F22 / #75). Loads the
3168 /// pending row, validates state, UPDATEs the resolution
3169 /// columns to 'dismissed', and writes a hash-chained
3170 /// `pending_policy_action_dismissed` audit row — all in one
3171 /// transaction. Inverse of `handle_confirm_pending_action`:
3172 /// no subject_actions row, no label emission, no strike-state
3173 /// recompute. The pending row stays in the table as forensic
3174 /// record (`confirmed_action_id` stays NULL; only confirmed
3175 /// pendings link forward to a materialized action).
3176 ///
3177 /// No `SubjectTakendown` defensive check (unlike confirm in
3178 /// #74): explicit dismissal is meaningful regardless of
3179 /// takedown state, and is in fact part of the cleanup path
3180 /// #76 will automate when a takedown lands.
3181 async fn handle_dismiss_pending_action(
3182 &self,
3183 req: DismissPendingActionRequest,
3184 ) -> Result<DismissedPendingAction> {
3185 let mut tx = self.pool.begin().await?;
3186
3187 // ---------- load + validate pending ----------
3188 let pending = sqlx::query!(
3189 r#"SELECT
3190 subject_did AS "subject_did!: String",
3191 subject_uri,
3192 action_type AS "action_type!: String",
3193 reason_codes AS "reason_codes!: String",
3194 triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
3195 resolution
3196 FROM pending_policy_actions
3197 WHERE id = ?1"#,
3198 req.pending_id,
3199 )
3200 .fetch_optional(&mut *tx)
3201 .await?;
3202 let pending = pending.ok_or(Error::PendingActionNotFound(req.pending_id))?;
3203 if pending.resolution.is_some() {
3204 return Err(Error::PendingAlreadyResolved(req.pending_id));
3205 }
3206
3207 let now_ms = epoch_ms_now();
3208
3209 // ---------- audit row ----------
3210 //
3211 // Reason JSON cross-references the pending row and echoes
3212 // the proposed-action shape (action_type, reason_codes)
3213 // so a forensic reader can reconstruct "moderator
3214 // dismissed a proposed indef_suspension for spam"
3215 // without joining to pending_policy_actions. subject_did
3216 // / actor_did are already on the audit_log row's target /
3217 // actor_did — not duplicated in the reason JSON.
3218 let reason_codes_for_audit: serde_json::Value =
3219 serde_json::from_str(&pending.reason_codes).unwrap_or(serde_json::Value::Null);
3220 let audit_reason = build_pending_dismissed_audit_reason(
3221 req.pending_id,
3222 &pending.triggered_by_policy_rule,
3223 &pending.action_type,
3224 reason_codes_for_audit,
3225 req.reason.as_deref(),
3226 );
3227 let audit_target = pending
3228 .subject_uri
3229 .clone()
3230 .unwrap_or_else(|| pending.subject_did.clone());
3231 crate::audit::append::append_in_tx(
3232 &mut tx,
3233 &crate::audit::append::AuditRowForAppend {
3234 created_at: now_ms,
3235 action: "pending_policy_action_dismissed".into(),
3236 actor_did: req.moderator_did.clone(),
3237 target: Some(audit_target),
3238 target_cid: None,
3239 outcome: "success".into(),
3240 reason: Some(audit_reason),
3241 },
3242 )
3243 .await?;
3244
3245 // ---------- UPDATE pending: NULL → 'dismissed' ----------
3246 //
3247 // The schema trigger (migration 0005) permits this single
3248 // NULL → non-NULL transition on resolution / resolved_at
3249 // / resolved_by_did; confirmed_action_id stays NULL by
3250 // design (no action was created). Setting all three
3251 // columns together so the row's "dismissed implies who
3252 // and when" invariant holds at every queryable instant.
3253 let dismissed_resolution = "dismissed";
3254 sqlx::query!(
3255 "UPDATE pending_policy_actions
3256 SET resolution = ?1,
3257 resolved_at = ?2,
3258 resolved_by_did = ?3
3259 WHERE id = ?4",
3260 dismissed_resolution,
3261 now_ms,
3262 req.moderator_did,
3263 req.pending_id,
3264 )
3265 .execute(&mut *tx)
3266 .await?;
3267
3268 tx.commit().await?;
3269
3270 Ok(DismissedPendingAction {
3271 pending_id: req.pending_id,
3272 resolved_at: rfc3339_from_epoch_ms(now_ms)?,
3273 })
3274 }
3275
3276 async fn heartbeat(&self) -> Result<()> {
3277 let now_ms = epoch_ms_now();
3278 sqlx::query!(
3279 "UPDATE server_instance_lease SET last_heartbeat = ?1
3280 WHERE id = 1 AND instance_id = ?2",
3281 now_ms,
3282 self.instance_id,
3283 )
3284 .execute(&self.pool)
3285 .await?;
3286 Ok(())
3287 }
3288
3289 async fn release_lease(&self) -> Result<()> {
3290 sqlx::query!(
3291 "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
3292 self.instance_id,
3293 )
3294 .execute(&self.pool)
3295 .await?;
3296 Ok(())
3297 }
3298}
3299
3300// ---------- Lease + key bootstrap ----------
3301
3302/// Acquire the §F5 single-instance lease.
3303///
3304/// Returns the new instance id on success. Returns
3305/// [`Error::LeaseHeld`] when an existing row's `last_heartbeat` is
3306/// still within [`LEASE_STALE_MS`] — i.e., another writer is live.
3307///
3308/// Exposed as `pub(crate)` so one-shot CLI tools that mutate
3309/// `audit_log` (`cairn audit-rebuild`, #40) can assert the same
3310/// "no live writer" invariant the long-running serve process holds.
3311/// The CLI tool releases via [`release_lease_by_id`] when done.
3312pub(crate) async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
3313 let now_ms = epoch_ms_now();
3314
3315 let existing =
3316 sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
3317 .fetch_optional(pool)
3318 .await?;
3319
3320 if let Some(row) = existing {
3321 let age_ms = (now_ms - row.last_heartbeat).max(0);
3322 if age_ms < LEASE_STALE_MS {
3323 return Err(Error::LeaseHeld {
3324 instance_id: row.instance_id,
3325 age_secs: (age_ms / 1000) as u64,
3326 });
3327 }
3328 }
3329
3330 let new_id = Uuid::new_v4().to_string();
3331 // INSERT OR REPLACE on id=1: takes over a stale lease row if present,
3332 // or creates the row on first startup. Legal because the staleness
3333 // check above confirms no live peer.
3334 sqlx::query!(
3335 "INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
3336 VALUES (1, ?1, ?2, ?2)
3337 ON CONFLICT(id) DO UPDATE SET
3338 instance_id = excluded.instance_id,
3339 acquired_at = excluded.acquired_at,
3340 last_heartbeat = excluded.last_heartbeat",
3341 new_id,
3342 now_ms,
3343 )
3344 .execute(pool)
3345 .await?;
3346
3347 Ok(new_id)
3348}
3349
3350/// Release a lease acquired via [`acquire_lease`] by its instance id.
3351/// Free-function counterpart to [`Writer::release_lease`] for callers
3352/// that don't own a `Writer` (one-shot CLI tools that just hold the
3353/// lease for the duration of a data migration).
3354pub(crate) async fn release_lease_by_id(pool: &Pool<Sqlite>, instance_id: &str) -> Result<()> {
3355 sqlx::query!(
3356 "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
3357 instance_id,
3358 )
3359 .execute(pool)
3360 .await?;
3361 Ok(())
3362}
3363
3364async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
3365 let kp = K256Keypair::from_private_key(key.expose_secret())?;
3366 let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
3367
3368 let existing =
3369 sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
3370 .fetch_optional(pool)
3371 .await?;
3372
3373 if let Some(row) = existing {
3374 if row.public_key_multibase != my_multibase {
3375 return Err(Error::Signing(format!(
3376 "signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
3377 row.public_key_multibase, my_multibase
3378 )));
3379 }
3380 return Ok(row.id);
3381 }
3382
3383 let now_ms = epoch_ms_now();
3384 let valid_from = rfc3339_from_epoch_ms(now_ms)?;
3385 let id = sqlx::query_scalar!(
3386 "INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
3387 VALUES (?1, ?2, NULL, ?3)
3388 RETURNING id",
3389 my_multibase,
3390 valid_from,
3391 now_ms,
3392 )
3393 .fetch_one(pool)
3394 .await?;
3395
3396 Ok(id)
3397}
3398
3399// ---------- Pure helpers (unit-testable without a DB) ----------
3400
3401async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
3402 // label_sequence has a single auto-increment column; `DEFAULT VALUES`
3403 // triggers a fresh allocation. Under AUTOINCREMENT the assigned seq
3404 // is never reused even across rolled-back transactions — this is what
3405 // gives us "strictly monotonic, no gaps from writer's perspective."
3406 let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
3407 .fetch_one(tx)
3408 .await?;
3409 Ok(seq)
3410}
3411
3412/// §6.1 monotonicity clamp. `prev_cts_str`, when present, is expected to
3413/// be in the writer's canonical `CTS_FORMAT`.
3414fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
3415 let effective_ms = match prev_cts_str {
3416 Some(s) => {
3417 let prev_ms = parse_rfc3339_ms(s)?;
3418 wall_now_ms.max(prev_ms + 1)
3419 }
3420 None => wall_now_ms,
3421 };
3422 rfc3339_from_epoch_ms(effective_ms)
3423}
3424
3425/// Epoch-ms → RFC-3339 Z with millisecond precision. `pub(crate)` so
3426/// peer modules (admin/audit_view, future retention sweep) share the
3427/// single formatter used on the writer's timestamp boundary.
3428pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
3429 let nanos: i128 = (ms as i128) * 1_000_000;
3430 let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
3431 .map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
3432 let formatted = dt
3433 .format(&CTS_FORMAT)
3434 .map_err(|e| Error::Signing(format!("format cts: {e}")))?;
3435 Ok(format!("{formatted}Z"))
3436}
3437
3438/// RFC-3339 Z with millisecond precision → epoch-ms. `pub(crate)` so
3439/// admin handlers can validate `since` / `until` query params against
3440/// the same parser the writer uses on its input boundary.
3441pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
3442 let stripped = s
3443 .strip_suffix('Z')
3444 .ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
3445 let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
3446 .map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
3447 let nanos = pdt.assume_utc().unix_timestamp_nanos();
3448 Ok((nanos / 1_000_000) as i64)
3449}
3450
3451/// Current wall-clock time as Unix epoch milliseconds. `pub(crate)` so
3452/// peer modules (server, future retention sweep) share one implementation
3453/// rather than each inlining `SystemTime::now()` with its own error
3454/// handling.
3455pub(crate) fn epoch_ms_now() -> i64 {
3456 SystemTime::now()
3457 .duration_since(UNIX_EPOCH)
3458 .expect("system clock before unix epoch")
3459 .as_millis() as i64
3460}
3461
3462fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
3463 let body = serde_json::json!({
3464 "val": val,
3465 "neg": neg,
3466 "moderator_reason": moderator_reason,
3467 });
3468 body.to_string()
3469}
3470
3471/// `report_resolved` audit reason JSON — see [`AUDIT_REASON_RESOLVE_REPORT`]
3472/// for the schema.
3473fn build_resolve_audit_reason(
3474 applied_label_val: Option<&str>,
3475 resolution_reason: Option<&str>,
3476) -> String {
3477 serde_json::json!({
3478 "applied_label_val": applied_label_val,
3479 "resolution_reason": resolution_reason,
3480 })
3481 .to_string()
3482}
3483
3484/// `retention_sweep` audit reason JSON — see
3485/// [`AUDIT_REASON_RETENTION_SWEEP`] for the schema. Shared with the
3486/// retentionSweep admin handler (audited only on the operator-
3487/// initiated path per Q6/D2).
3488pub(crate) fn build_retention_sweep_audit_reason(result: &SweepResult) -> String {
3489 serde_json::json!({
3490 "rows_deleted": result.rows_deleted,
3491 "batches": result.batches,
3492 "duration_ms": result.duration_ms,
3493 "retention_days_applied": result.retention_days_applied,
3494 })
3495 .to_string()
3496}
3497
3498/// Audit-log `reason` JSON schema for `subject_action_recorded`
3499/// (§F20 / #51 graduated-action moderation). Captures the action_id,
3500/// action_type, primary reason resolved at record time, the full
3501/// list of declared reasons, the strike value applied, and whether
3502/// dampening fired — enough for forensic reconstruction without
3503/// requiring a join to subject_actions.
3504///
3505/// ```json
3506/// {
3507/// "action_id": <i64>,
3508/// "action_type": "<warning|note|temp_suspension|indef_suspension|takedown>",
3509/// "primary_reason": "<identifier>",
3510/// "reason_codes": ["<id1>", "<id2>"],
3511/// "strike_value_base": <u32>,
3512/// "strike_value_applied": <u32>,
3513/// "was_dampened": <bool>
3514/// }
3515/// ```
3516#[doc(alias = "audit_log.reason.subject_action_recorded")]
3517pub 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 }";
3518
3519/// Audit-log `reason` JSON schema for `subject_action_revoked`
3520/// (§F20 / #51). Captures the revoked subject_actions row id and
3521/// the moderator-supplied rationale.
3522///
3523/// ```json
3524/// {
3525/// "action_id": <i64>,
3526/// "revoked_reason": "<free text>" | null
3527/// }
3528/// ```
3529#[doc(alias = "audit_log.reason.subject_action_revoked")]
3530pub const AUDIT_REASON_REVOKE_ACTION: &str =
3531 "subject_action_revoked: { action_id, revoked_reason }";
3532
3533/// Audit-log `reason` JSON schema for `pending_policy_action_confirmed`
3534/// (§F22 / #74). Single audit row commits with the
3535/// pending → confirmed transition: the new subject_actions INSERT,
3536/// label emission, the `pending_policy_actions` resolution UPDATE,
3537/// and the strike-state cache UPSERT all hash-chain together.
3538///
3539/// ```json
3540/// {
3541/// "pending_id": <i64>,
3542/// "action_id": <i64>,
3543/// "triggered_by_policy_rule": "<rule name>",
3544/// "action_type": "<warning|note|temp_suspension|indef_suspension|takedown>",
3545/// "primary_reason": "<identifier>",
3546/// "reason_codes": ["<id1>", "<id2>"],
3547/// "strike_value_base": <u32>,
3548/// "strike_value_applied": <u32>,
3549/// "was_dampened": <bool>,
3550/// "emitted_labels": [{"val": "<v>", "uri": "<u>"}, ...],
3551/// "moderator_note": "<free text>" | null
3552/// }
3553/// ```
3554#[doc(alias = "audit_log.reason.pending_policy_action_confirmed")]
3555pub 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 }";
3556
3557/// Audit-log `reason` JSON schema for `pending_policy_action_dismissed`
3558/// (§F22 / #75 + #76). Single audit row per dismissed pending,
3559/// hash-chained alongside any other writes in the transaction. No
3560/// subject_actions row is created and no labels emit on the
3561/// dismiss side; the dismissed pending stays in the table as
3562/// forensic record.
3563///
3564/// Two shapes share the same audit_log.action value
3565/// (`pending_policy_action_dismissed`), discriminated by the
3566/// `triggered_by` field in the reason JSON:
3567///
3568/// - `triggered_by = "moderator_dismissed"` (#75): explicit
3569/// moderator dismissal via `tools.cairn.admin.dismissPendingAction`.
3570/// Carries the moderator's optional rationale on
3571/// `moderator_reason`.
3572///
3573/// - `triggered_by = "takedown_terminal"` (#76): automatic cascade
3574/// when a takedown row INSERTs against the subject. Carries
3575/// `takedown_action_id` cross-referencing the triggering
3576/// `subject_actions` row; `moderator_reason` is omitted.
3577///
3578/// ```json
3579/// {
3580/// "pending_id": <i64>,
3581/// "triggered_by_policy_rule": "<rule name>",
3582/// "action_type": "<warning|note|temp_suspension|indef_suspension|takedown>",
3583/// "reason_codes": ["<id1>", "<id2>"],
3584/// "triggered_by": "moderator_dismissed" | "takedown_terminal",
3585/// "moderator_reason": "<free text>" | null, // moderator path only
3586/// "takedown_action_id": <i64> // takedown path only
3587/// }
3588/// ```
3589#[doc(alias = "audit_log.reason.pending_policy_action_dismissed")]
3590pub 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 }";
3591
3592/// `reporter_flagged` / `reporter_unflagged` audit reason JSON —
3593/// see [`AUDIT_REASON_FLAG_REPORTER`]. Shared with the flagReporter
3594/// handler (not the writer — flagReporter is a direct handler txn
3595/// per the #15 criteria).
3596pub(crate) fn build_flag_reporter_audit_reason(
3597 did: &str,
3598 suppressed: bool,
3599 moderator_reason: Option<&str>,
3600) -> String {
3601 serde_json::json!({
3602 "did": did,
3603 "suppressed": suppressed,
3604 "moderator_reason": moderator_reason,
3605 })
3606 .to_string()
3607}
3608
3609/// Idempotency guard for the recorder's action-label emission
3610/// step (#64). Returns `true` when a row already carries an
3611/// emitted action-label `val`, signalling the emission loop
3612/// should skip producing a duplicate (src, uri, val) record.
3613///
3614/// v1.5's [`Writer::handle_record_action`] always passes `None`
3615/// here in production (the row was just INSERTed and nothing
3616/// else has touched it inside the same transaction). The
3617/// function exists so future paths that operate on existing
3618/// rows — backfill migrations, retry helpers — can call into
3619/// the same gate logic, and so the gate's behavior is testable
3620/// without DB scaffolding.
3621fn should_skip_action_label_emission(existing_emitted_label_uri: Option<&str>) -> bool {
3622 existing_emitted_label_uri.is_some()
3623}
3624
3625/// Idempotency guard for the recorder's reason-label emission
3626/// step (#64). Returns `true` when a `subject_action_reason_labels`
3627/// row already exists for `(action_id, reason_code)`, signalling
3628/// the per-reason emission loop should skip this draft.
3629///
3630/// Same v1.5 production posture as
3631/// [`should_skip_action_label_emission`]: `existing` is always
3632/// empty when called from [`Writer::handle_record_action`]; the
3633/// function is the gate logic future paths can reuse and tests
3634/// can exercise.
3635fn should_skip_reason_emission(reason_code: &str, existing: &HashSet<&str>) -> bool {
3636 existing.contains(reason_code)
3637}
3638
3639/// `subject_action_recorded` audit reason JSON — see
3640/// [`AUDIT_REASON_RECORD_ACTION`] for the schema.
3641#[allow(clippy::too_many_arguments)]
3642fn build_record_action_audit_reason(
3643 action_id: i64,
3644 action_type: ActionType,
3645 primary_reason: &str,
3646 reason_codes: &[String],
3647 strike_value_base: u32,
3648 strike_value_applied: u32,
3649 was_dampened: bool,
3650 emitted_labels: &[serde_json::Value],
3651 actor_kind: &str,
3652 triggered_by_policy_rule: Option<&str>,
3653 policy_consequence: Option<serde_json::Value>,
3654) -> String {
3655 let mut obj = serde_json::json!({
3656 "action_id": action_id,
3657 "action_type": action_type.as_db_str(),
3658 "primary_reason": primary_reason,
3659 "reason_codes": reason_codes,
3660 "strike_value_base": strike_value_base,
3661 "strike_value_applied": strike_value_applied,
3662 "was_dampened": was_dampened,
3663 // v1.5 (#60): every emitted label as `{val, uri}`. Empty
3664 // when [label_emission].enabled = false, when the action
3665 // type is `note`, or when `warning` with the default
3666 // suppression gate. Captured pre-INSERT so the audit
3667 // hash chain locks the entire (action, labels) bundle
3668 // even though the actual label INSERTs land later in the
3669 // same tx.
3670 "emitted_labels": emitted_labels,
3671 // v1.6 (#73): provenance discriminator for the recorder
3672 // (moderator vs policy engine). Mirrors
3673 // subject_actions.actor_kind.
3674 "actor_kind": actor_kind,
3675 });
3676 if let Some(rule_name) = triggered_by_policy_rule {
3677 obj["triggered_by_policy_rule"] = serde_json::Value::String(rule_name.to_string());
3678 }
3679 // v1.6 (#73): when the recorded action crosses a
3680 // [policy_automation] threshold, the audit row's reason JSON
3681 // gains a `policy_consequence` field cross-referencing the
3682 // resulting auto-recorded action OR pending row. Omitted
3683 // when no rule fires.
3684 if let Some(pc) = policy_consequence {
3685 obj["policy_consequence"] = pc;
3686 }
3687 obj.to_string()
3688}
3689
3690/// `subject_action_revoked` audit reason JSON — see
3691/// [`AUDIT_REASON_REVOKE_ACTION`] for the schema.
3692///
3693/// `negated_labels` mirrors the emission-side `emitted_labels`
3694/// shape from [`build_record_action_audit_reason`]: an array of
3695/// `{val, uri}` entries, action label first followed by reason
3696/// labels in alphabetical order. Empty when the revoked action
3697/// had no emitted labels (note, suppressed warning, emission
3698/// disabled at record time, action recorded pre-v1.5).
3699fn build_revoke_action_audit_reason(
3700 action_id: i64,
3701 revoked_reason: Option<&str>,
3702 negated_labels: &[serde_json::Value],
3703) -> String {
3704 serde_json::json!({
3705 "action_id": action_id,
3706 "revoked_reason": revoked_reason,
3707 "negated_labels": negated_labels,
3708 })
3709 .to_string()
3710}
3711
3712/// `pending_policy_action_confirmed` audit reason JSON — see
3713/// [`AUDIT_REASON_CONFIRM_PENDING_ACTION`] for the schema. Cross-
3714/// references the originating pending row (`pending_id`,
3715/// `triggered_by_policy_rule`) and the materialized action
3716/// (`action_id`); echoes the moderator's optional `note` so
3717/// forensic readers can reconstruct the rationale without joining
3718/// to subject_actions.notes.
3719#[allow(clippy::too_many_arguments)]
3720fn build_pending_confirmed_audit_reason(
3721 pending_id: i64,
3722 triggered_by_policy_rule: &str,
3723 action_id: i64,
3724 action_type: ActionType,
3725 primary_reason: &str,
3726 reason_codes: &[String],
3727 strike_value_base: u32,
3728 strike_value_applied: u32,
3729 was_dampened: bool,
3730 emitted_labels: &[serde_json::Value],
3731 moderator_note: Option<&str>,
3732) -> String {
3733 serde_json::json!({
3734 "pending_id": pending_id,
3735 "action_id": action_id,
3736 "triggered_by_policy_rule": triggered_by_policy_rule,
3737 "action_type": action_type.as_db_str(),
3738 "primary_reason": primary_reason,
3739 "reason_codes": reason_codes,
3740 "strike_value_base": strike_value_base,
3741 "strike_value_applied": strike_value_applied,
3742 "was_dampened": was_dampened,
3743 "emitted_labels": emitted_labels,
3744 "moderator_note": moderator_note,
3745 })
3746 .to_string()
3747}
3748
3749/// `pending_policy_action_dismissed` audit reason JSON for the
3750/// explicit moderator-dismiss path (#75). Cross-references the
3751/// originating pending row (`pending_id`, `triggered_by_policy_rule`)
3752/// and echoes the proposed-action shape so forensic readers can
3753/// reconstruct "moderator dismissed proposed action X" without
3754/// joining back to `pending_policy_actions`. The `triggered_by`
3755/// field discriminates this shape from the takedown-cascade
3756/// dismissal (#76) — same audit_log.action value
3757/// (`pending_policy_action_dismissed`), different reason JSON
3758/// shapes.
3759///
3760/// `reason_codes_json` is the parsed JSON value from the pending
3761/// row's `reason_codes` column (already a JSON array of strings);
3762/// callers pass [`serde_json::Value::Null`] when the column is
3763/// somehow unparseable, which preserves the audit row rather
3764/// than failing the dismissal.
3765fn build_pending_dismissed_audit_reason(
3766 pending_id: i64,
3767 triggered_by_policy_rule: &str,
3768 action_type: &str,
3769 reason_codes_json: serde_json::Value,
3770 moderator_reason: Option<&str>,
3771) -> String {
3772 serde_json::json!({
3773 "pending_id": pending_id,
3774 "triggered_by_policy_rule": triggered_by_policy_rule,
3775 "action_type": action_type,
3776 "reason_codes": reason_codes_json,
3777 "moderator_reason": moderator_reason,
3778 "triggered_by": "moderator_dismissed",
3779 })
3780 .to_string()
3781}
3782
3783/// `pending_policy_action_dismissed` audit reason JSON for the
3784/// takedown-cascade auto-dismissal path (§F22 / #76). Same
3785/// audit_log.action value as the explicit moderator dismiss
3786/// (#75), discriminated via the `triggered_by` field. Cross-
3787/// references both the originating pending row (`pending_id`,
3788/// `triggered_by_policy_rule`) and the takedown that caused the
3789/// cascade (`takedown_action_id`); a forensic reader can walk
3790/// either direction.
3791///
3792/// `reason_codes_json` is the parsed JSON value from the pending
3793/// row's `reason_codes` column (already a JSON array of strings);
3794/// callers pass [`serde_json::Value::Null`] when the column is
3795/// somehow unparseable, which preserves the audit row rather
3796/// than failing the cascade.
3797fn build_pending_dismissed_on_takedown_audit_reason(
3798 pending_id: i64,
3799 triggered_by_policy_rule: &str,
3800 action_type: &str,
3801 reason_codes_json: serde_json::Value,
3802 takedown_action_id: i64,
3803) -> String {
3804 serde_json::json!({
3805 "pending_id": pending_id,
3806 "triggered_by_policy_rule": triggered_by_policy_rule,
3807 "action_type": action_type,
3808 "reason_codes": reason_codes_json,
3809 "triggered_by": "takedown_terminal",
3810 "takedown_action_id": takedown_action_id,
3811 })
3812 .to_string()
3813}
3814
3815/// Convert epoch-ms (i64) to [`SystemTime`]. The v1.4 calculators
3816/// take SystemTime; the schema stores epoch-ms; this is the
3817/// boundary helper. Negative values clamp to UNIX_EPOCH (defense:
3818/// schema columns are non-negative in practice).
3819fn epoch_ms_to_systemtime(ms: i64) -> SystemTime {
3820 if ms >= 0 {
3821 UNIX_EPOCH + Duration::from_millis(ms as u64)
3822 } else {
3823 UNIX_EPOCH
3824 }
3825}
3826
3827/// Inverse of [`epoch_ms_to_systemtime`]. Used by the emission
3828/// path (#60) to translate a [`LabelDraft`]'s `exp: SystemTime`
3829/// into the i64 ms that [`rfc3339_from_epoch_ms`] formats. Errors
3830/// when the value pre-dates UNIX_EPOCH — the emission core never
3831/// produces such a value (caller-supplied `expires_at` is
3832/// epoch-ms-derived), so this should be unreachable in practice.
3833fn systemtime_to_epoch_ms(st: SystemTime) -> Result<i64> {
3834 st.duration_since(UNIX_EPOCH)
3835 .map(|d| d.as_millis() as i64)
3836 .map_err(|e| Error::Signing(format!("SystemTime before unix epoch: {e}")))
3837}
3838
3839/// Route a subject string into (subject_did, subject_uri) per the
3840/// recorder's contract. DIDs go to subject_did with no URI; AT-URIs
3841/// extract the repo DID as subject_did and keep the full URI as
3842/// subject_uri. Anything else is rejected.
3843fn route_subject(subject: &str) -> Result<(String, Option<String>)> {
3844 if let Some(rest) = subject.strip_prefix("at://") {
3845 let repo = rest.split('/').next().unwrap_or("");
3846 if repo.is_empty() || !repo.starts_with("did:") {
3847 return Err(Error::Signing(format!(
3848 "subject at://-URI must start at://did:.../...; got {subject:?}"
3849 )));
3850 }
3851 Ok((repo.to_string(), Some(subject.to_string())))
3852 } else if subject.starts_with("did:") && subject.len() > "did:".len() {
3853 Ok((subject.to_string(), None))
3854 } else {
3855 Err(Error::Signing(format!(
3856 "subject must be a DID (`did:...`) or AT-URI (`at://did:...`); got {subject:?}"
3857 )))
3858 }
3859}
3860
3861/// Parse a narrow ISO-8601 duration subset into a [`Duration`].
3862/// Supports `P{n}D` (days), `P{n}W` (weeks), `PT{n}H` (hours),
3863/// `PT{n}M` (minutes), `PT{n}S` (seconds), and combinations like
3864/// `P1DT12H`. Years and months are NOT supported — moderation
3865/// suspensions are bounded enough that day/hour granularity covers
3866/// real cases without the calendar-arithmetic complexity of Y/M.
3867///
3868/// Returns the duration in seconds.
3869///
3870/// `pub(crate)` so [`crate::policy::automation`] (#71) can validate
3871/// `[policy_automation.rules.<name>].duration` at config-load time
3872/// against the same parser the recorder uses on its input
3873/// boundary — single source of truth for the supported duration
3874/// shape across the moderator-recorded path and the policy-engine
3875/// path.
3876pub(crate) fn parse_iso8601_duration(s: &str) -> Result<u64> {
3877 let body = s.strip_prefix('P').ok_or_else(|| {
3878 Error::Signing(format!(
3879 "duration {s:?} must start with 'P' (ISO-8601 duration form, e.g. P7D)"
3880 ))
3881 })?;
3882 if body.is_empty() {
3883 return Err(Error::Signing(format!(
3884 "duration {s:?} has no components after 'P'"
3885 )));
3886 }
3887
3888 // Split into date part (before T, if any) and time part.
3889 let (date_part, time_part) = match body.split_once('T') {
3890 Some((d, t)) => (d, Some(t)),
3891 None => (body, None),
3892 };
3893
3894 let mut total_secs: u64 = 0;
3895 let mut buf = String::new();
3896 for c in date_part.chars() {
3897 if c.is_ascii_digit() {
3898 buf.push(c);
3899 continue;
3900 }
3901 let n: u64 = buf
3902 .parse()
3903 .map_err(|_| Error::Signing(format!("duration {s:?}: malformed numeric run")))?;
3904 buf.clear();
3905 let mult = match c {
3906 'D' => 86_400u64,
3907 'W' => 7 * 86_400u64,
3908 'Y' | 'M' => {
3909 return Err(Error::Signing(format!(
3910 "duration {s:?}: years (Y) and months (M, without T prefix) not supported in v1.4"
3911 )));
3912 }
3913 other => {
3914 return Err(Error::Signing(format!(
3915 "duration {s:?}: unknown date-part unit {other:?}"
3916 )));
3917 }
3918 };
3919 total_secs = total_secs.saturating_add(n.saturating_mul(mult));
3920 }
3921 if !buf.is_empty() {
3922 return Err(Error::Signing(format!(
3923 "duration {s:?}: trailing digits without unit in date part"
3924 )));
3925 }
3926
3927 if let Some(time_body) = time_part {
3928 if time_body.is_empty() {
3929 return Err(Error::Signing(format!(
3930 "duration {s:?}: 'T' separator with no time components"
3931 )));
3932 }
3933 for c in time_body.chars() {
3934 if c.is_ascii_digit() {
3935 buf.push(c);
3936 continue;
3937 }
3938 let n: u64 = buf.parse().map_err(|_| {
3939 Error::Signing(format!(
3940 "duration {s:?}: malformed numeric run in time part"
3941 ))
3942 })?;
3943 buf.clear();
3944 let mult = match c {
3945 'H' => 3600u64,
3946 'M' => 60u64,
3947 'S' => 1u64,
3948 other => {
3949 return Err(Error::Signing(format!(
3950 "duration {s:?}: unknown time-part unit {other:?}"
3951 )));
3952 }
3953 };
3954 total_secs = total_secs.saturating_add(n.saturating_mul(mult));
3955 }
3956 if !buf.is_empty() {
3957 return Err(Error::Signing(format!(
3958 "duration {s:?}: trailing digits without unit in time part"
3959 )));
3960 }
3961 }
3962
3963 if total_secs == 0 {
3964 return Err(Error::Signing(format!(
3965 "duration {s:?}: parsed to zero — must be at least one second"
3966 )));
3967 }
3968 Ok(total_secs)
3969}
3970
3971/// Predict the next AUTOINCREMENT id for `subject_actions`. SQLite's
3972/// AUTOINCREMENT columns are backed by `sqlite_sequence`; the next
3973/// value is `max(seq, max(rowid)) + 1`. With BEGIN DEFERRED + the
3974/// writer task being the only in-process appender, the predict
3975/// step is racy only against external processes — and the ROLLBACK
3976/// path catches mismatch via the inserted_id check in
3977/// handle_record_action.
3978async fn predict_next_subject_action_id(tx: &mut sqlx::Transaction<'_, Sqlite>) -> Result<i64> {
3979 // sqlite_sequence.seq has no NOT NULL constraint; the type
3980 // override forces a non-nullable i64 because the column is
3981 // populated for every AUTOINCREMENT table that has had at
3982 // least one row inserted (NULL only happens if the row
3983 // doesn't exist, which fetch_optional handles).
3984 let row = sqlx::query!(
3985 r#"SELECT seq AS "seq!: i64" FROM sqlite_sequence WHERE name = 'subject_actions'"#
3986 )
3987 .fetch_optional(&mut **tx)
3988 .await?;
3989 Ok(row.map(|r| r.seq + 1).unwrap_or(1))
3990}
3991
3992/// Predict the next AUTOINCREMENT id for `pending_policy_actions`
3993/// (#73, v1.6). Same posture as
3994/// [`predict_next_subject_action_id`]: writer task is sole
3995/// appender; the predict-then-verify dance catches any race.
3996async fn predict_next_pending_action_id(tx: &mut sqlx::Transaction<'_, Sqlite>) -> Result<i64> {
3997 let row = sqlx::query!(
3998 r#"SELECT seq AS "seq!: i64" FROM sqlite_sequence WHERE name = 'pending_policy_actions'"#
3999 )
4000 .fetch_optional(&mut **tx)
4001 .await?;
4002 Ok(row.map(|r| r.seq + 1).unwrap_or(1))
4003}
4004
4005/// Load subject_actions rows for a subject, oldest-first, projecting
4006/// to [`ActionRecord`] for the v1.4 calculators (decay #50,
4007/// window #51). Filters to the columns the calculators consume —
4008/// the rest stay in the DB.
4009async fn load_subject_actions_for_calc(
4010 tx: &mut sqlx::Transaction<'_, Sqlite>,
4011 subject_did: &str,
4012) -> Result<Vec<ActionRecord>> {
4013 let rows = sqlx::query!(
4014 "SELECT action_type, strike_value_applied, was_dampened,
4015 effective_at, expires_at, revoked_at
4016 FROM subject_actions
4017 WHERE subject_did = ?1
4018 ORDER BY id ASC",
4019 subject_did,
4020 )
4021 .fetch_all(&mut **tx)
4022 .await?;
4023
4024 let mut out = Vec::with_capacity(rows.len());
4025 for r in rows {
4026 let action_type = ActionType::from_db_str(&r.action_type).ok_or_else(|| {
4027 Error::Signing(format!(
4028 "subject_actions row has invalid action_type {:?}",
4029 r.action_type
4030 ))
4031 })?;
4032 let strike_value_applied = u32::try_from(r.strike_value_applied).map_err(|_| {
4033 Error::Signing(format!(
4034 "subject_actions row strike_value_applied {} out of u32 range",
4035 r.strike_value_applied
4036 ))
4037 })?;
4038 out.push(ActionRecord {
4039 strike_value_applied,
4040 effective_at: epoch_ms_to_systemtime(r.effective_at),
4041 revoked_at: r.revoked_at.map(epoch_ms_to_systemtime),
4042 action_type,
4043 expires_at: r.expires_at.map(epoch_ms_to_systemtime),
4044 was_dampened: r.was_dampened != 0,
4045 });
4046 }
4047 Ok(out)
4048}
4049
4050/// Load subject_actions rows for a subject, oldest-first,
4051/// projecting to [`crate::policy::evaluator::ActionForPolicyEval`]
4052/// — the policy evaluator (#72) needs `triggered_by_policy_rule`
4053/// for idempotency detection but the v1.4 calculators don't.
4054/// Separate loader rather than extending [`ActionRecord`] to keep
4055/// the calculator-input projection narrow per the §F20 contract.
4056async fn load_subject_actions_for_policy_eval(
4057 tx: &mut sqlx::Transaction<'_, Sqlite>,
4058 subject_did: &str,
4059) -> Result<Vec<crate::policy::evaluator::ActionForPolicyEval>> {
4060 let rows = sqlx::query!(
4061 r#"SELECT
4062 action_type AS "action_type!: String",
4063 effective_at AS "effective_at!: i64",
4064 revoked_at,
4065 triggered_by_policy_rule
4066 FROM subject_actions
4067 WHERE subject_did = ?1
4068 ORDER BY id ASC"#,
4069 subject_did,
4070 )
4071 .fetch_all(&mut **tx)
4072 .await?;
4073
4074 let mut out = Vec::with_capacity(rows.len());
4075 for r in rows {
4076 let action_type = ActionType::from_db_str(&r.action_type).ok_or_else(|| {
4077 Error::Signing(format!(
4078 "subject_actions row has invalid action_type {:?}",
4079 r.action_type
4080 ))
4081 })?;
4082 out.push(crate::policy::evaluator::ActionForPolicyEval {
4083 effective_at: epoch_ms_to_systemtime(r.effective_at),
4084 action_type,
4085 revoked_at: r.revoked_at.map(epoch_ms_to_systemtime),
4086 triggered_by_policy_rule: r.triggered_by_policy_rule,
4087 });
4088 }
4089 Ok(out)
4090}
4091
4092/// Load pending_policy_actions rows for a subject (oldest-first
4093/// by triggered_at), projecting to
4094/// [`crate::policy::evaluator::PendingActionForPolicyEval`].
4095/// Includes resolved rows so the evaluator's idempotency check
4096/// can distinguish unresolved-vs-confirmed-vs-dismissed per
4097/// #72's contract.
4098async fn load_pending_for_policy_eval(
4099 tx: &mut sqlx::Transaction<'_, Sqlite>,
4100 subject_did: &str,
4101) -> Result<Vec<crate::policy::evaluator::PendingActionForPolicyEval>> {
4102 let rows = sqlx::query!(
4103 r#"SELECT
4104 triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
4105 resolution
4106 FROM pending_policy_actions
4107 WHERE subject_did = ?1
4108 ORDER BY triggered_at ASC"#,
4109 subject_did,
4110 )
4111 .fetch_all(&mut **tx)
4112 .await?;
4113
4114 let mut out = Vec::with_capacity(rows.len());
4115 for r in rows {
4116 let resolution = match r.resolution.as_deref() {
4117 None => None,
4118 Some("confirmed") => Some(crate::policy::evaluator::PendingResolution::Confirmed),
4119 Some("dismissed") => Some(crate::policy::evaluator::PendingResolution::Dismissed),
4120 Some(other) => {
4121 return Err(Error::Signing(format!(
4122 "pending_policy_actions row has invalid resolution {:?}",
4123 other
4124 )));
4125 }
4126 };
4127 out.push(crate::policy::evaluator::PendingActionForPolicyEval {
4128 triggered_by_policy_rule: r.triggered_by_policy_rule,
4129 resolution,
4130 });
4131 }
4132 Ok(out)
4133}
4134
4135/// Auto-dismiss every unresolved `pending_policy_actions` row for
4136/// a subject when a takedown lands against them (§F22 / #76). The
4137/// caller is the takedown's INSERT site (one of three: moderator-
4138/// recorded via [`Writer::handle_record_action`], policy-auto-
4139/// recorded via [`Writer::insert_policy_auto_action`], or pending-
4140/// confirmed via [`Writer::handle_confirm_pending_action`]) — each
4141/// invokes this helper inside its open transaction so the cascade
4142/// commits atomically with the takedown.
4143///
4144/// For each unresolved pending: append a hash-chained
4145/// `pending_policy_action_dismissed` audit row (with
4146/// `triggered_by: "takedown_terminal"` and a `takedown_action_id`
4147/// back-pointer to the triggering subject_actions row), then
4148/// UPDATE the pending's resolution columns to 'dismissed'. The
4149/// audit row is appended BEFORE the UPDATE so a forensic reader
4150/// cannot observe a dismissed pending without an audit row
4151/// explaining why.
4152///
4153/// Per the chainlink scope (#76):
4154/// - The pending row's `confirmed_action_id` stays NULL —
4155/// cascaded dismissals are not "confirmations under another
4156/// name."
4157/// - The pending row's `resolved_by_did` is the takedown's
4158/// actor_did — moderator DID for moderator-recorded /
4159/// confirm-flow takedowns; the synthetic policy DID for
4160/// policy-auto-recorded takedowns.
4161/// - Revoking the takedown later does NOT un-dismiss these
4162/// cascaded pendings; they stay dismissed as forensic record.
4163/// If the subject decay-and-recrosses post-revocation, the rule
4164/// re-fires through normal threshold-crossing logic per #72.
4165///
4166/// Returns the ids of dismissed pendings (possibly empty when the
4167/// subject had no unresolved pendings) so the caller can surface
4168/// the count in tracing or its own audit context.
4169async fn auto_dismiss_pendings_on_takedown(
4170 tx: &mut sqlx::Transaction<'_, Sqlite>,
4171 subject_did: &str,
4172 subject_uri: Option<&str>,
4173 triggering_takedown_id: i64,
4174 actor_did: &str,
4175 now_ms: i64,
4176) -> Result<Vec<i64>> {
4177 let rows = sqlx::query!(
4178 r#"SELECT
4179 id AS "id!: i64",
4180 triggered_by_policy_rule AS "triggered_by_policy_rule!: String",
4181 action_type AS "action_type!: String",
4182 reason_codes AS "reason_codes!: String"
4183 FROM pending_policy_actions
4184 WHERE subject_did = ?1 AND resolution IS NULL
4185 ORDER BY id ASC"#,
4186 subject_did,
4187 )
4188 .fetch_all(&mut **tx)
4189 .await?;
4190
4191 if rows.is_empty() {
4192 return Ok(Vec::new());
4193 }
4194
4195 let audit_target = subject_uri
4196 .map(str::to_string)
4197 .unwrap_or_else(|| subject_did.to_string());
4198 let dismissed_resolution = "dismissed";
4199 let mut dismissed_ids = Vec::with_capacity(rows.len());
4200
4201 for r in rows {
4202 let reason_codes_json: serde_json::Value =
4203 serde_json::from_str(&r.reason_codes).unwrap_or(serde_json::Value::Null);
4204 let audit_reason = build_pending_dismissed_on_takedown_audit_reason(
4205 r.id,
4206 &r.triggered_by_policy_rule,
4207 &r.action_type,
4208 reason_codes_json,
4209 triggering_takedown_id,
4210 );
4211 crate::audit::append::append_in_tx(
4212 tx,
4213 &crate::audit::append::AuditRowForAppend {
4214 created_at: now_ms,
4215 action: "pending_policy_action_dismissed".into(),
4216 actor_did: actor_did.to_string(),
4217 target: Some(audit_target.clone()),
4218 target_cid: None,
4219 outcome: "success".into(),
4220 reason: Some(audit_reason),
4221 },
4222 )
4223 .await?;
4224
4225 sqlx::query!(
4226 "UPDATE pending_policy_actions
4227 SET resolution = ?1,
4228 resolved_at = ?2,
4229 resolved_by_did = ?3
4230 WHERE id = ?4",
4231 dismissed_resolution,
4232 now_ms,
4233 actor_did,
4234 r.id,
4235 )
4236 .execute(&mut **tx)
4237 .await?;
4238
4239 dismissed_ids.push(r.id);
4240 }
4241
4242 Ok(dismissed_ids)
4243}
4244
4245#[cfg(test)]
4246mod tests {
4247 use super::*;
4248
4249 // 1_715_000_000_000 ms since epoch = 2024-05-06T12:53:20.000Z UTC.
4250 // The four tests below pin every branch of the clamp: no prior,
4251 // wall-ahead, wall-behind (use prev+1ms), wall-equal (advance anyway).
4252
4253 #[test]
4254 fn clamp_cts_uses_wall_clock_when_no_prior() {
4255 let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
4256 assert_eq!(out, "2024-05-06T12:53:20.000Z");
4257 }
4258
4259 #[test]
4260 fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
4261 // wall_clock one second behind prev: result is prev + 1ms, not wall.
4262 let prev = "2024-05-06T12:53:20.500Z";
4263 let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
4264 assert_eq!(out, "2024-05-06T12:53:20.501Z");
4265 }
4266
4267 #[test]
4268 fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
4269 let prev = "2024-05-06T12:53:20.500Z";
4270 let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
4271 assert_eq!(out, "2024-05-06T12:53:22.000Z");
4272 }
4273
4274 #[test]
4275 fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
4276 // wall_clock == prev exactly: must advance by 1ms (strictly greater).
4277 let prev = "2024-05-06T12:53:20.500Z";
4278 let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
4279 assert_eq!(out, "2024-05-06T12:53:20.501Z");
4280 }
4281
4282 #[test]
4283 fn build_audit_reason_shape_matches_documented_schema() {
4284 let json = build_audit_reason("spam", false, Some("user reported"));
4285 let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
4286 assert_eq!(v["val"], "spam");
4287 assert_eq!(v["neg"], false);
4288 assert_eq!(v["moderator_reason"], "user reported");
4289 }
4290
4291 #[test]
4292 fn build_audit_reason_null_when_moderator_reason_absent() {
4293 let json = build_audit_reason("spam", true, None);
4294 let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
4295 assert!(v["moderator_reason"].is_null());
4296 }
4297
4298 #[test]
4299 fn rfc3339_roundtrip() {
4300 let s = "2026-04-22T12:00:00.938Z";
4301 let ms = parse_rfc3339_ms(s).expect("parse");
4302 let back = rfc3339_from_epoch_ms(ms).expect("format");
4303 assert_eq!(back, s);
4304 }
4305
4306 // ---------- route_subject (#51) ----------
4307
4308 #[test]
4309 fn route_subject_did_only_returns_did_no_uri() {
4310 let (did, uri) = route_subject("did:plc:abc123").unwrap();
4311 assert_eq!(did, "did:plc:abc123");
4312 assert!(uri.is_none());
4313 }
4314
4315 #[test]
4316 fn route_subject_at_uri_extracts_repo_did_and_keeps_full_uri() {
4317 let (did, uri) = route_subject("at://did:plc:abc/app.bsky.feed.post/3xx").unwrap();
4318 assert_eq!(did, "did:plc:abc");
4319 assert_eq!(
4320 uri.as_deref(),
4321 Some("at://did:plc:abc/app.bsky.feed.post/3xx")
4322 );
4323 }
4324
4325 #[test]
4326 fn route_subject_rejects_at_uri_with_non_did_repo() {
4327 let err = route_subject("at://handle.example/app.bsky.feed.post/3xx").unwrap_err();
4328 assert!(matches!(err, Error::Signing(_)));
4329 }
4330
4331 #[test]
4332 fn route_subject_rejects_arbitrary_string() {
4333 assert!(route_subject("not a subject").is_err());
4334 assert!(route_subject("did:").is_err());
4335 assert!(route_subject("").is_err());
4336 }
4337
4338 // ---------- parse_iso8601_duration (#51) ----------
4339
4340 #[test]
4341 fn duration_p7d_is_seven_days() {
4342 assert_eq!(parse_iso8601_duration("P7D").unwrap(), 7 * 86_400);
4343 }
4344
4345 #[test]
4346 fn duration_p1w_is_seven_days() {
4347 assert_eq!(parse_iso8601_duration("P1W").unwrap(), 7 * 86_400);
4348 }
4349
4350 #[test]
4351 fn duration_pt12h_is_twelve_hours() {
4352 assert_eq!(parse_iso8601_duration("PT12H").unwrap(), 12 * 3600);
4353 }
4354
4355 #[test]
4356 fn duration_pt30m_is_thirty_minutes() {
4357 assert_eq!(parse_iso8601_duration("PT30M").unwrap(), 30 * 60);
4358 }
4359
4360 #[test]
4361 fn duration_pt45s_is_forty_five_seconds() {
4362 assert_eq!(parse_iso8601_duration("PT45S").unwrap(), 45);
4363 }
4364
4365 #[test]
4366 fn duration_p1d_t12h_combines() {
4367 assert_eq!(
4368 parse_iso8601_duration("P1DT12H").unwrap(),
4369 86_400 + 12 * 3600
4370 );
4371 }
4372
4373 #[test]
4374 fn duration_no_p_prefix_rejected() {
4375 assert!(parse_iso8601_duration("7D").is_err());
4376 assert!(parse_iso8601_duration("").is_err());
4377 }
4378
4379 #[test]
4380 fn duration_p_alone_rejected() {
4381 assert!(parse_iso8601_duration("P").is_err());
4382 }
4383
4384 #[test]
4385 fn duration_year_rejected_in_v14() {
4386 let err = parse_iso8601_duration("P1Y").unwrap_err();
4387 let msg = format!("{err}");
4388 assert!(msg.contains("not supported"));
4389 }
4390
4391 #[test]
4392 fn duration_unknown_unit_rejected() {
4393 assert!(parse_iso8601_duration("P5X").is_err());
4394 assert!(parse_iso8601_duration("PT5X").is_err());
4395 }
4396
4397 #[test]
4398 fn duration_zero_rejected() {
4399 assert!(parse_iso8601_duration("P0D").is_err());
4400 }
4401
4402 #[test]
4403 fn duration_trailing_digits_without_unit_rejected() {
4404 assert!(parse_iso8601_duration("P5").is_err());
4405 assert!(parse_iso8601_duration("PT12").is_err());
4406 }
4407
4408 #[test]
4409 fn duration_t_separator_with_no_time_rejected() {
4410 assert!(parse_iso8601_duration("P1DT").is_err());
4411 }
4412
4413 // ---------- audit reason builders (#51) ----------
4414
4415 #[test]
4416 fn record_action_audit_reason_shape() {
4417 let emitted = vec![
4418 serde_json::json!({"val": "!hide", "uri": "did:plc:subject0000000000000000"}),
4419 serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
4420 ];
4421 let json = build_record_action_audit_reason(
4422 42,
4423 ActionType::TempSuspension,
4424 "spam",
4425 &["spam".to_string(), "harassment".to_string()],
4426 4,
4427 2,
4428 true,
4429 &emitted,
4430 "moderator",
4431 None,
4432 None,
4433 );
4434 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4435 assert_eq!(v["action_id"], 42);
4436 assert_eq!(v["action_type"], "temp_suspension");
4437 assert_eq!(v["primary_reason"], "spam");
4438 assert_eq!(v["reason_codes"], serde_json::json!(["spam", "harassment"]));
4439 assert_eq!(v["strike_value_base"], 4);
4440 assert_eq!(v["strike_value_applied"], 2);
4441 assert_eq!(v["was_dampened"], true);
4442 assert_eq!(v["emitted_labels"][0]["val"], "!hide");
4443 assert_eq!(v["emitted_labels"][1]["val"], "reason-spam");
4444 assert_eq!(v["actor_kind"], "moderator");
4445 assert!(v.get("triggered_by_policy_rule").is_none());
4446 assert!(v.get("policy_consequence").is_none());
4447 }
4448
4449 #[test]
4450 fn record_action_audit_reason_empty_emitted_labels() {
4451 let json = build_record_action_audit_reason(
4452 7,
4453 ActionType::Note,
4454 "spam",
4455 &["spam".to_string()],
4456 0,
4457 0,
4458 false,
4459 &[],
4460 "moderator",
4461 None,
4462 None,
4463 );
4464 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4465 assert_eq!(v["emitted_labels"], serde_json::json!([]));
4466 }
4467
4468 #[test]
4469 fn record_action_audit_reason_with_policy_consequence_auto() {
4470 let pc = serde_json::json!({
4471 "rule_fired": "warn_at_5",
4472 "mode": "auto",
4473 "auto_action_id": 43,
4474 });
4475 let json = build_record_action_audit_reason(
4476 42,
4477 ActionType::Warning,
4478 "spam",
4479 &["spam".to_string()],
4480 0,
4481 0,
4482 false,
4483 &[],
4484 "moderator",
4485 None,
4486 Some(pc),
4487 );
4488 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4489 assert_eq!(v["policy_consequence"]["rule_fired"], "warn_at_5");
4490 assert_eq!(v["policy_consequence"]["mode"], "auto");
4491 assert_eq!(v["policy_consequence"]["auto_action_id"], 43);
4492 }
4493
4494 #[test]
4495 fn record_action_audit_reason_for_policy_recorded_action() {
4496 let json = build_record_action_audit_reason(
4497 43,
4498 ActionType::Warning,
4499 "policy-threshold",
4500 &["policy-threshold".to_string()],
4501 0,
4502 0,
4503 false,
4504 &[],
4505 "policy",
4506 Some("warn_at_5"),
4507 None,
4508 );
4509 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4510 assert_eq!(v["actor_kind"], "policy");
4511 assert_eq!(v["triggered_by_policy_rule"], "warn_at_5");
4512 }
4513
4514 #[test]
4515 fn revoke_action_audit_reason_with_reason() {
4516 let json = build_revoke_action_audit_reason(7, Some("appeal granted"), &[]);
4517 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4518 assert_eq!(v["action_id"], 7);
4519 assert_eq!(v["revoked_reason"], "appeal granted");
4520 assert_eq!(v["negated_labels"], serde_json::json!([]));
4521 }
4522
4523 #[test]
4524 fn revoke_action_audit_reason_without_reason_is_null() {
4525 let json = build_revoke_action_audit_reason(7, None, &[]);
4526 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4527 assert!(v["revoked_reason"].is_null());
4528 }
4529
4530 #[test]
4531 fn revoke_action_audit_reason_with_negated_labels() {
4532 let labels = vec![
4533 serde_json::json!({"val": "!takedown", "uri": "did:plc:subject0000000000000000"}),
4534 serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
4535 ];
4536 let json = build_revoke_action_audit_reason(11, None, &labels);
4537 let v: serde_json::Value = serde_json::from_str(&json).unwrap();
4538 assert_eq!(v["negated_labels"][0]["val"], "!takedown");
4539 assert_eq!(v["negated_labels"][1]["val"], "reason-spam");
4540 }
4541
4542 // ---------- emission idempotency guards (#64) ----------
4543
4544 #[test]
4545 fn skip_action_label_when_row_already_carries_an_emitted_val() {
4546 // The defense-in-depth case the guard exists for: an
4547 // existing emitted_label_uri value means the action label
4548 // was already produced on the wire; the emission loop
4549 // must skip to avoid a duplicate (src, uri, val) record.
4550 assert!(should_skip_action_label_emission(Some("!takedown")));
4551 }
4552
4553 #[test]
4554 fn emit_action_label_when_row_has_no_emitted_val() {
4555 // v1.5's recordAction always reaches this branch in
4556 // production (fresh INSERT, NULL column). Pin the
4557 // happy-path so a future refactor that flips the
4558 // sense of the check breaks here.
4559 assert!(!should_skip_action_label_emission(None));
4560 }
4561
4562 #[test]
4563 fn skip_reason_emission_when_already_linked() {
4564 let existing: HashSet<&str> = ["spam", "harassment"].into_iter().collect();
4565 assert!(should_skip_reason_emission("spam", &existing));
4566 assert!(should_skip_reason_emission("harassment", &existing));
4567 }
4568
4569 #[test]
4570 fn emit_reason_when_not_already_linked() {
4571 let existing: HashSet<&str> = ["spam"].into_iter().collect();
4572 assert!(!should_skip_reason_emission("hate-speech", &existing));
4573 // Empty existing set (the v1.5 production case): never skip.
4574 let empty: HashSet<&str> = HashSet::new();
4575 assert!(!should_skip_reason_emission("anything", &empty));
4576 }
4577
4578 #[test]
4579 fn reason_emission_skip_check_is_case_sensitive() {
4580 // Reason codes are operator-vocabulary identifiers from
4581 // [moderation_reasons] (#47). Case sensitivity matches
4582 // SQLite's default text comparison and the recorder's
4583 // primary-reason resolution; pinning here defends
4584 // against an accidental case-fold in the guard during
4585 // refactoring.
4586 let existing: HashSet<&str> = ["SPAM"].into_iter().collect();
4587 assert!(!should_skip_reason_emission("spam", &existing));
4588 }
4589}