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    "report_resolved",
171    "reporter_flagged",
172    "reporter_unflagged",
173    "retention_sweep",
174    "subject_action_recorded",
175    "subject_action_revoked",
176];
177
178/// Closed set of `audit_log.outcome` values matching the SQL `CHECK`
179/// constraint in `migrations/0001_init.sql`. listAuditLog validates the
180/// `outcome` query param against this set; the SQL constraint itself is
181/// the durable source of truth — this slice mirrors it so the handler
182/// can reject invalid values pre-query without round-tripping to SQLite.
183pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
184
185/// RFC-3339 with millisecond precision. `Z` is appended by
186/// [`rfc3339_from_epoch_ms`] and stripped by [`parse_rfc3339_ms`] — kept
187/// out of the format description because `time::OffsetDateTime::parse`
188/// can't infer a UTC offset from the literal character `Z`, and
189/// symmetric Z-handling at the string boundary is simpler than switching
190/// parsing-only to `well_known::Rfc3339`.
191const CTS_FORMAT: &[FormatItem<'_>] =
192    format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
193
194/// Moderator request to apply a label (positive event).
195#[derive(Debug, Clone)]
196pub struct ApplyLabelRequest {
197    /// The DID of the moderator issuing the request. Becomes
198    /// `audit_log.actor_did`. `src` on the label itself is the writer's
199    /// service DID, not this.
200    pub actor_did: String,
201    /// AT-URI or DID of the subject being labeled.
202    pub uri: String,
203    /// Optional record-version pin. Present means "this specific version";
204    /// absent means "all versions / the account" per §6.1.
205    pub cid: Option<String>,
206    /// Label value, ≤128 bytes (validated on insert by schema CHECK and
207    /// by caller-side input validation on the XRPC boundary).
208    pub val: String,
209    /// Optional expiration timestamp (RFC-3339 Z). Stored only; expiry
210    /// enforcement is v1.1 (§F7).
211    pub exp: Option<String>,
212    /// Free-text reason recorded to `audit_log` only. Not signed, not
213    /// included in the label record.
214    pub moderator_reason: Option<String>,
215}
216
217/// Moderator request to negate (withdraw) a previously-applied label.
218///
219/// Uniqueness is at the `(src, uri, val)` tuple. The negation copies the
220/// most-recent applied event's `cid` so the negation pins the same record
221/// version — callers don't supply it.
222#[derive(Debug, Clone)]
223pub struct NegateLabelRequest {
224    /// Moderator DID issuing the negation. Becomes
225    /// `audit_log.actor_did`.
226    pub actor_did: String,
227    /// AT-URI or DID of the subject whose label is being withdrawn.
228    /// The `cid` pinning (if any) is copied from the prior apply.
229    pub uri: String,
230    /// Label value being negated. Must match an existing applied
231    /// label on `(src, uri, val)` or the call errors with
232    /// `LabelNotFound`.
233    pub val: String,
234    /// Free-text reason recorded to `audit_log` only.
235    pub moderator_reason: Option<String>,
236}
237
238/// A committed label event. Returned by `apply_label` / `negate_label` and
239/// broadcast to subscribers (subscribeLabels consumers in #7). Wire
240/// serialization is that consumer's concern — the LabelEvent itself
241/// carries no serde derive because the canonical encoding for the wire is
242/// the DAG-CBOR path in `crate::signing`, not serde JSON.
243#[derive(Debug, Clone)]
244pub struct LabelEvent {
245    /// Frame sequence number from `label_sequence`. Strictly monotonic.
246    pub seq: i64,
247    /// The full signed label (`sig` populated).
248    pub label: Label,
249}
250
251/// Internal write command. One variant per public method.
252enum WriteCommand {
253    Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
254    Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
255    ResolveReport(
256        ResolveReportRequest,
257        oneshot::Sender<Result<ResolvedReport>>,
258    ),
259    Sweep(SweepRequest, oneshot::Sender<Result<SweepBatchResult>>),
260    /// Append a hash-chained audit row (#39). Used by in-process
261    /// callers that don't have an existing transaction (e.g.,
262    /// `retentionSweep` after the sweep itself completes). Callers
263    /// that already hold a transaction (writer-internal handlers,
264    /// `flag_reporter`) call [`crate::audit::append::append_in_tx`]
265    /// directly within that transaction instead.
266    AppendAudit(
267        crate::audit::append::AuditRowForAppend,
268        oneshot::Sender<Result<i64>>,
269    ),
270    /// Record a graduated-action subject_actions row (§F20 / #51).
271    /// Single-transaction: input validation, strike calculation,
272    /// subject_actions INSERT, strike_state UPSERT, audit_log
273    /// hash-chain append.
274    RecordAction(RecordActionRequest, oneshot::Sender<Result<RecordedAction>>),
275    /// Revoke a previously-recorded subject_actions row (§F20 /
276    /// #51). Single-transaction: row lookup + state-check,
277    /// revoked_at UPDATE, strike_state recompute, audit_log
278    /// hash-chain append.
279    RevokeAction(RevokeActionRequest, oneshot::Sender<Result<RevokedAction>>),
280    Shutdown(oneshot::Sender<Result<()>>),
281}
282
283/// Inline label-application sub-object for [`ResolveReportRequest`].
284/// Structurally aligned with [`ApplyLabelRequest`] minus the
285/// moderator-reason (that field on the outer request already captures
286/// the resolution rationale; the label's own audit reason is derived
287/// from the resolution flow).
288#[derive(Debug, Clone)]
289pub struct ApplyLabelInline {
290    /// Subject the label applies to (AT-URI or DID).
291    pub uri: String,
292    /// Optional record-version pin; `None` targets the account / all
293    /// versions of the record.
294    pub cid: Option<String>,
295    /// Label value to apply (≤128 bytes).
296    pub val: String,
297    /// Optional expiration timestamp (RFC-3339 Z). Stored only;
298    /// enforcement is v1.1.
299    pub exp: Option<String>,
300}
301
302/// Resolution outcome (#27): the implicit "with-label vs without-label"
303/// semantic, made explicit. The wire shape is `applyLabel: Option<…>`
304/// per the `resolveReport` lexicon; this enum is internal and the
305/// handler maps `None → Dismiss` / `Some(_) → ApplyLabel(_)` at the
306/// HTTP boundary.
307#[derive(Debug, Clone)]
308pub enum ResolutionAction {
309    /// Resolve without emitting a label (the operator-UX "dismiss"
310    /// flow). `resolution_label` on the row stays NULL.
311    Dismiss,
312    /// Resolve and emit a label in the same transaction. The label's
313    /// `val` is also recorded as `resolution_label` on the report row.
314    ApplyLabel(ApplyLabelInline),
315}
316
317impl ResolutionAction {
318    /// Borrow the inner `ApplyLabelInline` if this is an `ApplyLabel`
319    /// variant. Convenience for the writer's UPDATE / audit code that
320    /// needs both the optional label *and* its `val` projection.
321    pub fn as_apply(&self) -> Option<&ApplyLabelInline> {
322        match self {
323            ResolutionAction::Dismiss => None,
324            ResolutionAction::ApplyLabel(a) => Some(a),
325        }
326    }
327}
328
329/// Request to resolve a report (§F12 `resolveReport`). The optional
330/// label is applied **in the same transaction** as the report status
331/// update and both audit rows — §F5 single-writer invariant plus §F12
332/// atomicity requirement documented in the `resolveReport` lexicon.
333#[derive(Debug, Clone)]
334pub struct ResolveReportRequest {
335    /// Moderator DID issuing the resolution. Becomes
336    /// `audit_log.actor_did` on both the label-applied (if any) and
337    /// report_resolved rows.
338    pub actor_did: String,
339    /// Primary key of the report being resolved.
340    pub report_id: i64,
341    /// Whether the resolution emits a label or just closes the report.
342    /// Replaces the pre-#27 `apply_label: Option<ApplyLabelInline>`
343    /// representation; semantically identical, named for the operator
344    /// UX (dismiss vs apply-label).
345    pub action: ResolutionAction,
346    /// Free-text resolution rationale recorded to audit_log.
347    pub resolution_reason: Option<String>,
348}
349
350/// Result of a successful resolve. `label_event` is `Some(..)` iff
351/// the request carried `apply_label`; the broadcast has already
352/// happened inside the writer task post-commit.
353#[derive(Debug, Clone)]
354pub struct ResolvedReport {
355    /// The updated report row (status now `resolved`).
356    pub report: crate::report::Report,
357    /// The emitted label event, if the resolution included an
358    /// `apply_label`. Signed, broadcast, and committed as part of
359    /// the same transaction as the report UPDATE.
360    pub label_event: Option<LabelEvent>,
361}
362
363/// Trigger for the retention sweep (§F4). Carries no per-call
364/// parameters today — the cutoff comes from `retention_days` baked
365/// into the writer at spawn time, not from the request — but is
366/// kept as a typed unit so a future "sweep with explicit override
367/// for this one run" remains a non-breaking change to the variant.
368#[derive(Debug, Clone, Default)]
369pub struct SweepRequest;
370
371/// Per-batch outcome from one sweep dispatch through the writer's
372/// internal `WriteCommand::Sweep` channel. The writer task processes
373/// ONE batch per command so its main `select!` can interleave
374/// incoming label writes between batches (§F4 + §F5 — single-writer
375/// invariant + bounded latency). Callers who want a full sweep loop
376/// until [`Self::has_more`] is `false`; the [`WriterHandle::sweep`]
377/// convenience wrapper does this internally.
378#[derive(Debug, Clone)]
379pub struct SweepBatchResult {
380    /// Rows deleted in this batch. Zero means "no more old rows
381    /// match the cutoff" — caller stops looping.
382    pub rows_deleted: i64,
383    /// `true` when the batch hit the configured `sweep_batch_size`
384    /// limit and a follow-up batch may find more rows. `false`
385    /// indicates the batch was partial (last batch) and the sweep
386    /// is complete.
387    pub has_more: bool,
388    /// Cutoff days actually applied. `None` when the writer was
389    /// spawned with `retention_days = None` — the sweep is a no-op
390    /// in that configuration and `rows_deleted` is always 0.
391    pub retention_days_applied: Option<u32>,
392}
393
394/// Request to record a new subject_actions row (§F20 / #51 graduated-
395/// action moderation). The handler validates the inputs, computes
396/// the strike value at action time via the v1.4 calculators
397/// (#48/#49/#50/#51), inserts the row, updates the strike-state
398/// cache, and writes a hash-chained audit_log row — all in a
399/// single transaction per §F5 atomicity.
400#[derive(Debug, Clone)]
401pub struct RecordActionRequest {
402    /// Raw subject — DID (`did:plc:...`, `did:web:...`) or AT-URI
403    /// (`at://did:.../col/r`). The handler routes to subject_did vs
404    /// subject_uri based on prefix; for AT-URIs the parent repo DID
405    /// is extracted as subject_did and strike accounting rolls up
406    /// to the account.
407    pub subject: String,
408    /// DID of the moderator/admin recording the action (JWT iss
409    /// from XRPC; CLI caller DID otherwise). Becomes
410    /// subject_actions.actor_did and audit_log.actor_did.
411    pub actor_did: String,
412    /// Graduated-action category. The handler enforces:
413    /// `temp_suspension` requires `duration_iso`; everything else
414    /// rejects it.
415    pub action_type: ActionType,
416    /// Operator-vocabulary identifiers from `[moderation_reasons]`.
417    /// Must be non-empty. Multi-reason resolution: severe wins; else
418    /// highest base_weight wins; ties → first-listed.
419    pub reason_codes: Vec<String>,
420    /// ISO-8601 duration string (e.g. `P7D`). Required for
421    /// `temp_suspension`; rejected for other types. Stored verbatim
422    /// on the row for display; `expires_at` is the canonical
423    /// "when does it end" surface.
424    pub duration_iso: Option<String>,
425    /// Optional moderator-facing rationale.
426    pub notes: Option<String>,
427    /// Optional report row ids that motivated this action.
428    pub report_ids: Vec<i64>,
429}
430
431/// Result of a successful [`WriterHandle::record_action`]. The
432/// admin handler echoes these fields verbatim in the
433/// `tools.cairn.admin.recordAction` response; the CLI surfaces them
434/// in human and JSON output.
435#[derive(Debug, Clone)]
436pub struct RecordedAction {
437    /// Inserted subject_actions row id.
438    pub action_id: i64,
439    /// Reason's `base_weight` before dampening. `0` for warning/note.
440    pub strike_value_base: u32,
441    /// Strike weight actually applied after dampening.
442    pub strike_value_applied: u32,
443    /// `true` iff the dampening curve was consulted (see #49).
444    pub was_dampened: bool,
445    /// Subject's `current_strike_count` BEFORE this action — frozen
446    /// for forensic history.
447    pub strikes_at_time_of_action: u32,
448}
449
450/// Request to revoke a previously-recorded subject_actions row
451/// (§F20 / #51). Sets the row's revoked_at/revoked_by_did/
452/// revoked_reason columns (the schema's no-update-except-revoke
453/// trigger permits exactly this NULL→non-NULL transition);
454/// recomputes the strike-state cache; writes a hash-chained
455/// audit_log row. All in a single transaction.
456#[derive(Debug, Clone)]
457pub struct RevokeActionRequest {
458    /// subject_actions.id to revoke. Errors with
459    /// [`Error::ActionNotFound`] when the row doesn't exist;
460    /// [`Error::ActionAlreadyRevoked`] when revoked_at is already
461    /// non-NULL.
462    pub action_id: i64,
463    /// DID of the moderator/admin performing the revocation.
464    pub revoked_by_did: String,
465    /// Optional rationale stored on the row's revoked_reason column.
466    pub revoked_reason: Option<String>,
467}
468
469/// Result of a successful [`WriterHandle::revoke_action`].
470#[derive(Debug, Clone)]
471pub struct RevokedAction {
472    /// The revoked row's id (echoed for confirmation).
473    pub action_id: i64,
474    /// Wall-clock the revocation took effect, as RFC-3339 Z
475    /// (matches the [`AUDIT_REASON_RECORD_ACTION`] surface for
476    /// consumer convenience).
477    pub revoked_at: String,
478}
479
480/// Aggregate result of a full sweep run (returned by
481/// [`WriterHandle::sweep`] after looping over batches).
482#[derive(Debug, Clone)]
483pub struct SweepResult {
484    /// Total rows deleted across all batches.
485    pub rows_deleted: i64,
486    /// Number of batches issued.
487    pub batches: u64,
488    /// Wall-clock duration of the full sweep, in milliseconds.
489    pub duration_ms: u64,
490    /// Cutoff days actually applied, or `None` when the writer's
491    /// `retention_days` is `None` (sweep is a no-op).
492    pub retention_days_applied: Option<u32>,
493}
494
495/// Cheap handle to the writer task. Clones share the same underlying
496/// mpsc channel; a drop of the last clone is a silent shutdown signal
497/// to the writer task (receiver closes). For a clean shutdown that
498/// releases the lease row, call [`WriterHandle::shutdown`].
499#[derive(Debug, Clone)]
500pub struct WriterHandle {
501    tx: mpsc::Sender<WriteCommand>,
502    broadcast_tx: broadcast::Sender<LabelEvent>,
503    /// Flipped to `true` when the writer task starts its shutdown path.
504    /// Exposed via [`WriterHandle::shutdown_signal`] so peer components
505    /// (the subscribeLabels handler, future maintenance tasks) can close
506    /// their own resources cleanly. Using a watch channel rather than the
507    /// broadcast-channel close signal because clones of `WriterHandle`
508    /// keep `broadcast_tx` alive — receivers would otherwise never see
509    /// `RecvError::Closed`.
510    shutdown_rx: watch::Receiver<bool>,
511}
512
513impl WriterHandle {
514    /// Submit an apply-label request. Resolves when the writer has
515    /// committed the transaction and broadcast the event, or returns an
516    /// `Err` describing why the write was rejected.
517    pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
518        let (reply_tx, reply_rx) = oneshot::channel();
519        self.tx
520            .send(WriteCommand::Apply(req, reply_tx))
521            .await
522            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
523        reply_rx
524            .await
525            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
526    }
527
528    /// Submit a negate-label request. Returns [`Error::LabelNotFound`] if
529    /// no applied label currently exists for `(service_did, uri, val)`.
530    pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
531        let (reply_tx, reply_rx) = oneshot::channel();
532        self.tx
533            .send(WriteCommand::Negate(req, reply_tx))
534            .await
535            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
536        reply_rx
537            .await
538            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
539    }
540
541    /// Resolve a report, optionally applying a label in the same
542    /// atomic transaction. The label INSERT + audit row + report
543    /// UPDATE + resolution audit row all commit together or not at
544    /// all; the label broadcast fires post-commit inside the writer
545    /// task (§F5 + §F12 atomicity contract).
546    ///
547    /// Errors:
548    /// - [`Error::ReportNotFound`] — `report_id` doesn't exist.
549    /// - [`Error::ReportAlreadyResolved`] — report is not in the
550    ///   `pending` state. Handler maps to a generic `InvalidRequest`.
551    pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
552        let (reply_tx, reply_rx) = oneshot::channel();
553        self.tx
554            .send(WriteCommand::ResolveReport(req, reply_tx))
555            .await
556            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
557        reply_rx
558            .await
559            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
560    }
561
562    /// Trigger a §F4 retention sweep through the writer task. Loops
563    /// over single-batch dispatches until a batch returns
564    /// `has_more = false`, aggregating rows + duration.
565    /// Other writer commands interleave between batches (single-writer
566    /// invariant + bounded latency).
567    ///
568    /// Returns `rows_deleted = 0, batches = 0` when the writer was
569    /// spawned with `retention_days = None` (sweep configured off):
570    /// the first batch returns immediately and the loop exits with
571    /// the no-op result.
572    pub async fn sweep(&self, _req: SweepRequest) -> Result<SweepResult> {
573        let start = std::time::Instant::now();
574        let mut total_rows: i64 = 0;
575        let mut batches: u64 = 0;
576
577        loop {
578            let (reply_tx, reply_rx) = oneshot::channel();
579            self.tx
580                .send(WriteCommand::Sweep(SweepRequest, reply_tx))
581                .await
582                .map_err(|_| Error::Signing("writer task is shut down".into()))?;
583            let batch = reply_rx
584                .await
585                .map_err(|_| Error::Signing("writer dropped reply channel".into()))??;
586
587            total_rows += batch.rows_deleted;
588            // Count the no-op early-exit batch too — it represents
589            // the round-trip the caller paid for. Useful for tracing.
590            batches += 1;
591
592            if !batch.has_more {
593                return Ok(SweepResult {
594                    rows_deleted: total_rows,
595                    batches,
596                    duration_ms: start.elapsed().as_millis() as u64,
597                    retention_days_applied: batch.retention_days_applied,
598                });
599            }
600        }
601    }
602
603    /// Record a graduated-action moderation event (§F20 / #51).
604    /// The writer validates inputs, computes the strike value at
605    /// action time via the v1.4 calculators (#48/#49/#50/#51),
606    /// inserts the subject_actions row, updates the strike-state
607    /// cache, and writes a hash-chained audit_log row — all in a
608    /// single transaction. Errors:
609    ///
610    /// - [`Error::ReasonNotFound`] when a reason identifier is not
611    ///   declared in `[moderation_reasons]`.
612    /// - [`Error::DurationRequiredForTempSuspension`] for
613    ///   `temp_suspension` without `duration_iso`.
614    /// - [`Error::DurationOnlyForTempSuspension`] for non-temp
615    ///   types with `duration_iso` set.
616    /// - [`Error::Signing`] for malformed subject / duration / etc.
617    pub async fn record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
618        let (reply_tx, reply_rx) = oneshot::channel();
619        self.tx
620            .send(WriteCommand::RecordAction(req, reply_tx))
621            .await
622            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
623        reply_rx
624            .await
625            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
626    }
627
628    /// Revoke a previously-recorded action (§F20 / #51). Errors:
629    ///
630    /// - [`Error::ActionNotFound`] when no row matches `action_id`.
631    /// - [`Error::ActionAlreadyRevoked`] when the row's revoked_at
632    ///   is already non-NULL — the schema trigger forbids
633    ///   re-revocation; the handler catches it before the UPDATE.
634    pub async fn revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
635        let (reply_tx, reply_rx) = oneshot::channel();
636        self.tx
637            .send(WriteCommand::RevokeAction(req, reply_tx))
638            .await
639            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
640        reply_rx
641            .await
642            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
643    }
644
645    /// Append a hash-chained audit row through the writer task (#39).
646    /// Used by in-process callers that don't already hold a transaction
647    /// — e.g., `retentionSweep`'s post-sweep audit row. Callers that
648    /// have an open transaction (writer-internal handlers,
649    /// `flag_reporter`) use [`crate::audit::append::append_in_tx`]
650    /// directly so the audit row commits atomically with the rest of
651    /// their work. Cross-process CLIs (publish/unpublish-service-
652    /// record) use [`crate::audit::append::append_via_pool`].
653    ///
654    /// Returns the inserted `audit_log.id`.
655    pub async fn append_audit(&self, row: crate::audit::append::AuditRowForAppend) -> Result<i64> {
656        let (reply_tx, reply_rx) = oneshot::channel();
657        self.tx
658            .send(WriteCommand::AppendAudit(row, reply_tx))
659            .await
660            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
661        reply_rx
662            .await
663            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
664    }
665
666    /// Subscribe to committed events. The returned receiver lags past
667    /// the internal broadcast buffer; the consumer (subscribeLabels, #7)
668    /// turns `RecvError::Lagged` into a connection close.
669    pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
670        self.broadcast_tx.subscribe()
671    }
672
673    /// Count of live broadcast receivers. Exposed for test sync points
674    /// (the WS handler subscribes asynchronously after upgrade; tests
675    /// that want to emit an event "after the subscriber is ready" poll
676    /// this until it reaches the expected count). Not a production hot
677    /// path — `broadcast::Sender::receiver_count` walks an atomic chain.
678    pub fn receiver_count(&self) -> usize {
679        self.broadcast_tx.receiver_count()
680    }
681
682    /// Observe writer lifecycle. The returned receiver starts at `false`
683    /// and flips to `true` once exactly, when the writer has accepted a
684    /// shutdown command and is about to release the lease. Holders close
685    /// downstream connections gracefully on the `true` transition.
686    pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
687        self.shutdown_rx.clone()
688    }
689
690    /// Explicit shutdown. Drains in-flight writes, releases the lease
691    /// row, stops the heartbeat, and returns. Idempotent-ish: calling on
692    /// an already-shut-down writer returns a channel-closed error, not a
693    /// panic.
694    pub async fn shutdown(&self) -> Result<()> {
695        let (reply_tx, reply_rx) = oneshot::channel();
696        self.tx
697            .send(WriteCommand::Shutdown(reply_tx))
698            .await
699            .map_err(|_| Error::Signing("writer task is already shut down".into()))?;
700        reply_rx
701            .await
702            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
703    }
704}
705
706/// Start the writer task. Acquires the instance lease, bootstraps the
707/// `signing_keys` row if empty, and spawns the event loop + heartbeat.
708///
709/// `service_did` stamps `labels.src` on every emitted event. `key` must
710/// be the private key whose public form is recorded (or will be recorded)
711/// in `signing_keys`.
712///
713/// `retention_days` is the §F4 retention cutoff in days. `None` disables
714/// the retention sweep entirely (`WriteCommand::Sweep` becomes a no-op
715/// returning `rows_deleted = 0`); `Some(N)` lets the sweep delete labels
716/// whose `created_at` is older than `now - N days`. Source of truth is
717/// [`crate::SubscribeConfig::retention_days`] — pass it through verbatim.
718///
719/// `retention` is the sweep execution policy (schedule + batching);
720/// distinct from `retention_days` per the [§F4 split](crate::RetentionConfig).
721///
722/// `reason_vocabulary` and `strike_policy` are the resolved v1.4
723/// moderation surface (#47 / #48). The recorder (#51 RecordAction
724/// handler) consults them at action time to validate reason codes
725/// and compute strike values. They're held by the writer task so
726/// the recorder doesn't have to thread them through every call.
727/// Operators restart `cairn serve` to change either — same posture
728/// as `[labeler]` config.
729///
730/// `label_emission_policy` is the resolved `[label_emission]`
731/// surface (#58, v1.5). The recorder (#60) consults it post-INSERT
732/// to translate the freshly-recorded action into ATProto labels
733/// emitted in the same transaction. Held by the writer task for
734/// the same reason as the v1.4 surfaces above.
735#[allow(clippy::too_many_arguments)]
736pub async fn spawn(
737    pool: Pool<Sqlite>,
738    key: SigningKey,
739    service_did: String,
740    retention_days: Option<u32>,
741    retention: RetentionConfig,
742    reason_vocabulary: ReasonVocabulary,
743    strike_policy: StrikePolicy,
744    label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
745) -> Result<WriterHandle> {
746    let instance_id = acquire_lease(&pool).await?;
747    let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
748
749    let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
750    let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
751    let (shutdown_tx, shutdown_rx) = watch::channel(false);
752
753    let writer = Writer {
754        pool,
755        key,
756        service_did,
757        signing_key_id,
758        instance_id,
759        rx,
760        broadcast_tx: broadcast_tx.clone(),
761        shutdown_tx,
762        retention_days,
763        retention,
764        reason_vocabulary,
765        strike_policy,
766        label_emission_policy,
767    };
768
769    tokio::spawn(writer.run());
770
771    Ok(WriterHandle {
772        tx,
773        broadcast_tx,
774        shutdown_rx,
775    })
776}
777
778// ---------- Writer task ----------
779
780struct Writer {
781    pool: Pool<Sqlite>,
782    key: SigningKey,
783    service_did: String,
784    signing_key_id: i64,
785    instance_id: String,
786    rx: mpsc::Receiver<WriteCommand>,
787    broadcast_tx: broadcast::Sender<LabelEvent>,
788    shutdown_tx: watch::Sender<bool>,
789    /// Retention cutoff in days. `None` makes the sweep a no-op.
790    /// Source of truth is [`crate::SubscribeConfig::retention_days`].
791    retention_days: Option<u32>,
792    /// Sweep execution policy (schedule + batching).
793    retention: RetentionConfig,
794    /// Resolved [moderation_reasons] vocabulary (#47). Consulted by
795    /// the RecordAction handler to validate reason codes and look
796    /// up base_weight / severe.
797    reason_vocabulary: ReasonVocabulary,
798    /// Resolved [strike_policy] (#48). Consulted by the RecordAction
799    /// handler for the dampening curve, threshold, and decay
800    /// parameters.
801    strike_policy: StrikePolicy,
802    /// Resolved [label_emission] (#58, v1.5). Consulted by the
803    /// RecordAction handler post-INSERT to translate the action into
804    /// ATProto labels, emitted in the same tx.
805    label_emission_policy: crate::labels::policy::LabelEmissionPolicy,
806}
807
808/// Internal accumulator for an in-flight scheduled sweep. Lives only
809/// while [`Writer::run`] is mid-sweep; absent between sweeps.
810struct SweepRunState {
811    started_at: Instant,
812    rows: i64,
813    batches: u64,
814}
815
816impl Writer {
817    /// Compute the next absolute [`Instant`] at which the scheduled
818    /// sweep should fire, or `None` when the sweep is disabled.
819    ///
820    /// "Disabled" = `sweep_enabled = false` OR `retention_days = None`.
821    /// In the second case the sweep would be a no-op anyway; skipping
822    /// the timer entirely keeps the run loop quiet.
823    ///
824    /// Targets the next occurrence of `sweep_run_at_utc_hour:00:00`
825    /// UTC. If we're already past that hour today, target tomorrow at
826    /// the same hour.
827    fn compute_next_sweep_fire(&self) -> Option<Instant> {
828        if !self.retention.sweep_enabled || self.retention_days.is_none() {
829            return None;
830        }
831        let now_utc = OffsetDateTime::now_utc();
832        let target_hour = self.retention.sweep_run_at_utc_hour;
833        let target_today_time = time::Time::from_hms(target_hour, 0, 0)
834            .expect("hour validated < 24 by Config::validate");
835        let target_today = now_utc.replace_time(target_today_time);
836        let target_dt = if target_today <= now_utc {
837            target_today + time::Duration::days(1)
838        } else {
839            target_today
840        };
841        let wait = target_dt - now_utc;
842        let wait_secs = wait.whole_seconds().max(0) as u64;
843        Some(Instant::now() + Duration::from_secs(wait_secs))
844    }
845
846    async fn run(mut self) {
847        let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
848        heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
849        // First tick fires immediately; skip it — we just acquired the lease
850        // with a fresh timestamp, so touching it again is redundant.
851        heartbeat_timer.tick().await;
852
853        // Scheduled-sweep wiring (§F4). The check timer wakes the run
854        // loop once a minute to evaluate "is it time to start a sweep?";
855        // when a sweep starts, `sweep_state` is populated and the
856        // immediate-batch arm runs one batch per loop iteration until
857        // `has_more = false`. Inter-batch yielding to incoming commands
858        // is automatic — the biased select prefers `rx.recv()` and
859        // `heartbeat_timer.tick()` over running another batch.
860        let mut sweep_check_timer = interval(SWEEP_CHECK_INTERVAL);
861        sweep_check_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
862        sweep_check_timer.tick().await;
863
864        let mut next_scheduled_fire: Option<Instant> = self.compute_next_sweep_fire();
865        let mut sweep_state: Option<SweepRunState> = None;
866
867        loop {
868            let sweep_in_progress = sweep_state.is_some();
869            tokio::select! {
870                biased;
871                // Prefer draining write commands over heartbeats so a
872                // steady-state of inbound work does not starve itself
873                // behind a background housekeeping task.
874                cmd = self.rx.recv() => {
875                    match cmd {
876                        Some(WriteCommand::Apply(req, reply)) => {
877                            let res = self.handle_apply(req).await;
878                            // Caller may have cancelled; dropping the reply is fine.
879                            let _ = reply.send(res);
880                        }
881                        Some(WriteCommand::Negate(req, reply)) => {
882                            let res = self.handle_negate(req).await;
883                            let _ = reply.send(res);
884                        }
885                        Some(WriteCommand::ResolveReport(req, reply)) => {
886                            let res = self.handle_resolve_report(req).await;
887                            let _ = reply.send(res);
888                        }
889                        Some(WriteCommand::Sweep(req, reply)) => {
890                            // Single-batch dispatch — the caller loops
891                            // until has_more=false. Per-batch return
892                            // lets the writer's biased select pick up
893                            // pending Apply / Negate / ResolveReport
894                            // commands between batches (§F4 + §F5).
895                            let res = self.handle_sweep(req).await;
896                            let _ = reply.send(res);
897                        }
898                        Some(WriteCommand::AppendAudit(row, reply)) => {
899                            let res = self.handle_append_audit(row).await;
900                            let _ = reply.send(res);
901                        }
902                        Some(WriteCommand::RecordAction(req, reply)) => {
903                            let res = self.handle_record_action(req).await;
904                            let _ = reply.send(res);
905                        }
906                        Some(WriteCommand::RevokeAction(req, reply)) => {
907                            let res = self.handle_revoke_action(req).await;
908                            let _ = reply.send(res);
909                        }
910                        Some(WriteCommand::Shutdown(reply)) => {
911                            // Flip the shutdown watch *before* releasing the
912                            // lease so subscriber tasks see the signal while
913                            // the DB is still accessible for their close-
914                            // frame sends.
915                            let _ = self.shutdown_tx.send(true);
916                            let res = self.release_lease().await;
917                            let _ = reply.send(res);
918                            return;
919                        }
920                        None => {
921                            // All handles dropped without explicit shutdown.
922                            // Best-effort lease release so a same-process
923                            // restart doesn't trip the 60s wait.
924                            let _ = self.shutdown_tx.send(true);
925                            if let Err(e) = self.release_lease().await {
926                                tracing::error!("lease release on handle drop: {e}");
927                            }
928                            return;
929                        }
930                    }
931                }
932                _ = heartbeat_timer.tick() => {
933                    if let Err(e) = self.heartbeat().await {
934                        // Transient DB errors are logged, not fatal. A
935                        // sustained failure is caught at the next-instance
936                        // startup check — not at this writer's expense.
937                        tracing::error!("lease heartbeat failed: {e}");
938                    }
939                }
940                _ = sweep_check_timer.tick() => {
941                    if sweep_state.is_none()
942                        && let Some(fire_at) = next_scheduled_fire
943                        && Instant::now() >= fire_at
944                    {
945                        sweep_state = Some(SweepRunState {
946                            started_at: Instant::now(),
947                            rows: 0,
948                            batches: 0,
949                        });
950                        next_scheduled_fire = Some(fire_at + Duration::from_secs(86_400));
951                        tracing::info!(
952                            retention_days = ?self.retention_days,
953                            sweep_batch_size = self.retention.sweep_batch_size,
954                            "scheduled retention sweep starting"
955                        );
956                    }
957                }
958                // Always-ready arm gated on sweep_in_progress. With biased
959                // order this ranks below rx.recv() + the two timers, so
960                // normal commands and heartbeats interleave naturally
961                // between batches (§F4 inter-batch yield, §F5 single-
962                // writer invariant).
963                _ = std::future::ready(()), if sweep_in_progress => {
964                    match self.handle_sweep(SweepRequest).await {
965                        Ok(batch) => {
966                            let state = sweep_state
967                                .as_mut()
968                                .expect("sweep_in_progress => sweep_state Some");
969                            state.rows += batch.rows_deleted;
970                            state.batches += 1;
971                            if !batch.has_more {
972                                let final_state = sweep_state.take().expect("just set");
973                                tracing::info!(
974                                    rows_deleted = final_state.rows,
975                                    batches = final_state.batches,
976                                    duration_ms = final_state.started_at.elapsed().as_millis() as u64,
977                                    retention_days_applied = ?batch.retention_days_applied,
978                                    "scheduled retention sweep complete"
979                                );
980                            }
981                        }
982                        Err(e) => {
983                            let final_state = sweep_state.take();
984                            tracing::error!(
985                                error = %e,
986                                batches_done = final_state.as_ref().map(|s| s.batches).unwrap_or(0),
987                                rows_so_far = final_state.as_ref().map(|s| s.rows).unwrap_or(0),
988                                "scheduled retention sweep batch failed; aborting run (will retry on next schedule)"
989                            );
990                        }
991                    }
992                }
993            }
994        }
995    }
996
997    /// Process one batch of the §F4 retention sweep. Single-batch
998    /// dispatch — see [`WriteCommand::Sweep`] for the loop ordering.
999    ///
1000    /// Returns immediately with `rows_deleted = 0, has_more = false`
1001    /// when `retention_days` is `None`. Otherwise issues one
1002    /// `DELETE FROM labels WHERE created_at < cutoff LIMIT N` in its
1003    /// own transaction; sets `has_more = true` iff the batch hit the
1004    /// `sweep_batch_size` limit (suggesting more rows may match).
1005    ///
1006    /// Errors are propagated as `Err(_)` rather than swallowed: a
1007    /// transient DB failure should surface to the operator on a
1008    /// manual sweep, and to the schedule-loop logger on a scheduled
1009    /// sweep. Idempotency (Q5) means the next sweep retries cleanly.
1010    async fn handle_sweep(&self, _req: SweepRequest) -> Result<SweepBatchResult> {
1011        let Some(days) = self.retention_days else {
1012            return Ok(SweepBatchResult {
1013                rows_deleted: 0,
1014                has_more: false,
1015                retention_days_applied: None,
1016            });
1017        };
1018
1019        let cutoff_ms = epoch_ms_now() - (days as i64) * 86_400_000;
1020        let limit = self.retention.sweep_batch_size;
1021
1022        let mut tx = self.pool.begin().await?;
1023        // SQLite's DELETE doesn't support LIMIT without the
1024        // SQLITE_ENABLE_UPDATE_DELETE_LIMIT compile flag (off in
1025        // bundled builds). The rowid IN (SELECT ... LIMIT) trick is
1026        // the canonical portable workaround.
1027        let result = sqlx::query!(
1028            "DELETE FROM labels WHERE rowid IN (
1029               SELECT rowid FROM labels WHERE created_at < ?1 LIMIT ?2
1030             )",
1031            cutoff_ms,
1032            limit,
1033        )
1034        .execute(&mut *tx)
1035        .await?;
1036        tx.commit().await?;
1037
1038        let rows = result.rows_affected() as i64;
1039        Ok(SweepBatchResult {
1040            rows_deleted: rows,
1041            has_more: rows >= limit,
1042            retention_days_applied: Some(days),
1043        })
1044    }
1045
1046    /// Handler for [`WriteCommand::AppendAudit`]. Opens its own
1047    /// transaction (BEGIN DEFERRED — fine because the writer task is
1048    /// the only in-process audit appender during a `cairn serve`
1049    /// session, and cross-process appenders use BEGIN IMMEDIATE on
1050    /// their side), inserts the audit row with a freshly-computed
1051    /// hash, commits.
1052    async fn handle_append_audit(
1053        &self,
1054        row: crate::audit::append::AuditRowForAppend,
1055    ) -> Result<i64> {
1056        let mut tx = self.pool.begin().await?;
1057        let id = crate::audit::append::append_in_tx(&mut tx, &row).await?;
1058        tx.commit().await?;
1059        Ok(id)
1060    }
1061
1062    async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
1063        let mut tx = self.pool.begin().await?;
1064        let created_at = epoch_ms_now();
1065        let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
1066        tx.commit().await?;
1067        // No-receivers is not a write failure (§plan point G).
1068        let _ = self.broadcast_tx.send(event.clone());
1069        Ok(event)
1070    }
1071
1072    /// Inner label-application pipeline: reserve seq → clamp cts →
1073    /// sign → INSERT label → INSERT audit. Does NOT commit the
1074    /// transaction and does NOT broadcast — caller handles both.
1075    ///
1076    /// Extracted from `handle_apply` so `handle_resolve_report` can
1077    /// reuse it inside the same transaction as the report update,
1078    /// preserving §F5's single-writer-owns-label-emission invariant
1079    /// (this helper only runs on the writer task; no other code path
1080    /// can access seq allocation or cts clamping).
1081    async fn apply_label_inner(
1082        &self,
1083        tx: &mut sqlx::Transaction<'_, Sqlite>,
1084        req: &ApplyLabelRequest,
1085        created_at_ms: i64,
1086    ) -> Result<LabelEvent> {
1087        let event = self
1088            .sign_and_persist_label(
1089                tx,
1090                &req.val,
1091                &req.uri,
1092                req.cid.as_deref(),
1093                false,
1094                req.exp.as_deref(),
1095                created_at_ms,
1096            )
1097            .await?;
1098
1099        let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
1100        crate::audit::append::append_in_tx(
1101            tx,
1102            &crate::audit::append::AuditRowForAppend {
1103                created_at: created_at_ms,
1104                action: "label_applied".into(),
1105                actor_did: req.actor_did.clone(),
1106                target: Some(req.uri.clone()),
1107                target_cid: req.cid.clone(),
1108                outcome: "success".into(),
1109                reason: Some(audit_reason),
1110            },
1111        )
1112        .await?;
1113
1114        Ok(event)
1115    }
1116
1117    /// Tx-scoped sign-and-persist core for one label row: reserve
1118    /// seq → fetch prev cts → clamp → build wire-level [`Label`] →
1119    /// sign → INSERT into `labels`. Returns the seq + signed label.
1120    /// Does NOT write audit rows and does NOT broadcast — callers
1121    /// own both (the post-tx broadcast and any audit row layered
1122    /// over the bare label INSERT).
1123    ///
1124    /// Reused by [`Self::apply_label_inner`] (which adds a
1125    /// `label_applied` audit row) and by
1126    /// [`Self::handle_record_action`] (which writes one consolidated
1127    /// `subject_action_recorded` audit row capturing the action plus
1128    /// every emitted label, rather than per-label rows). Both reach
1129    /// it only from the single writer task, preserving §F5.
1130    #[allow(clippy::too_many_arguments)]
1131    async fn sign_and_persist_label(
1132        &self,
1133        tx: &mut sqlx::Transaction<'_, Sqlite>,
1134        val: &str,
1135        uri: &str,
1136        cid: Option<&str>,
1137        neg: bool,
1138        exp: Option<&str>,
1139        created_at_ms: i64,
1140    ) -> Result<LabelEvent> {
1141        let seq = reserve_seq(tx).await?;
1142        let prev_cts: Option<String> = sqlx::query_scalar!(
1143            // sqlx type override — MAX() strips column origin metadata so
1144            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
1145            // forces the nullable wrapper (Option<String>).
1146            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
1147               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
1148            self.service_did,
1149            uri,
1150            val,
1151        )
1152        .fetch_one(&mut **tx)
1153        .await?;
1154
1155        let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
1156
1157        let cid_owned = cid.map(str::to_string);
1158        let exp_owned = exp.map(str::to_string);
1159        let val_owned = val.to_string();
1160        let uri_owned = uri.to_string();
1161
1162        let mut label = Label {
1163            ver: 1,
1164            src: self.service_did.clone(),
1165            uri: uri_owned,
1166            cid: cid_owned,
1167            val: val_owned,
1168            neg,
1169            cts,
1170            exp: exp_owned,
1171            sig: None,
1172        };
1173        label.sig = Some(sign_label(&self.key, &label)?);
1174        let sig_bytes = label.sig.expect("just set").to_vec();
1175
1176        let neg_int: i64 = if neg { 1 } else { 0 };
1177        sqlx::query!(
1178            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
1179             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1180            seq,
1181            label.ver,
1182            label.src,
1183            label.uri,
1184            label.cid,
1185            label.val,
1186            neg_int,
1187            label.cts,
1188            label.exp,
1189            sig_bytes,
1190            self.signing_key_id,
1191            created_at_ms,
1192        )
1193        .execute(&mut **tx)
1194        .await?;
1195
1196        Ok(LabelEvent { seq, label })
1197    }
1198
1199    /// Tx-scoped wrapper around [`Self::sign_and_persist_label`] that
1200    /// converts a [`LabelDraft`] (the emission core's purpose-shaped
1201    /// output, with `cts/exp` as `SystemTime`) into the wire-level
1202    /// arguments. Used by [`Self::handle_record_action`] to sign +
1203    /// persist every draft produced by `resolve_action_labels` /
1204    /// `resolve_reason_labels`. The draft's `cts` is informational
1205    /// here — the actual cts on the labels row comes from
1206    /// `clamp_cts(created_at_ms, prev_cts)` inside
1207    /// [`Self::sign_and_persist_label`]. In v1.5 the recorder calls
1208    /// the resolvers with `now = effective_at`, so the two are
1209    /// equal pre-clamp.
1210    async fn sign_and_persist_label_from_draft(
1211        &self,
1212        tx: &mut sqlx::Transaction<'_, Sqlite>,
1213        draft: &LabelDraft,
1214        created_at_ms: i64,
1215    ) -> Result<LabelEvent> {
1216        let exp_str = match draft.exp {
1217            Some(st) => Some(rfc3339_from_epoch_ms(systemtime_to_epoch_ms(st)?)?),
1218            None => None,
1219        };
1220        self.sign_and_persist_label(
1221            tx,
1222            &draft.val,
1223            &draft.uri,
1224            draft.cid.as_deref(),
1225            draft.neg,
1226            exp_str.as_deref(),
1227            created_at_ms,
1228        )
1229        .await
1230    }
1231
1232    async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
1233        let mut tx = self.pool.begin().await?;
1234
1235        // Most-recent event for the tuple. Not-found OR latest-is-already-
1236        // a-negation both surface as LabelNotFound (§F6).
1237        let latest = sqlx::query!(
1238            "SELECT neg, cid FROM labels
1239             WHERE src = ?1 AND uri = ?2 AND val = ?3
1240             ORDER BY seq DESC LIMIT 1",
1241            self.service_did,
1242            req.uri,
1243            req.val,
1244        )
1245        .fetch_optional(&mut *tx)
1246        .await?;
1247
1248        let cid = match latest {
1249            Some(row) if row.neg == 0 => row.cid,
1250            _ => {
1251                return Err(Error::LabelNotFound {
1252                    src: self.service_did.clone(),
1253                    uri: req.uri,
1254                    val: req.val,
1255                });
1256            }
1257        };
1258
1259        let seq = reserve_seq(&mut tx).await?;
1260        // Prev cts lookup is the same query as apply; tuple uniqueness is
1261        // on (src, uri, val) regardless of neg.
1262        let prev_cts: Option<String> = sqlx::query_scalar!(
1263            // sqlx type override — MAX() strips column origin metadata so
1264            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
1265            // forces the nullable wrapper (Option<String>).
1266            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
1267               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
1268            self.service_did,
1269            req.uri,
1270            req.val,
1271        )
1272        .fetch_one(&mut *tx)
1273        .await?;
1274
1275        let wall_now_ms = epoch_ms_now();
1276        let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
1277
1278        let mut label = Label {
1279            ver: 1,
1280            src: self.service_did.clone(),
1281            uri: req.uri.clone(),
1282            cid: cid.clone(),
1283            val: req.val.clone(),
1284            neg: true,
1285            cts: cts.clone(),
1286            exp: None, // Negations never carry an expiry.
1287            sig: None,
1288        };
1289        label.sig = Some(sign_label(&self.key, &label)?);
1290        let sig_bytes = label.sig.expect("just set").to_vec();
1291
1292        let created_at = wall_now_ms;
1293        let neg_int: i64 = 1;
1294        sqlx::query!(
1295            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
1296             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1297            seq,
1298            label.ver,
1299            label.src,
1300            label.uri,
1301            label.cid,
1302            label.val,
1303            neg_int,
1304            label.cts,
1305            label.exp,
1306            sig_bytes,
1307            self.signing_key_id,
1308            created_at,
1309        )
1310        .execute(&mut *tx)
1311        .await?;
1312
1313        let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
1314        crate::audit::append::append_in_tx(
1315            &mut tx,
1316            &crate::audit::append::AuditRowForAppend {
1317                created_at,
1318                action: "label_negated".into(),
1319                actor_did: req.actor_did.clone(),
1320                target: Some(req.uri.clone()),
1321                target_cid: cid.clone(),
1322                outcome: "success".into(),
1323                reason: Some(audit_reason),
1324            },
1325        )
1326        .await?;
1327
1328        tx.commit().await?;
1329
1330        let event = LabelEvent { seq, label };
1331        let _ = self.broadcast_tx.send(event.clone());
1332        Ok(event)
1333    }
1334
1335    /// Atomic resolveReport flow — §F12 requires label INSERT +
1336    /// audit(label_applied) + report UPDATE + audit(report_resolved)
1337    /// all commit together or not at all. Everything happens inside
1338    /// one SQLite transaction on the writer task, preserving §F5's
1339    /// single-writer invariant for the optional label emission.
1340    ///
1341    /// Returns [`Error::ReportNotFound`] if the report doesn't
1342    /// exist, [`Error::ReportAlreadyResolved`] if it is not in
1343    /// `status='pending'`. Label-value validation is a handler
1344    /// concern (the handler consults `AdminConfig.label_values`
1345    /// before sending the command here — saves a writer round-trip
1346    /// on invalid input and keeps the anti-leak message local).
1347    async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
1348        use crate::report::ReportStatus;
1349
1350        let mut tx = self.pool.begin().await?;
1351
1352        // 1. Load the report; confirm pending status.
1353        let current = sqlx::query_as!(
1354            crate::report::Report,
1355            r#"SELECT
1356                 id                 AS "id!: i64",
1357                 created_at         AS "created_at!: String",
1358                 reported_by        AS "reported_by!: String",
1359                 reason_type        AS "reason_type!: String",
1360                 reason,
1361                 subject_type       AS "subject_type!: String",
1362                 subject_did        AS "subject_did!: String",
1363                 subject_uri,
1364                 subject_cid,
1365                 status             AS "status!: ReportStatus",
1366                 resolved_at,
1367                 resolved_by,
1368                 resolution_label,
1369                 resolution_reason
1370               FROM reports WHERE id = ?1"#,
1371            req.report_id,
1372        )
1373        .fetch_optional(&mut *tx)
1374        .await?;
1375
1376        let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
1377        if report.status != ReportStatus::Pending {
1378            return Err(Error::ReportAlreadyResolved { id: req.report_id });
1379        }
1380
1381        let created_at = epoch_ms_now();
1382
1383        // 2. Optional label-apply BEFORE the report UPDATE so audit
1384        // row insertion order reflects logical sequence
1385        // (label_applied, then report_resolved — §G per #15 criteria).
1386        let label_event = if let Some(apply) = req.action.as_apply() {
1387            let apply_req = ApplyLabelRequest {
1388                actor_did: req.actor_did.clone(),
1389                uri: apply.uri.clone(),
1390                cid: apply.cid.clone(),
1391                val: apply.val.clone(),
1392                exp: apply.exp.clone(),
1393                // The label's own audit_reason JSON captures {val,
1394                // neg, moderator_reason}; the resolution-level reason
1395                // goes on the report_resolved audit row below.
1396                moderator_reason: None,
1397            };
1398            Some(
1399                self.apply_label_inner(&mut tx, &apply_req, created_at)
1400                    .await?,
1401            )
1402        } else {
1403            None
1404        };
1405
1406        // 3. UPDATE reports.
1407        let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
1408        let resolution_label = req.action.as_apply().map(|a| a.val.clone());
1409        let resolved_status = ReportStatus::Resolved;
1410        sqlx::query!(
1411            "UPDATE reports SET
1412                status = ?1,
1413                resolved_at = ?2,
1414                resolved_by = ?3,
1415                resolution_label = ?4,
1416                resolution_reason = ?5
1417             WHERE id = ?6",
1418            resolved_status,
1419            resolved_at_rfc,
1420            req.actor_did,
1421            resolution_label,
1422            req.resolution_reason,
1423            req.report_id,
1424        )
1425        .execute(&mut *tx)
1426        .await?;
1427
1428        // 4. Audit: report_resolved.
1429        let audit_reason = build_resolve_audit_reason(
1430            req.action.as_apply().map(|a| a.val.as_str()),
1431            req.resolution_reason.as_deref(),
1432        );
1433        crate::audit::append::append_in_tx(
1434            &mut tx,
1435            &crate::audit::append::AuditRowForAppend {
1436                created_at,
1437                action: "report_resolved".into(),
1438                actor_did: req.actor_did.clone(),
1439                target: Some(req.report_id.to_string()),
1440                target_cid: None,
1441                outcome: "success".into(),
1442                reason: Some(audit_reason),
1443            },
1444        )
1445        .await?;
1446
1447        tx.commit().await?;
1448
1449        // 5. Broadcast (post-commit, per §F5 broadcast-after-commit rule).
1450        if let Some(event) = &label_event {
1451            let _ = self.broadcast_tx.send(event.clone());
1452        }
1453
1454        // 6. Mutate the loaded report struct to reflect the committed
1455        // state and return it. Avoids a re-fetch round-trip since we
1456        // know exactly what changed.
1457        report.status = ReportStatus::Resolved;
1458        report.resolved_at = Some(resolved_at_rfc);
1459        report.resolved_by = Some(req.actor_did);
1460        report.resolution_label = resolution_label;
1461        report.resolution_reason = req.resolution_reason;
1462
1463        Ok(ResolvedReport {
1464            report,
1465            label_event,
1466        })
1467    }
1468
1469    /// Atomic recordAction flow (§F20 / #51). Validates inputs,
1470    /// resolves the primary reason, computes strike values via the
1471    /// v1.4 calculators (#48/#49/#50/#51), inserts the audit_log
1472    /// row first (so subject_actions can carry audit_log_id at
1473    /// INSERT — the schema's no-update-except-revoke trigger
1474    /// forbids backfilling it later), inserts the subject_actions
1475    /// row, then UPSERTs subject_strike_state. All in one
1476    /// transaction.
1477    async fn handle_record_action(&self, req: RecordActionRequest) -> Result<RecordedAction> {
1478        // ---------- pre-flight validation ----------
1479
1480        if req.reason_codes.is_empty() {
1481            return Err(Error::Signing(
1482                "recordAction: reason_codes must be non-empty".into(),
1483            ));
1484        }
1485        // Note/Warning don't take reasons in the §F20 design intent,
1486        // but the issue body says "warning requires --reason" for
1487        // the audit trail. Accept reasons for all types; only
1488        // strike-bearing types actually use them for strike
1489        // calculation.
1490        match req.action_type {
1491            ActionType::TempSuspension => {
1492                if req.duration_iso.as_deref().unwrap_or("").is_empty() {
1493                    return Err(Error::DurationRequiredForTempSuspension);
1494                }
1495            }
1496            _ => {
1497                if req.duration_iso.is_some() {
1498                    return Err(Error::DurationOnlyForTempSuspension);
1499                }
1500            }
1501        }
1502
1503        let (subject_did, subject_uri) = route_subject(&req.subject)?;
1504
1505        // Resolve primary reason against the writer's vocabulary.
1506        // Vocabulary lookups are cheap; doing it before opening the
1507        // transaction keeps a bad-reason failure from acquiring the
1508        // SQLite write lock unnecessarily.
1509        let primary = resolve_primary_reason(&req.reason_codes, &self.reason_vocabulary)?;
1510
1511        // Parse duration (cheap; pre-tx for the same reason).
1512        let duration_secs = match (req.action_type, req.duration_iso.as_deref()) {
1513            (ActionType::TempSuspension, Some(iso)) => Some(parse_iso8601_duration(iso)?),
1514            _ => None,
1515        };
1516
1517        // ---------- transaction ----------
1518
1519        let mut tx = self.pool.begin().await?;
1520        let created_at = epoch_ms_now();
1521        let effective_at = created_at;
1522        let expires_at: Option<i64> = duration_secs.map(|s| effective_at + (s as i64) * 1000);
1523
1524        // Load history for strike + position calculation. id-ascending
1525        // (oldest-first) per the calculators' contract.
1526        let history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
1527
1528        // Compute current count via the decay calculator.
1529        let now_systemtime = epoch_ms_to_systemtime(created_at);
1530        let pre_state = calculate_strike_state(&history, &self.strike_policy, now_systemtime);
1531        let strikes_at_time_of_action = pre_state.current_count;
1532
1533        // Compute position and strike value.
1534        let position = compute_position_in_window(&history, &self.strike_policy, now_systemtime);
1535        let calc: StrikeApplication = strike_calculate(
1536            strikes_at_time_of_action,
1537            &primary,
1538            &self.strike_policy,
1539            position,
1540        );
1541
1542        // Note/Warning don't carry strikes regardless of reason
1543        // (§F20 design + #48 semantics). Override with zero values
1544        // and was_dampened=false; defense-in-depth for an operator
1545        // who attaches a strike-bearing reason to a Note row.
1546        let (strike_base, strike_applied, was_dampened) = match req.action_type {
1547            ActionType::Note | ActionType::Warning => (0u32, 0u32, false),
1548            _ => (calc.base_weight, calc.applied, calc.was_dampened),
1549        };
1550
1551        // Compute label-emission drafts before the audit row so the
1552        // audit reason JSON captures every (val, uri) tuple this
1553        // action will produce. The drafts are pure data — no I/O,
1554        // no signing — so committing them to the audit row before
1555        // signing is safe; the actual label INSERTs land below in
1556        // the same tx, and any failure rolls everything back.
1557        //
1558        // LabelDraft.cid is None for v1.5: subject CIDs aren't yet
1559        // plumbed through RecordActionRequest. The
1560        // ActionForEmission.cid field is reserved for the wiring a
1561        // future ticket adds via admin XRPC + CLI.
1562        let action_for_emission = ActionForEmission {
1563            action_type: req.action_type,
1564            expires_at: expires_at.map(epoch_ms_to_systemtime),
1565            subject_did: subject_did.clone(),
1566            subject_uri: subject_uri.clone(),
1567            reason_codes: req.reason_codes.clone(),
1568            cid: None,
1569        };
1570        let action_drafts = resolve_action_labels(
1571            &action_for_emission,
1572            &self.label_emission_policy,
1573            now_systemtime,
1574        );
1575        let reason_drafts = resolve_reason_labels(
1576            &action_for_emission,
1577            &self.label_emission_policy,
1578            now_systemtime,
1579        );
1580        let emitted_labels_for_audit: Vec<serde_json::Value> = action_drafts
1581            .iter()
1582            .chain(reason_drafts.iter())
1583            .map(|d| serde_json::json!({"val": d.val, "uri": d.uri}))
1584            .collect();
1585
1586        // Build audit_log reason JSON before the INSERT so the
1587        // audit row can carry the resolved values for forensic
1588        // reconstruction.
1589        // action_id is unknown until subject_actions INSERT, so we
1590        // patch it in after that step.
1591        let reason_codes_json = serde_json::to_string(&req.reason_codes)
1592            .map_err(|e| Error::Signing(format!("serialize reason_codes: {e}")))?;
1593        let report_ids_json = if req.report_ids.is_empty() {
1594            None
1595        } else {
1596            Some(
1597                serde_json::to_string(&req.report_ids)
1598                    .map_err(|e| Error::Signing(format!("serialize report_ids: {e}")))?,
1599            )
1600        };
1601
1602        // Append audit_log row first. target = subject_did (the
1603        // account being moderated); target_cid = None. We don't
1604        // know action_id yet — pass 0 as a placeholder; the actual
1605        // id is patched into the reason JSON after the
1606        // subject_actions INSERT, but only the reason JSON gets
1607        // the patch. The audit row's hash chain locks at this
1608        // INSERT, so the action_id is captured at write time.
1609        //
1610        // Order rationale: audit_log first so subject_actions
1611        // can carry audit_log_id at INSERT (the trigger forbids
1612        // UPDATE of audit_log_id). The audit row's reason JSON
1613        // captures the action_id IT will reference, computed
1614        // post-INSERT via the next-id query.
1615        let next_action_id = predict_next_subject_action_id(&mut tx).await?;
1616        let audit_reason = build_record_action_audit_reason(
1617            next_action_id,
1618            req.action_type,
1619            &primary.identifier,
1620            &req.reason_codes,
1621            strike_base,
1622            strike_applied,
1623            was_dampened,
1624            &emitted_labels_for_audit,
1625        );
1626        let audit_target = subject_uri.clone().unwrap_or_else(|| subject_did.clone());
1627        let audit_log_id = crate::audit::append::append_in_tx(
1628            &mut tx,
1629            &crate::audit::append::AuditRowForAppend {
1630                created_at,
1631                action: "subject_action_recorded".into(),
1632                actor_did: req.actor_did.clone(),
1633                target: Some(audit_target),
1634                target_cid: None,
1635                outcome: "success".into(),
1636                reason: Some(audit_reason),
1637            },
1638        )
1639        .await?;
1640
1641        // INSERT subject_actions. action_type goes through
1642        // ActionType::as_db_str so the SQL CHECK constraint binds
1643        // the same set as the Rust enum.
1644        let action_type_str = req.action_type.as_db_str();
1645        let was_dampened_int: i64 = if was_dampened { 1 } else { 0 };
1646        let strikes_at_time_i64 = strikes_at_time_of_action as i64;
1647        let strike_base_i64 = strike_base as i64;
1648        let strike_applied_i64 = strike_applied as i64;
1649        let inserted_id = sqlx::query_scalar!(
1650            "INSERT INTO subject_actions (
1651                subject_did, subject_uri, actor_did, action_type, reason_codes,
1652                duration, effective_at, expires_at, notes, report_ids,
1653                strike_value_base, strike_value_applied, was_dampened,
1654                strikes_at_time_of_action, audit_log_id, created_at
1655             ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)
1656             RETURNING id",
1657            subject_did,
1658            subject_uri,
1659            req.actor_did,
1660            action_type_str,
1661            reason_codes_json,
1662            req.duration_iso,
1663            effective_at,
1664            expires_at,
1665            req.notes,
1666            report_ids_json,
1667            strike_base_i64,
1668            strike_applied_i64,
1669            was_dampened_int,
1670            strikes_at_time_i64,
1671            audit_log_id,
1672            created_at,
1673        )
1674        .fetch_one(&mut *tx)
1675        .await?;
1676
1677        if inserted_id != next_action_id {
1678            // Audit row already committed; we'd have to invalidate
1679            // it. With BEGIN DEFERRED + the writer task being the
1680            // only in-process appender, the predict-vs-insert
1681            // race is closed — but if it ever fires, surface as
1682            // an internal error rather than a silent mismatch.
1683            return Err(Error::Signing(format!(
1684                "subject_actions inserted at id {inserted_id} but predicted {next_action_id}; audit chain captured the predicted id (corrupted)"
1685            )));
1686        }
1687
1688        // ---------- idempotency guard (#64, v1.5) ----------
1689        //
1690        // Defense-in-depth: skip emission for any state already
1691        // present on the row. v1.5's recordAction flow always
1692        // finds these queries returning the empty/NULL case (the
1693        // INSERT just landed and nothing else has touched the row
1694        // yet), so the guard is structurally a no-op in
1695        // production. It exists to protect against future paths
1696        // that might invoke emission against a row already
1697        // carrying linkage state — backfill migrations, retry
1698        // helpers, alternate code paths — and against bugs that
1699        // would otherwise silently produce duplicate (src, uri,
1700        // val) records on the wire.
1701        //
1702        // The audit row's `emitted_labels` was already written
1703        // above (intent, derived from the drafts). In v1.5 normal
1704        // flow intent matches reality. If a future defensive
1705        // scenario fires the guard and skips emissions, the audit
1706        // row will claim more emissions than the labels table
1707        // holds. That divergence is the future-implementer's
1708        // responsibility to handle (e.g., by adding a re-emission
1709        // entry point that constructs its own audit row); v1.5
1710        // accepts the limitation because the divergence is
1711        // structurally unreachable from this code path.
1712        //
1713        // The guard's "fire" branch is structurally unreachable
1714        // from the public XRPC + writer-task API in v1.5 (every
1715        // recordAction goes through this same INSERT). That makes
1716        // it untestable end-to-end without a re-emission entry
1717        // point we don't ship; the pure decision logic is
1718        // exercised via unit tests on
1719        // `should_skip_action_label_emission` /
1720        // `should_skip_reason_emission` instead.
1721        let existing_action_label_val: Option<String> = sqlx::query_scalar!(
1722            "SELECT emitted_label_uri FROM subject_actions WHERE id = ?1",
1723            inserted_id,
1724        )
1725        .fetch_one(&mut *tx)
1726        .await?;
1727        let existing_reason_codes: Vec<String> = sqlx::query_scalar!(
1728            "SELECT reason_code FROM subject_action_reason_labels WHERE action_id = ?1",
1729            inserted_id,
1730        )
1731        .fetch_all(&mut *tx)
1732        .await?;
1733        let existing_reason_set: HashSet<&str> =
1734            existing_reason_codes.iter().map(String::as_str).collect();
1735
1736        // ---------- emission (#60, v1.5) ----------
1737        //
1738        // Sign + persist each draft (action label first so its
1739        // seq < reason labels'), then UPDATE
1740        // subject_actions.emitted_label_uri, then write per-reason
1741        // linkage rows. All in the same tx as the action INSERT
1742        // and the audit row. Failure here rolls back everything
1743        // up to and including the audit row, so the audit chain
1744        // never claims emission that didn't happen.
1745        //
1746        // Empty action_drafts (note / suppressed warning / policy
1747        // disabled) cleanly short-circuits both loops, leaves
1748        // emitted_label_uri NULL, and writes no reason-label
1749        // linkage rows — the action still records, just without
1750        // an ATProto label tail.
1751        let mut label_events: Vec<LabelEvent> =
1752            Vec::with_capacity(action_drafts.len() + reason_drafts.len());
1753        let mut action_label_val: Option<String> = None;
1754        if !should_skip_action_label_emission(existing_action_label_val.as_deref()) {
1755            for draft in &action_drafts {
1756                let event = self
1757                    .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
1758                    .await?;
1759                if action_label_val.is_none() {
1760                    action_label_val = Some(draft.val.clone());
1761                }
1762                label_events.push(event);
1763            }
1764
1765            if let Some(ref val) = action_label_val {
1766                // emitted_label_uri stores the action label's `val`,
1767                // not a URI: ATProto labels have no canonical URI; the
1768                // discriminator within (src=service_did, uri=subject,
1769                // val) is just the val. Column name predates the
1770                // realization (#57); locked once shipped.
1771                sqlx::query!(
1772                    "UPDATE subject_actions SET emitted_label_uri = ?1 WHERE id = ?2",
1773                    val,
1774                    inserted_id,
1775                )
1776                .execute(&mut *tx)
1777                .await?;
1778            }
1779        }
1780
1781        // Per-(action, reason) linkage rows. resolve_reason_labels
1782        // emits drafts in `reason_codes` order, so a parallel
1783        // iteration recovers the source reason_code without
1784        // re-parsing the label val. Reasons whose linkage row
1785        // already exists are skipped per the idempotency guard
1786        // above.
1787        for (draft, reason_code) in reason_drafts.iter().zip(req.reason_codes.iter()) {
1788            if should_skip_reason_emission(reason_code, &existing_reason_set) {
1789                continue;
1790            }
1791            let event = self
1792                .sign_and_persist_label_from_draft(&mut tx, draft, effective_at)
1793                .await?;
1794            sqlx::query!(
1795                "INSERT INTO subject_action_reason_labels
1796                   (action_id, reason_code, emitted_label_uri, emitted_at)
1797                 VALUES (?1, ?2, ?3, ?4)",
1798                inserted_id,
1799                reason_code,
1800                draft.val,
1801                effective_at,
1802            )
1803            .execute(&mut *tx)
1804            .await?;
1805            label_events.push(event);
1806        }
1807
1808        // Recompute strike state for the cache. Reload history with
1809        // the new row so the cache reflects post-insert reality.
1810        let post_history = load_subject_actions_for_calc(&mut tx, &subject_did).await?;
1811        let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
1812
1813        let post_count_i64 = post_state.current_count as i64;
1814        sqlx::query!(
1815            "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
1816             VALUES (?1, ?2, ?3, ?3)
1817             ON CONFLICT(subject_did) DO UPDATE SET
1818                 current_strike_count = excluded.current_strike_count,
1819                 last_action_at = excluded.last_action_at,
1820                 last_recompute_at = excluded.last_recompute_at",
1821            subject_did,
1822            post_count_i64,
1823            created_at,
1824        )
1825        .execute(&mut *tx)
1826        .await?;
1827
1828        tx.commit().await?;
1829
1830        // Broadcast each emitted label to subscribeLabels consumers.
1831        // Order matches the persistence order (action label first,
1832        // then reasons). No-receivers is not a write failure
1833        // (§plan point G).
1834        for event in label_events {
1835            let _ = self.broadcast_tx.send(event);
1836        }
1837
1838        Ok(RecordedAction {
1839            action_id: inserted_id,
1840            strike_value_base: strike_base,
1841            strike_value_applied: strike_applied,
1842            was_dampened,
1843            strikes_at_time_of_action,
1844        })
1845    }
1846
1847    /// Atomic revokeAction flow (§F20 / #51). Looks up the row,
1848    /// verifies state, sets the revoked_* columns (the schema's
1849    /// no-update-except-revoke trigger permits this NULL→non-NULL
1850    /// transition), recomputes strike_state for the subject, and
1851    /// appends an audit_log row. All in one transaction.
1852    async fn handle_revoke_action(&self, req: RevokeActionRequest) -> Result<RevokedAction> {
1853        let mut tx = self.pool.begin().await?;
1854
1855        // Look up the row + check current state. Pull the fields
1856        // needed for both v1.4 revocation (subject_did, revoked_at)
1857        // and v1.5 #62 negation (subject_uri, emitted_label_uri).
1858        let row = sqlx::query!(
1859            "SELECT subject_did, subject_uri, emitted_label_uri, revoked_at
1860             FROM subject_actions WHERE id = ?1",
1861            req.action_id,
1862        )
1863        .fetch_optional(&mut *tx)
1864        .await?;
1865        let row = row.ok_or(Error::ActionNotFound(req.action_id))?;
1866        if row.revoked_at.is_some() {
1867            return Err(Error::ActionAlreadyRevoked(req.action_id));
1868        }
1869
1870        let revoked_at_ms = epoch_ms_now();
1871
1872        // UPDATE the revocation columns. The schema trigger permits
1873        // NULL→non-NULL on these three columns; any other column
1874        // change in this UPDATE would abort.
1875        sqlx::query!(
1876            "UPDATE subject_actions
1877             SET revoked_at = ?1, revoked_by_did = ?2, revoked_reason = ?3
1878             WHERE id = ?4",
1879            revoked_at_ms,
1880            req.revoked_by_did,
1881            req.revoked_reason,
1882            req.action_id,
1883        )
1884        .execute(&mut *tx)
1885        .await?;
1886
1887        // ---------- negation (#62, v1.5) ----------
1888        //
1889        // Read the action's emitted-label linkage from storage
1890        // (NOT from current policy — see module docs for the
1891        // val-from-storage rule). For each emitted val, sign and
1892        // persist a fresh label row with neg=true targeting the
1893        // same (src, uri, val) tuple. The original rows stay in
1894        // place; ATProto consumers honor the latest record per
1895        // tuple and so see the negation. Reason linkage rows
1896        // (subject_action_reason_labels) are PRESERVED — they're
1897        // forensic record of "at this point these labels were
1898        // emitted," not a cache of what's currently in force.
1899        //
1900        // exp = None on every negation. The original temp_suspension
1901        // label may have carried an expiry (its exp said "stop
1902        // honoring me at this wall-clock"), but the negation itself
1903        // is a permanent statement that supersedes the original;
1904        // expiring the negation would resurrect the original label
1905        // in consumer caches.
1906        //
1907        // Action that was never emitted (note, suppressed warning,
1908        // emission disabled at recording time) → both the action
1909        // label fetch and the reason-label fetch yield zero rows
1910        // and the negation step is a no-op. The revocation
1911        // audit row still lands with negated_labels = [].
1912        let label_uri = row
1913            .subject_uri
1914            .clone()
1915            .unwrap_or_else(|| row.subject_did.clone());
1916
1917        // Reason-label linkage rows ordered by reason_code for
1918        // deterministic audit shape. Original emission order isn't
1919        // recoverable (storage doesn't preserve req.reason_codes
1920        // ordering), and ATProto consumers don't care about
1921        // negation order; alphabetical is the cheapest stable
1922        // ordering for forensic readers.
1923        let reason_linkage = sqlx::query!(
1924            "SELECT reason_code, emitted_label_uri
1925             FROM subject_action_reason_labels
1926             WHERE action_id = ?1
1927             ORDER BY reason_code ASC",
1928            req.action_id,
1929        )
1930        .fetch_all(&mut *tx)
1931        .await?;
1932
1933        let mut negated_for_audit: Vec<serde_json::Value> = Vec::new();
1934        if let Some(ref val) = row.emitted_label_uri {
1935            negated_for_audit.push(serde_json::json!({
1936                "val": val,
1937                "uri": label_uri,
1938            }));
1939        }
1940        for r in &reason_linkage {
1941            negated_for_audit.push(serde_json::json!({
1942                "val": r.emitted_label_uri,
1943                "uri": label_uri,
1944            }));
1945        }
1946
1947        // Audit row first — captures the negation set in the hash
1948        // chain before the label INSERTs land. Mirrors #60's
1949        // emission-side ordering: a forensic reader who sees the
1950        // audit row knows what labels SHOULD have been written;
1951        // the labels table is the witness. tx atomicity guarantees
1952        // they agree.
1953        let audit_reason = build_revoke_action_audit_reason(
1954            req.action_id,
1955            req.revoked_reason.as_deref(),
1956            &negated_for_audit,
1957        );
1958        crate::audit::append::append_in_tx(
1959            &mut tx,
1960            &crate::audit::append::AuditRowForAppend {
1961                created_at: revoked_at_ms,
1962                action: "subject_action_revoked".into(),
1963                actor_did: req.revoked_by_did.clone(),
1964                target: Some(req.action_id.to_string()),
1965                target_cid: None,
1966                outcome: "success".into(),
1967                reason: Some(audit_reason),
1968            },
1969        )
1970        .await?;
1971
1972        // Sign + persist each negation label. Action label first so
1973        // its seq < reason negations'.
1974        let mut negation_events: Vec<LabelEvent> = Vec::with_capacity(1 + reason_linkage.len());
1975        if let Some(ref val) = row.emitted_label_uri {
1976            let event = self
1977                .sign_and_persist_label(
1978                    &mut tx,
1979                    val,
1980                    &label_uri,
1981                    None, // cid: not plumbed in v1.5
1982                    true, // neg
1983                    None, // exp: negations don't expire
1984                    revoked_at_ms,
1985                )
1986                .await?;
1987            negation_events.push(event);
1988        }
1989        for r in &reason_linkage {
1990            let event = self
1991                .sign_and_persist_label(
1992                    &mut tx,
1993                    &r.emitted_label_uri,
1994                    &label_uri,
1995                    None,
1996                    true,
1997                    None,
1998                    revoked_at_ms,
1999                )
2000                .await?;
2001            negation_events.push(event);
2002        }
2003
2004        // Recompute strike state for the subject. The decay
2005        // calculator excludes revoked rows from current_count.
2006        let post_history = load_subject_actions_for_calc(&mut tx, &row.subject_did).await?;
2007        let now_systemtime = epoch_ms_to_systemtime(revoked_at_ms);
2008        let post_state = calculate_strike_state(&post_history, &self.strike_policy, now_systemtime);
2009
2010        let post_count_i64 = post_state.current_count as i64;
2011        sqlx::query!(
2012            "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
2013             VALUES (?1, ?2, ?3, ?3)
2014             ON CONFLICT(subject_did) DO UPDATE SET
2015                 current_strike_count = excluded.current_strike_count,
2016                 last_recompute_at = excluded.last_recompute_at",
2017            row.subject_did,
2018            post_count_i64,
2019            revoked_at_ms,
2020        )
2021        .execute(&mut *tx)
2022        .await?;
2023
2024        tx.commit().await?;
2025
2026        // Broadcast each negation to subscribeLabels consumers
2027        // post-commit. Same ordering and no-receivers-is-fine
2028        // posture as the emission path (§plan point G).
2029        for event in negation_events {
2030            let _ = self.broadcast_tx.send(event);
2031        }
2032
2033        Ok(RevokedAction {
2034            action_id: req.action_id,
2035            revoked_at: rfc3339_from_epoch_ms(revoked_at_ms)?,
2036        })
2037    }
2038
2039    async fn heartbeat(&self) -> Result<()> {
2040        let now_ms = epoch_ms_now();
2041        sqlx::query!(
2042            "UPDATE server_instance_lease SET last_heartbeat = ?1
2043             WHERE id = 1 AND instance_id = ?2",
2044            now_ms,
2045            self.instance_id,
2046        )
2047        .execute(&self.pool)
2048        .await?;
2049        Ok(())
2050    }
2051
2052    async fn release_lease(&self) -> Result<()> {
2053        sqlx::query!(
2054            "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
2055            self.instance_id,
2056        )
2057        .execute(&self.pool)
2058        .await?;
2059        Ok(())
2060    }
2061}
2062
2063// ---------- Lease + key bootstrap ----------
2064
2065/// Acquire the §F5 single-instance lease.
2066///
2067/// Returns the new instance id on success. Returns
2068/// [`Error::LeaseHeld`] when an existing row's `last_heartbeat` is
2069/// still within [`LEASE_STALE_MS`] — i.e., another writer is live.
2070///
2071/// Exposed as `pub(crate)` so one-shot CLI tools that mutate
2072/// `audit_log` (`cairn audit-rebuild`, #40) can assert the same
2073/// "no live writer" invariant the long-running serve process holds.
2074/// The CLI tool releases via [`release_lease_by_id`] when done.
2075pub(crate) async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
2076    let now_ms = epoch_ms_now();
2077
2078    let existing =
2079        sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
2080            .fetch_optional(pool)
2081            .await?;
2082
2083    if let Some(row) = existing {
2084        let age_ms = (now_ms - row.last_heartbeat).max(0);
2085        if age_ms < LEASE_STALE_MS {
2086            return Err(Error::LeaseHeld {
2087                instance_id: row.instance_id,
2088                age_secs: (age_ms / 1000) as u64,
2089            });
2090        }
2091    }
2092
2093    let new_id = Uuid::new_v4().to_string();
2094    // INSERT OR REPLACE on id=1: takes over a stale lease row if present,
2095    // or creates the row on first startup. Legal because the staleness
2096    // check above confirms no live peer.
2097    sqlx::query!(
2098        "INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
2099         VALUES (1, ?1, ?2, ?2)
2100         ON CONFLICT(id) DO UPDATE SET
2101             instance_id = excluded.instance_id,
2102             acquired_at = excluded.acquired_at,
2103             last_heartbeat = excluded.last_heartbeat",
2104        new_id,
2105        now_ms,
2106    )
2107    .execute(pool)
2108    .await?;
2109
2110    Ok(new_id)
2111}
2112
2113/// Release a lease acquired via [`acquire_lease`] by its instance id.
2114/// Free-function counterpart to [`Writer::release_lease`] for callers
2115/// that don't own a `Writer` (one-shot CLI tools that just hold the
2116/// lease for the duration of a data migration).
2117pub(crate) async fn release_lease_by_id(pool: &Pool<Sqlite>, instance_id: &str) -> Result<()> {
2118    sqlx::query!(
2119        "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
2120        instance_id,
2121    )
2122    .execute(pool)
2123    .await?;
2124    Ok(())
2125}
2126
2127async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
2128    let kp = K256Keypair::from_private_key(key.expose_secret())?;
2129    let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
2130
2131    let existing =
2132        sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
2133            .fetch_optional(pool)
2134            .await?;
2135
2136    if let Some(row) = existing {
2137        if row.public_key_multibase != my_multibase {
2138            return Err(Error::Signing(format!(
2139                "signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
2140                row.public_key_multibase, my_multibase
2141            )));
2142        }
2143        return Ok(row.id);
2144    }
2145
2146    let now_ms = epoch_ms_now();
2147    let valid_from = rfc3339_from_epoch_ms(now_ms)?;
2148    let id = sqlx::query_scalar!(
2149        "INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
2150         VALUES (?1, ?2, NULL, ?3)
2151         RETURNING id",
2152        my_multibase,
2153        valid_from,
2154        now_ms,
2155    )
2156    .fetch_one(pool)
2157    .await?;
2158
2159    Ok(id)
2160}
2161
2162// ---------- Pure helpers (unit-testable without a DB) ----------
2163
2164async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
2165    // label_sequence has a single auto-increment column; `DEFAULT VALUES`
2166    // triggers a fresh allocation. Under AUTOINCREMENT the assigned seq
2167    // is never reused even across rolled-back transactions — this is what
2168    // gives us "strictly monotonic, no gaps from writer's perspective."
2169    let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
2170        .fetch_one(tx)
2171        .await?;
2172    Ok(seq)
2173}
2174
2175/// §6.1 monotonicity clamp. `prev_cts_str`, when present, is expected to
2176/// be in the writer's canonical `CTS_FORMAT`.
2177fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
2178    let effective_ms = match prev_cts_str {
2179        Some(s) => {
2180            let prev_ms = parse_rfc3339_ms(s)?;
2181            wall_now_ms.max(prev_ms + 1)
2182        }
2183        None => wall_now_ms,
2184    };
2185    rfc3339_from_epoch_ms(effective_ms)
2186}
2187
2188/// Epoch-ms → RFC-3339 Z with millisecond precision. `pub(crate)` so
2189/// peer modules (admin/audit_view, future retention sweep) share the
2190/// single formatter used on the writer's timestamp boundary.
2191pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
2192    let nanos: i128 = (ms as i128) * 1_000_000;
2193    let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
2194        .map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
2195    let formatted = dt
2196        .format(&CTS_FORMAT)
2197        .map_err(|e| Error::Signing(format!("format cts: {e}")))?;
2198    Ok(format!("{formatted}Z"))
2199}
2200
2201/// RFC-3339 Z with millisecond precision → epoch-ms. `pub(crate)` so
2202/// admin handlers can validate `since` / `until` query params against
2203/// the same parser the writer uses on its input boundary.
2204pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
2205    let stripped = s
2206        .strip_suffix('Z')
2207        .ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
2208    let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
2209        .map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
2210    let nanos = pdt.assume_utc().unix_timestamp_nanos();
2211    Ok((nanos / 1_000_000) as i64)
2212}
2213
2214/// Current wall-clock time as Unix epoch milliseconds. `pub(crate)` so
2215/// peer modules (server, future retention sweep) share one implementation
2216/// rather than each inlining `SystemTime::now()` with its own error
2217/// handling.
2218pub(crate) fn epoch_ms_now() -> i64 {
2219    SystemTime::now()
2220        .duration_since(UNIX_EPOCH)
2221        .expect("system clock before unix epoch")
2222        .as_millis() as i64
2223}
2224
2225fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
2226    let body = serde_json::json!({
2227        "val": val,
2228        "neg": neg,
2229        "moderator_reason": moderator_reason,
2230    });
2231    body.to_string()
2232}
2233
2234/// `report_resolved` audit reason JSON — see [`AUDIT_REASON_RESOLVE_REPORT`]
2235/// for the schema.
2236fn build_resolve_audit_reason(
2237    applied_label_val: Option<&str>,
2238    resolution_reason: Option<&str>,
2239) -> String {
2240    serde_json::json!({
2241        "applied_label_val": applied_label_val,
2242        "resolution_reason": resolution_reason,
2243    })
2244    .to_string()
2245}
2246
2247/// `retention_sweep` audit reason JSON — see
2248/// [`AUDIT_REASON_RETENTION_SWEEP`] for the schema. Shared with the
2249/// retentionSweep admin handler (audited only on the operator-
2250/// initiated path per Q6/D2).
2251pub(crate) fn build_retention_sweep_audit_reason(result: &SweepResult) -> String {
2252    serde_json::json!({
2253        "rows_deleted": result.rows_deleted,
2254        "batches": result.batches,
2255        "duration_ms": result.duration_ms,
2256        "retention_days_applied": result.retention_days_applied,
2257    })
2258    .to_string()
2259}
2260
2261/// Audit-log `reason` JSON schema for `subject_action_recorded`
2262/// (§F20 / #51 graduated-action moderation). Captures the action_id,
2263/// action_type, primary reason resolved at record time, the full
2264/// list of declared reasons, the strike value applied, and whether
2265/// dampening fired — enough for forensic reconstruction without
2266/// requiring a join to subject_actions.
2267///
2268/// ```json
2269/// {
2270///   "action_id": <i64>,
2271///   "action_type": "<warning|note|temp_suspension|indef_suspension|takedown>",
2272///   "primary_reason": "<identifier>",
2273///   "reason_codes": ["<id1>", "<id2>"],
2274///   "strike_value_base": <u32>,
2275///   "strike_value_applied": <u32>,
2276///   "was_dampened": <bool>
2277/// }
2278/// ```
2279#[doc(alias = "audit_log.reason.subject_action_recorded")]
2280pub 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 }";
2281
2282/// Audit-log `reason` JSON schema for `subject_action_revoked`
2283/// (§F20 / #51). Captures the revoked subject_actions row id and
2284/// the moderator-supplied rationale.
2285///
2286/// ```json
2287/// {
2288///   "action_id": <i64>,
2289///   "revoked_reason": "<free text>" | null
2290/// }
2291/// ```
2292#[doc(alias = "audit_log.reason.subject_action_revoked")]
2293pub const AUDIT_REASON_REVOKE_ACTION: &str =
2294    "subject_action_revoked: { action_id, revoked_reason }";
2295
2296/// `reporter_flagged` / `reporter_unflagged` audit reason JSON —
2297/// see [`AUDIT_REASON_FLAG_REPORTER`]. Shared with the flagReporter
2298/// handler (not the writer — flagReporter is a direct handler txn
2299/// per the #15 criteria).
2300pub(crate) fn build_flag_reporter_audit_reason(
2301    did: &str,
2302    suppressed: bool,
2303    moderator_reason: Option<&str>,
2304) -> String {
2305    serde_json::json!({
2306        "did": did,
2307        "suppressed": suppressed,
2308        "moderator_reason": moderator_reason,
2309    })
2310    .to_string()
2311}
2312
2313/// Idempotency guard for the recorder's action-label emission
2314/// step (#64). Returns `true` when a row already carries an
2315/// emitted action-label `val`, signalling the emission loop
2316/// should skip producing a duplicate (src, uri, val) record.
2317///
2318/// v1.5's [`Writer::handle_record_action`] always passes `None`
2319/// here in production (the row was just INSERTed and nothing
2320/// else has touched it inside the same transaction). The
2321/// function exists so future paths that operate on existing
2322/// rows — backfill migrations, retry helpers — can call into
2323/// the same gate logic, and so the gate's behavior is testable
2324/// without DB scaffolding.
2325fn should_skip_action_label_emission(existing_emitted_label_uri: Option<&str>) -> bool {
2326    existing_emitted_label_uri.is_some()
2327}
2328
2329/// Idempotency guard for the recorder's reason-label emission
2330/// step (#64). Returns `true` when a `subject_action_reason_labels`
2331/// row already exists for `(action_id, reason_code)`, signalling
2332/// the per-reason emission loop should skip this draft.
2333///
2334/// Same v1.5 production posture as
2335/// [`should_skip_action_label_emission`]: `existing` is always
2336/// empty when called from [`Writer::handle_record_action`]; the
2337/// function is the gate logic future paths can reuse and tests
2338/// can exercise.
2339fn should_skip_reason_emission(reason_code: &str, existing: &HashSet<&str>) -> bool {
2340    existing.contains(reason_code)
2341}
2342
2343/// `subject_action_recorded` audit reason JSON — see
2344/// [`AUDIT_REASON_RECORD_ACTION`] for the schema.
2345#[allow(clippy::too_many_arguments)]
2346fn build_record_action_audit_reason(
2347    action_id: i64,
2348    action_type: ActionType,
2349    primary_reason: &str,
2350    reason_codes: &[String],
2351    strike_value_base: u32,
2352    strike_value_applied: u32,
2353    was_dampened: bool,
2354    emitted_labels: &[serde_json::Value],
2355) -> String {
2356    serde_json::json!({
2357        "action_id": action_id,
2358        "action_type": action_type.as_db_str(),
2359        "primary_reason": primary_reason,
2360        "reason_codes": reason_codes,
2361        "strike_value_base": strike_value_base,
2362        "strike_value_applied": strike_value_applied,
2363        "was_dampened": was_dampened,
2364        // v1.5 (#60): every emitted label as `{val, uri}`. Empty
2365        // when [label_emission].enabled = false, when the action
2366        // type is `note`, or when `warning` with the default
2367        // suppression gate. Captured pre-INSERT so the audit
2368        // hash chain locks the entire (action, labels) bundle
2369        // even though the actual label INSERTs land later in the
2370        // same tx.
2371        "emitted_labels": emitted_labels,
2372    })
2373    .to_string()
2374}
2375
2376/// `subject_action_revoked` audit reason JSON — see
2377/// [`AUDIT_REASON_REVOKE_ACTION`] for the schema.
2378///
2379/// `negated_labels` mirrors the emission-side `emitted_labels`
2380/// shape from [`build_record_action_audit_reason`]: an array of
2381/// `{val, uri}` entries, action label first followed by reason
2382/// labels in alphabetical order. Empty when the revoked action
2383/// had no emitted labels (note, suppressed warning, emission
2384/// disabled at record time, action recorded pre-v1.5).
2385fn build_revoke_action_audit_reason(
2386    action_id: i64,
2387    revoked_reason: Option<&str>,
2388    negated_labels: &[serde_json::Value],
2389) -> String {
2390    serde_json::json!({
2391        "action_id": action_id,
2392        "revoked_reason": revoked_reason,
2393        "negated_labels": negated_labels,
2394    })
2395    .to_string()
2396}
2397
2398/// Convert epoch-ms (i64) to [`SystemTime`]. The v1.4 calculators
2399/// take SystemTime; the schema stores epoch-ms; this is the
2400/// boundary helper. Negative values clamp to UNIX_EPOCH (defense:
2401/// schema columns are non-negative in practice).
2402fn epoch_ms_to_systemtime(ms: i64) -> SystemTime {
2403    if ms >= 0 {
2404        UNIX_EPOCH + Duration::from_millis(ms as u64)
2405    } else {
2406        UNIX_EPOCH
2407    }
2408}
2409
2410/// Inverse of [`epoch_ms_to_systemtime`]. Used by the emission
2411/// path (#60) to translate a [`LabelDraft`]'s `exp: SystemTime`
2412/// into the i64 ms that [`rfc3339_from_epoch_ms`] formats. Errors
2413/// when the value pre-dates UNIX_EPOCH — the emission core never
2414/// produces such a value (caller-supplied `expires_at` is
2415/// epoch-ms-derived), so this should be unreachable in practice.
2416fn systemtime_to_epoch_ms(st: SystemTime) -> Result<i64> {
2417    st.duration_since(UNIX_EPOCH)
2418        .map(|d| d.as_millis() as i64)
2419        .map_err(|e| Error::Signing(format!("SystemTime before unix epoch: {e}")))
2420}
2421
2422/// Route a subject string into (subject_did, subject_uri) per the
2423/// recorder's contract. DIDs go to subject_did with no URI; AT-URIs
2424/// extract the repo DID as subject_did and keep the full URI as
2425/// subject_uri. Anything else is rejected.
2426fn route_subject(subject: &str) -> Result<(String, Option<String>)> {
2427    if let Some(rest) = subject.strip_prefix("at://") {
2428        let repo = rest.split('/').next().unwrap_or("");
2429        if repo.is_empty() || !repo.starts_with("did:") {
2430            return Err(Error::Signing(format!(
2431                "subject at://-URI must start at://did:.../...; got {subject:?}"
2432            )));
2433        }
2434        Ok((repo.to_string(), Some(subject.to_string())))
2435    } else if subject.starts_with("did:") && subject.len() > "did:".len() {
2436        Ok((subject.to_string(), None))
2437    } else {
2438        Err(Error::Signing(format!(
2439            "subject must be a DID (`did:...`) or AT-URI (`at://did:...`); got {subject:?}"
2440        )))
2441    }
2442}
2443
2444/// Parse a narrow ISO-8601 duration subset into a [`Duration`].
2445/// Supports `P{n}D` (days), `P{n}W` (weeks), `PT{n}H` (hours),
2446/// `PT{n}M` (minutes), `PT{n}S` (seconds), and combinations like
2447/// `P1DT12H`. Years and months are NOT supported — moderation
2448/// suspensions are bounded enough that day/hour granularity covers
2449/// real cases without the calendar-arithmetic complexity of Y/M.
2450///
2451/// Returns the duration in seconds.
2452fn parse_iso8601_duration(s: &str) -> Result<u64> {
2453    let body = s.strip_prefix('P').ok_or_else(|| {
2454        Error::Signing(format!(
2455            "duration {s:?} must start with 'P' (ISO-8601 duration form, e.g. P7D)"
2456        ))
2457    })?;
2458    if body.is_empty() {
2459        return Err(Error::Signing(format!(
2460            "duration {s:?} has no components after 'P'"
2461        )));
2462    }
2463
2464    // Split into date part (before T, if any) and time part.
2465    let (date_part, time_part) = match body.split_once('T') {
2466        Some((d, t)) => (d, Some(t)),
2467        None => (body, None),
2468    };
2469
2470    let mut total_secs: u64 = 0;
2471    let mut buf = String::new();
2472    for c in date_part.chars() {
2473        if c.is_ascii_digit() {
2474            buf.push(c);
2475            continue;
2476        }
2477        let n: u64 = buf
2478            .parse()
2479            .map_err(|_| Error::Signing(format!("duration {s:?}: malformed numeric run")))?;
2480        buf.clear();
2481        let mult = match c {
2482            'D' => 86_400u64,
2483            'W' => 7 * 86_400u64,
2484            'Y' | 'M' => {
2485                return Err(Error::Signing(format!(
2486                    "duration {s:?}: years (Y) and months (M, without T prefix) not supported in v1.4"
2487                )));
2488            }
2489            other => {
2490                return Err(Error::Signing(format!(
2491                    "duration {s:?}: unknown date-part unit {other:?}"
2492                )));
2493            }
2494        };
2495        total_secs = total_secs.saturating_add(n.saturating_mul(mult));
2496    }
2497    if !buf.is_empty() {
2498        return Err(Error::Signing(format!(
2499            "duration {s:?}: trailing digits without unit in date part"
2500        )));
2501    }
2502
2503    if let Some(time_body) = time_part {
2504        if time_body.is_empty() {
2505            return Err(Error::Signing(format!(
2506                "duration {s:?}: 'T' separator with no time components"
2507            )));
2508        }
2509        for c in time_body.chars() {
2510            if c.is_ascii_digit() {
2511                buf.push(c);
2512                continue;
2513            }
2514            let n: u64 = buf.parse().map_err(|_| {
2515                Error::Signing(format!(
2516                    "duration {s:?}: malformed numeric run in time part"
2517                ))
2518            })?;
2519            buf.clear();
2520            let mult = match c {
2521                'H' => 3600u64,
2522                'M' => 60u64,
2523                'S' => 1u64,
2524                other => {
2525                    return Err(Error::Signing(format!(
2526                        "duration {s:?}: unknown time-part unit {other:?}"
2527                    )));
2528                }
2529            };
2530            total_secs = total_secs.saturating_add(n.saturating_mul(mult));
2531        }
2532        if !buf.is_empty() {
2533            return Err(Error::Signing(format!(
2534                "duration {s:?}: trailing digits without unit in time part"
2535            )));
2536        }
2537    }
2538
2539    if total_secs == 0 {
2540        return Err(Error::Signing(format!(
2541            "duration {s:?}: parsed to zero — must be at least one second"
2542        )));
2543    }
2544    Ok(total_secs)
2545}
2546
2547/// Predict the next AUTOINCREMENT id for `subject_actions`. SQLite's
2548/// AUTOINCREMENT columns are backed by `sqlite_sequence`; the next
2549/// value is `max(seq, max(rowid)) + 1`. With BEGIN DEFERRED + the
2550/// writer task being the only in-process appender, the predict
2551/// step is racy only against external processes — and the ROLLBACK
2552/// path catches mismatch via the inserted_id check in
2553/// handle_record_action.
2554async fn predict_next_subject_action_id(tx: &mut sqlx::Transaction<'_, Sqlite>) -> Result<i64> {
2555    // sqlite_sequence.seq has no NOT NULL constraint; the type
2556    // override forces a non-nullable i64 because the column is
2557    // populated for every AUTOINCREMENT table that has had at
2558    // least one row inserted (NULL only happens if the row
2559    // doesn't exist, which fetch_optional handles).
2560    let row = sqlx::query!(
2561        r#"SELECT seq AS "seq!: i64" FROM sqlite_sequence WHERE name = 'subject_actions'"#
2562    )
2563    .fetch_optional(&mut **tx)
2564    .await?;
2565    Ok(row.map(|r| r.seq + 1).unwrap_or(1))
2566}
2567
2568/// Load subject_actions rows for a subject, oldest-first, projecting
2569/// to [`ActionRecord`] for the v1.4 calculators (decay #50,
2570/// window #51). Filters to the columns the calculators consume —
2571/// the rest stay in the DB.
2572async fn load_subject_actions_for_calc(
2573    tx: &mut sqlx::Transaction<'_, Sqlite>,
2574    subject_did: &str,
2575) -> Result<Vec<ActionRecord>> {
2576    let rows = sqlx::query!(
2577        "SELECT action_type, strike_value_applied, was_dampened,
2578                effective_at, expires_at, revoked_at
2579         FROM subject_actions
2580         WHERE subject_did = ?1
2581         ORDER BY id ASC",
2582        subject_did,
2583    )
2584    .fetch_all(&mut **tx)
2585    .await?;
2586
2587    let mut out = Vec::with_capacity(rows.len());
2588    for r in rows {
2589        let action_type = ActionType::from_db_str(&r.action_type).ok_or_else(|| {
2590            Error::Signing(format!(
2591                "subject_actions row has invalid action_type {:?}",
2592                r.action_type
2593            ))
2594        })?;
2595        let strike_value_applied = u32::try_from(r.strike_value_applied).map_err(|_| {
2596            Error::Signing(format!(
2597                "subject_actions row strike_value_applied {} out of u32 range",
2598                r.strike_value_applied
2599            ))
2600        })?;
2601        out.push(ActionRecord {
2602            strike_value_applied,
2603            effective_at: epoch_ms_to_systemtime(r.effective_at),
2604            revoked_at: r.revoked_at.map(epoch_ms_to_systemtime),
2605            action_type,
2606            expires_at: r.expires_at.map(epoch_ms_to_systemtime),
2607            was_dampened: r.was_dampened != 0,
2608        });
2609    }
2610    Ok(out)
2611}
2612
2613#[cfg(test)]
2614mod tests {
2615    use super::*;
2616
2617    // 1_715_000_000_000 ms since epoch = 2024-05-06T12:53:20.000Z UTC.
2618    // The four tests below pin every branch of the clamp: no prior,
2619    // wall-ahead, wall-behind (use prev+1ms), wall-equal (advance anyway).
2620
2621    #[test]
2622    fn clamp_cts_uses_wall_clock_when_no_prior() {
2623        let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
2624        assert_eq!(out, "2024-05-06T12:53:20.000Z");
2625    }
2626
2627    #[test]
2628    fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
2629        // wall_clock one second behind prev: result is prev + 1ms, not wall.
2630        let prev = "2024-05-06T12:53:20.500Z";
2631        let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
2632        assert_eq!(out, "2024-05-06T12:53:20.501Z");
2633    }
2634
2635    #[test]
2636    fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
2637        let prev = "2024-05-06T12:53:20.500Z";
2638        let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
2639        assert_eq!(out, "2024-05-06T12:53:22.000Z");
2640    }
2641
2642    #[test]
2643    fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
2644        // wall_clock == prev exactly: must advance by 1ms (strictly greater).
2645        let prev = "2024-05-06T12:53:20.500Z";
2646        let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
2647        assert_eq!(out, "2024-05-06T12:53:20.501Z");
2648    }
2649
2650    #[test]
2651    fn build_audit_reason_shape_matches_documented_schema() {
2652        let json = build_audit_reason("spam", false, Some("user reported"));
2653        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
2654        assert_eq!(v["val"], "spam");
2655        assert_eq!(v["neg"], false);
2656        assert_eq!(v["moderator_reason"], "user reported");
2657    }
2658
2659    #[test]
2660    fn build_audit_reason_null_when_moderator_reason_absent() {
2661        let json = build_audit_reason("spam", true, None);
2662        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
2663        assert!(v["moderator_reason"].is_null());
2664    }
2665
2666    #[test]
2667    fn rfc3339_roundtrip() {
2668        let s = "2026-04-22T12:00:00.938Z";
2669        let ms = parse_rfc3339_ms(s).expect("parse");
2670        let back = rfc3339_from_epoch_ms(ms).expect("format");
2671        assert_eq!(back, s);
2672    }
2673
2674    // ---------- route_subject (#51) ----------
2675
2676    #[test]
2677    fn route_subject_did_only_returns_did_no_uri() {
2678        let (did, uri) = route_subject("did:plc:abc123").unwrap();
2679        assert_eq!(did, "did:plc:abc123");
2680        assert!(uri.is_none());
2681    }
2682
2683    #[test]
2684    fn route_subject_at_uri_extracts_repo_did_and_keeps_full_uri() {
2685        let (did, uri) = route_subject("at://did:plc:abc/app.bsky.feed.post/3xx").unwrap();
2686        assert_eq!(did, "did:plc:abc");
2687        assert_eq!(
2688            uri.as_deref(),
2689            Some("at://did:plc:abc/app.bsky.feed.post/3xx")
2690        );
2691    }
2692
2693    #[test]
2694    fn route_subject_rejects_at_uri_with_non_did_repo() {
2695        let err = route_subject("at://handle.example/app.bsky.feed.post/3xx").unwrap_err();
2696        assert!(matches!(err, Error::Signing(_)));
2697    }
2698
2699    #[test]
2700    fn route_subject_rejects_arbitrary_string() {
2701        assert!(route_subject("not a subject").is_err());
2702        assert!(route_subject("did:").is_err());
2703        assert!(route_subject("").is_err());
2704    }
2705
2706    // ---------- parse_iso8601_duration (#51) ----------
2707
2708    #[test]
2709    fn duration_p7d_is_seven_days() {
2710        assert_eq!(parse_iso8601_duration("P7D").unwrap(), 7 * 86_400);
2711    }
2712
2713    #[test]
2714    fn duration_p1w_is_seven_days() {
2715        assert_eq!(parse_iso8601_duration("P1W").unwrap(), 7 * 86_400);
2716    }
2717
2718    #[test]
2719    fn duration_pt12h_is_twelve_hours() {
2720        assert_eq!(parse_iso8601_duration("PT12H").unwrap(), 12 * 3600);
2721    }
2722
2723    #[test]
2724    fn duration_pt30m_is_thirty_minutes() {
2725        assert_eq!(parse_iso8601_duration("PT30M").unwrap(), 30 * 60);
2726    }
2727
2728    #[test]
2729    fn duration_pt45s_is_forty_five_seconds() {
2730        assert_eq!(parse_iso8601_duration("PT45S").unwrap(), 45);
2731    }
2732
2733    #[test]
2734    fn duration_p1d_t12h_combines() {
2735        assert_eq!(
2736            parse_iso8601_duration("P1DT12H").unwrap(),
2737            86_400 + 12 * 3600
2738        );
2739    }
2740
2741    #[test]
2742    fn duration_no_p_prefix_rejected() {
2743        assert!(parse_iso8601_duration("7D").is_err());
2744        assert!(parse_iso8601_duration("").is_err());
2745    }
2746
2747    #[test]
2748    fn duration_p_alone_rejected() {
2749        assert!(parse_iso8601_duration("P").is_err());
2750    }
2751
2752    #[test]
2753    fn duration_year_rejected_in_v14() {
2754        let err = parse_iso8601_duration("P1Y").unwrap_err();
2755        let msg = format!("{err}");
2756        assert!(msg.contains("not supported"));
2757    }
2758
2759    #[test]
2760    fn duration_unknown_unit_rejected() {
2761        assert!(parse_iso8601_duration("P5X").is_err());
2762        assert!(parse_iso8601_duration("PT5X").is_err());
2763    }
2764
2765    #[test]
2766    fn duration_zero_rejected() {
2767        assert!(parse_iso8601_duration("P0D").is_err());
2768    }
2769
2770    #[test]
2771    fn duration_trailing_digits_without_unit_rejected() {
2772        assert!(parse_iso8601_duration("P5").is_err());
2773        assert!(parse_iso8601_duration("PT12").is_err());
2774    }
2775
2776    #[test]
2777    fn duration_t_separator_with_no_time_rejected() {
2778        assert!(parse_iso8601_duration("P1DT").is_err());
2779    }
2780
2781    // ---------- audit reason builders (#51) ----------
2782
2783    #[test]
2784    fn record_action_audit_reason_shape() {
2785        let emitted = vec![
2786            serde_json::json!({"val": "!hide", "uri": "did:plc:subject0000000000000000"}),
2787            serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
2788        ];
2789        let json = build_record_action_audit_reason(
2790            42,
2791            ActionType::TempSuspension,
2792            "spam",
2793            &["spam".to_string(), "harassment".to_string()],
2794            4,
2795            2,
2796            true,
2797            &emitted,
2798        );
2799        let v: serde_json::Value = serde_json::from_str(&json).unwrap();
2800        assert_eq!(v["action_id"], 42);
2801        assert_eq!(v["action_type"], "temp_suspension");
2802        assert_eq!(v["primary_reason"], "spam");
2803        assert_eq!(v["reason_codes"], serde_json::json!(["spam", "harassment"]));
2804        assert_eq!(v["strike_value_base"], 4);
2805        assert_eq!(v["strike_value_applied"], 2);
2806        assert_eq!(v["was_dampened"], true);
2807        assert_eq!(v["emitted_labels"][0]["val"], "!hide");
2808        assert_eq!(v["emitted_labels"][1]["val"], "reason-spam");
2809    }
2810
2811    #[test]
2812    fn record_action_audit_reason_empty_emitted_labels() {
2813        let json = build_record_action_audit_reason(
2814            7,
2815            ActionType::Note,
2816            "spam",
2817            &["spam".to_string()],
2818            0,
2819            0,
2820            false,
2821            &[],
2822        );
2823        let v: serde_json::Value = serde_json::from_str(&json).unwrap();
2824        assert_eq!(v["emitted_labels"], serde_json::json!([]));
2825    }
2826
2827    #[test]
2828    fn revoke_action_audit_reason_with_reason() {
2829        let json = build_revoke_action_audit_reason(7, Some("appeal granted"), &[]);
2830        let v: serde_json::Value = serde_json::from_str(&json).unwrap();
2831        assert_eq!(v["action_id"], 7);
2832        assert_eq!(v["revoked_reason"], "appeal granted");
2833        assert_eq!(v["negated_labels"], serde_json::json!([]));
2834    }
2835
2836    #[test]
2837    fn revoke_action_audit_reason_without_reason_is_null() {
2838        let json = build_revoke_action_audit_reason(7, None, &[]);
2839        let v: serde_json::Value = serde_json::from_str(&json).unwrap();
2840        assert!(v["revoked_reason"].is_null());
2841    }
2842
2843    #[test]
2844    fn revoke_action_audit_reason_with_negated_labels() {
2845        let labels = vec![
2846            serde_json::json!({"val": "!takedown", "uri": "did:plc:subject0000000000000000"}),
2847            serde_json::json!({"val": "reason-spam", "uri": "did:plc:subject0000000000000000"}),
2848        ];
2849        let json = build_revoke_action_audit_reason(11, None, &labels);
2850        let v: serde_json::Value = serde_json::from_str(&json).unwrap();
2851        assert_eq!(v["negated_labels"][0]["val"], "!takedown");
2852        assert_eq!(v["negated_labels"][1]["val"], "reason-spam");
2853    }
2854
2855    // ---------- emission idempotency guards (#64) ----------
2856
2857    #[test]
2858    fn skip_action_label_when_row_already_carries_an_emitted_val() {
2859        // The defense-in-depth case the guard exists for: an
2860        // existing emitted_label_uri value means the action label
2861        // was already produced on the wire; the emission loop
2862        // must skip to avoid a duplicate (src, uri, val) record.
2863        assert!(should_skip_action_label_emission(Some("!takedown")));
2864    }
2865
2866    #[test]
2867    fn emit_action_label_when_row_has_no_emitted_val() {
2868        // v1.5's recordAction always reaches this branch in
2869        // production (fresh INSERT, NULL column). Pin the
2870        // happy-path so a future refactor that flips the
2871        // sense of the check breaks here.
2872        assert!(!should_skip_action_label_emission(None));
2873    }
2874
2875    #[test]
2876    fn skip_reason_emission_when_already_linked() {
2877        let existing: HashSet<&str> = ["spam", "harassment"].into_iter().collect();
2878        assert!(should_skip_reason_emission("spam", &existing));
2879        assert!(should_skip_reason_emission("harassment", &existing));
2880    }
2881
2882    #[test]
2883    fn emit_reason_when_not_already_linked() {
2884        let existing: HashSet<&str> = ["spam"].into_iter().collect();
2885        assert!(!should_skip_reason_emission("hate-speech", &existing));
2886        // Empty existing set (the v1.5 production case): never skip.
2887        let empty: HashSet<&str> = HashSet::new();
2888        assert!(!should_skip_reason_emission("anything", &empty));
2889    }
2890
2891    #[test]
2892    fn reason_emission_skip_check_is_case_sensitive() {
2893        // Reason codes are operator-vocabulary identifiers from
2894        // [moderation_reasons] (#47). Case sensitivity matches
2895        // SQLite's default text comparison and the recorder's
2896        // primary-reason resolution; pinning here defends
2897        // against an accidental case-fold in the guard during
2898        // refactoring.
2899        let existing: HashSet<&str> = ["SPAM"].into_iter().collect();
2900        assert!(!should_skip_reason_emission("spam", &existing));
2901    }
2902}