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}