Skip to main content

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}