Skip to main content

cairn_mod/
writer.rs

1//! Single-writer task (§F5 / §10 architecture).
2//!
3//! All mutations to `labels`, `label_sequence`, `audit_log`, and
4//! `server_instance_lease` flow through one task that owns the write pool.
5//! Other components obtain a cheaply-cloneable [`WriterHandle`] and submit
6//! [`ApplyLabelRequest`] / [`NegateLabelRequest`] over an mpsc channel;
7//! replies travel on per-request oneshot channels so the caller may cancel
8//! (drop the future) without corrupting writer state — the transaction
9//! still commits, the reply is discarded.
10//!
11//! Invariants this module enforces, any of which is load-bearing:
12//!
13//! - **Single instance.** On startup [`spawn`] either acquires the lease
14//!   row in `server_instance_lease` (id = 1) or returns
15//!   [`Error::LeaseHeld`]. A background heartbeat updates `last_heartbeat`
16//!   every 10s; a second Cairn started against the same file sees a
17//!   <60s-old heartbeat and refuses. The 60s threshold is 6× the heartbeat
18//!   interval — tolerant of a single missed tick, strict enough that a
19//!   legitimately-orphaned lease (hard crash, unclean shutdown) ages out
20//!   before the operator manually restarts.
21//!
22//! - **Monotonic seq.** Sequence values come from
23//!   `label_sequence.seq` (AUTOINCREMENT) reserved inside the same
24//!   transaction as the `labels` INSERT. Under AUTOINCREMENT, SQLite
25//!   never reuses a seq value even across rollbacks.
26//!
27//! - **Monotonic cts.** For each `(src, uri, val)` tuple the writer
28//!   clamps to `max(wall_now, prev_cts + 1ms)` (§6.1) using the
29//!   `labels_tuple_idx` composite index for the prev-cts lookup.
30//!
31//! - **Audit-per-write.** Every successful label INSERT produces exactly
32//!   one `audit_log` row **in the same transaction**. The audit row's
33//!   `reason` column carries a JSON object documenting the event
34//!   ({"val", "neg", "moderator_reason"}). See [`AUDIT_REASON_SCHEMA`].
35//!
36//! - **Sign-then-insert.** `sign_label` runs between sequence reservation
37//!   and the labels INSERT. The signed bytes include `ver` but not `seq`
38//!   or `signing_key_id` (§6.2 step 1, §6.1: seq lives on the frame, not
39//!   the label).
40
41use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
42
43use proto_blue_crypto::{K256Keypair, Keypair as _, format_multikey};
44use sqlx::{Pool, Sqlite};
45use time::format_description::FormatItem;
46use time::macros::format_description;
47use time::{OffsetDateTime, PrimitiveDateTime};
48use tokio::sync::{broadcast, mpsc, oneshot, watch};
49use tokio::time::{MissedTickBehavior, interval};
50use uuid::Uuid;
51
52use crate::error::{Error, Result};
53use crate::label::Label;
54use crate::server::RetentionConfig;
55use crate::signing::sign_label;
56use crate::signing_key::SigningKey;
57
58/// Lease freshness threshold (§F5). A lease younger than this is held by
59/// a live peer; younger than 10s would be flaky under a single missed
60/// heartbeat, older than ~2× the value risks a genuine zombie blocking a
61/// legitimate restart for too long.
62pub(crate) const LEASE_STALE_MS: i64 = 60_000;
63
64/// Heartbeat interval. 10s × 6 = 60s staleness budget per the threshold.
65const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
66
67/// Granularity at which the writer checks "is it time to start the
68/// scheduled sweep?". Daily fires bucketed by UTC hour means we don't
69/// need finer than this; 60s keeps the check cost negligible (one
70/// `Instant::now()` comparison + an `Option` peek per minute).
71const SWEEP_CHECK_INTERVAL: Duration = Duration::from_secs(60);
72
73/// Upper bound on queued write commands. Moderators rarely drive more
74/// than single-digit req/s even on busy operators; 64 is a comfortable
75/// overshoot that still makes backpressure visible quickly under a bug.
76const COMMAND_BUFFER: usize = 64;
77
78/// Broadcast buffer for [`LabelEvent`] fan-out. Slow subscribers hit
79/// `RecvError::Lagged` past this; the subscribeLabels endpoint (#7) turns
80/// that into a client-connection close.
81const BROADCAST_BUFFER: usize = 1024;
82
83/// Audit-log `reason` JSON schema for `label_applied` / `label_negated`.
84/// Centralized as a doc constant so future actions (signing_key_added,
85/// report_resolved, etc.) have an obvious place to register their shape.
86///
87/// ```json
88/// {
89///   "val": "<label value>",
90///   "neg": true | false,
91///   "moderator_reason": "<free text>" | null
92/// }
93/// ```
94#[doc(alias = "audit_log.reason")]
95pub const AUDIT_REASON_SCHEMA: &str =
96    "label_applied / label_negated: { val, neg, moderator_reason }";
97
98/// Audit-log `reason` JSON schema for `report_resolved` (§F12 resolveReport
99/// atomicity — two rows land in one transaction: the inner label_applied
100/// per [`AUDIT_REASON_SCHEMA`] plus this one).
101///
102/// ```json
103/// {
104///   "applied_label_val": "<val>" | null,
105///   "resolution_reason": "<free text>" | null
106/// }
107/// ```
108#[doc(alias = "audit_log.reason.report_resolved")]
109pub const AUDIT_REASON_RESOLVE_REPORT: &str =
110    "report_resolved: { applied_label_val, resolution_reason }";
111
112/// Audit-log `reason` JSON schema for `reporter_flagged` /
113/// `reporter_unflagged` (§F12 flagReporter).
114///
115/// ```json
116/// {
117///   "did": "<flagged DID>",
118///   "suppressed": true | false,
119///   "moderator_reason": "<free text>" | null
120/// }
121/// ```
122#[doc(alias = "audit_log.reason.flag_reporter")]
123pub const AUDIT_REASON_FLAG_REPORTER: &str =
124    "reporter_flagged / reporter_unflagged: { did, suppressed, moderator_reason }";
125
126/// Audit-log `reason` JSON schema for `retention_sweep` (§F4 — written
127/// only by the operator-initiated admin path; the scheduled-fire
128/// path does NOT audit per Q6/D2). Captures the result of the sweep
129/// run for reconstruction-friendly ops queries.
130///
131/// ```json
132/// {
133///   "rows_deleted": <i64>,
134///   "batches": <u64>,
135///   "duration_ms": <u64>,
136///   "retention_days_applied": <u32> | null
137/// }
138/// ```
139#[doc(alias = "audit_log.reason.retention_sweep")]
140pub const AUDIT_REASON_RETENTION_SWEEP: &str =
141    "retention_sweep: { rows_deleted, batches, duration_ms, retention_days_applied }";
142
143/// Closed set of `audit_log.action` values emitted by Cairn write paths.
144///
145/// §F10: audit rows only for moderation decisions, not for input/operational
146/// events — `createReport` (input) intentionally does NOT audit. The
147/// listAuditLog handler validates its `action` query param against this
148/// exact set; unknown values return `InvalidRequest` rather than silently
149/// matching zero rows.
150///
151/// Adding a new moderation action requires: (1) the handler writes an
152/// `INSERT INTO audit_log (action, ...)` row, (2) the value is added here,
153/// (3) the `reason` schema for the new action is documented as a const
154/// alongside [`AUDIT_REASON_SCHEMA`] and friends.
155pub const AUDIT_ACTION_VALUES: &[&str] = &[
156    "label_applied",
157    "label_negated",
158    "report_resolved",
159    "reporter_flagged",
160    "reporter_unflagged",
161    "retention_sweep",
162];
163
164/// Closed set of `audit_log.outcome` values matching the SQL `CHECK`
165/// constraint in `migrations/0001_init.sql`. listAuditLog validates the
166/// `outcome` query param against this set; the SQL constraint itself is
167/// the durable source of truth — this slice mirrors it so the handler
168/// can reject invalid values pre-query without round-tripping to SQLite.
169pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
170
171/// RFC-3339 with millisecond precision. `Z` is appended by
172/// [`rfc3339_from_epoch_ms`] and stripped by [`parse_rfc3339_ms`] — kept
173/// out of the format description because `time::OffsetDateTime::parse`
174/// can't infer a UTC offset from the literal character `Z`, and
175/// symmetric Z-handling at the string boundary is simpler than switching
176/// parsing-only to `well_known::Rfc3339`.
177const CTS_FORMAT: &[FormatItem<'_>] =
178    format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
179
180/// Moderator request to apply a label (positive event).
181#[derive(Debug, Clone)]
182pub struct ApplyLabelRequest {
183    /// The DID of the moderator issuing the request. Becomes
184    /// `audit_log.actor_did`. `src` on the label itself is the writer's
185    /// service DID, not this.
186    pub actor_did: String,
187    /// AT-URI or DID of the subject being labeled.
188    pub uri: String,
189    /// Optional record-version pin. Present means "this specific version";
190    /// absent means "all versions / the account" per §6.1.
191    pub cid: Option<String>,
192    /// Label value, ≤128 bytes (validated on insert by schema CHECK and
193    /// by caller-side input validation on the XRPC boundary).
194    pub val: String,
195    /// Optional expiration timestamp (RFC-3339 Z). Stored only; expiry
196    /// enforcement is v1.1 (§F7).
197    pub exp: Option<String>,
198    /// Free-text reason recorded to `audit_log` only. Not signed, not
199    /// included in the label record.
200    pub moderator_reason: Option<String>,
201}
202
203/// Moderator request to negate (withdraw) a previously-applied label.
204///
205/// Uniqueness is at the `(src, uri, val)` tuple. The negation copies the
206/// most-recent applied event's `cid` so the negation pins the same record
207/// version — callers don't supply it.
208#[derive(Debug, Clone)]
209pub struct NegateLabelRequest {
210    /// Moderator DID issuing the negation. Becomes
211    /// `audit_log.actor_did`.
212    pub actor_did: String,
213    /// AT-URI or DID of the subject whose label is being withdrawn.
214    /// The `cid` pinning (if any) is copied from the prior apply.
215    pub uri: String,
216    /// Label value being negated. Must match an existing applied
217    /// label on `(src, uri, val)` or the call errors with
218    /// `LabelNotFound`.
219    pub val: String,
220    /// Free-text reason recorded to `audit_log` only.
221    pub moderator_reason: Option<String>,
222}
223
224/// A committed label event. Returned by `apply_label` / `negate_label` and
225/// broadcast to subscribers (subscribeLabels consumers in #7). Wire
226/// serialization is that consumer's concern — the LabelEvent itself
227/// carries no serde derive because the canonical encoding for the wire is
228/// the DAG-CBOR path in `crate::signing`, not serde JSON.
229#[derive(Debug, Clone)]
230pub struct LabelEvent {
231    /// Frame sequence number from `label_sequence`. Strictly monotonic.
232    pub seq: i64,
233    /// The full signed label (`sig` populated).
234    pub label: Label,
235}
236
237/// Internal write command. One variant per public method.
238enum WriteCommand {
239    Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
240    Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
241    ResolveReport(
242        ResolveReportRequest,
243        oneshot::Sender<Result<ResolvedReport>>,
244    ),
245    Sweep(SweepRequest, oneshot::Sender<Result<SweepBatchResult>>),
246    /// Append a hash-chained audit row (#39). Used by in-process
247    /// callers that don't have an existing transaction (e.g.,
248    /// `retentionSweep` after the sweep itself completes). Callers
249    /// that already hold a transaction (writer-internal handlers,
250    /// `flag_reporter`) call [`crate::audit::append::append_in_tx`]
251    /// directly within that transaction instead.
252    AppendAudit(
253        crate::audit::append::AuditRowForAppend,
254        oneshot::Sender<Result<i64>>,
255    ),
256    Shutdown(oneshot::Sender<Result<()>>),
257}
258
259/// Inline label-application sub-object for [`ResolveReportRequest`].
260/// Structurally aligned with [`ApplyLabelRequest`] minus the
261/// moderator-reason (that field on the outer request already captures
262/// the resolution rationale; the label's own audit reason is derived
263/// from the resolution flow).
264#[derive(Debug, Clone)]
265pub struct ApplyLabelInline {
266    /// Subject the label applies to (AT-URI or DID).
267    pub uri: String,
268    /// Optional record-version pin; `None` targets the account / all
269    /// versions of the record.
270    pub cid: Option<String>,
271    /// Label value to apply (≤128 bytes).
272    pub val: String,
273    /// Optional expiration timestamp (RFC-3339 Z). Stored only;
274    /// enforcement is v1.1.
275    pub exp: Option<String>,
276}
277
278/// Resolution outcome (#27): the implicit "with-label vs without-label"
279/// semantic, made explicit. The wire shape is `applyLabel: Option<…>`
280/// per the `resolveReport` lexicon; this enum is internal and the
281/// handler maps `None → Dismiss` / `Some(_) → ApplyLabel(_)` at the
282/// HTTP boundary.
283#[derive(Debug, Clone)]
284pub enum ResolutionAction {
285    /// Resolve without emitting a label (the operator-UX "dismiss"
286    /// flow). `resolution_label` on the row stays NULL.
287    Dismiss,
288    /// Resolve and emit a label in the same transaction. The label's
289    /// `val` is also recorded as `resolution_label` on the report row.
290    ApplyLabel(ApplyLabelInline),
291}
292
293impl ResolutionAction {
294    /// Borrow the inner `ApplyLabelInline` if this is an `ApplyLabel`
295    /// variant. Convenience for the writer's UPDATE / audit code that
296    /// needs both the optional label *and* its `val` projection.
297    pub fn as_apply(&self) -> Option<&ApplyLabelInline> {
298        match self {
299            ResolutionAction::Dismiss => None,
300            ResolutionAction::ApplyLabel(a) => Some(a),
301        }
302    }
303}
304
305/// Request to resolve a report (§F12 `resolveReport`). The optional
306/// label is applied **in the same transaction** as the report status
307/// update and both audit rows — §F5 single-writer invariant plus §F12
308/// atomicity requirement documented in the `resolveReport` lexicon.
309#[derive(Debug, Clone)]
310pub struct ResolveReportRequest {
311    /// Moderator DID issuing the resolution. Becomes
312    /// `audit_log.actor_did` on both the label-applied (if any) and
313    /// report_resolved rows.
314    pub actor_did: String,
315    /// Primary key of the report being resolved.
316    pub report_id: i64,
317    /// Whether the resolution emits a label or just closes the report.
318    /// Replaces the pre-#27 `apply_label: Option<ApplyLabelInline>`
319    /// representation; semantically identical, named for the operator
320    /// UX (dismiss vs apply-label).
321    pub action: ResolutionAction,
322    /// Free-text resolution rationale recorded to audit_log.
323    pub resolution_reason: Option<String>,
324}
325
326/// Result of a successful resolve. `label_event` is `Some(..)` iff
327/// the request carried `apply_label`; the broadcast has already
328/// happened inside the writer task post-commit.
329#[derive(Debug, Clone)]
330pub struct ResolvedReport {
331    /// The updated report row (status now `resolved`).
332    pub report: crate::report::Report,
333    /// The emitted label event, if the resolution included an
334    /// `apply_label`. Signed, broadcast, and committed as part of
335    /// the same transaction as the report UPDATE.
336    pub label_event: Option<LabelEvent>,
337}
338
339/// Trigger for the retention sweep (§F4). Carries no per-call
340/// parameters today — the cutoff comes from `retention_days` baked
341/// into the writer at spawn time, not from the request — but is
342/// kept as a typed unit so a future "sweep with explicit override
343/// for this one run" remains a non-breaking change to the variant.
344#[derive(Debug, Clone, Default)]
345pub struct SweepRequest;
346
347/// Per-batch outcome from one sweep dispatch through the writer's
348/// internal `WriteCommand::Sweep` channel. The writer task processes
349/// ONE batch per command so its main `select!` can interleave
350/// incoming label writes between batches (§F4 + §F5 — single-writer
351/// invariant + bounded latency). Callers who want a full sweep loop
352/// until [`Self::has_more`] is `false`; the [`WriterHandle::sweep`]
353/// convenience wrapper does this internally.
354#[derive(Debug, Clone)]
355pub struct SweepBatchResult {
356    /// Rows deleted in this batch. Zero means "no more old rows
357    /// match the cutoff" — caller stops looping.
358    pub rows_deleted: i64,
359    /// `true` when the batch hit the configured `sweep_batch_size`
360    /// limit and a follow-up batch may find more rows. `false`
361    /// indicates the batch was partial (last batch) and the sweep
362    /// is complete.
363    pub has_more: bool,
364    /// Cutoff days actually applied. `None` when the writer was
365    /// spawned with `retention_days = None` — the sweep is a no-op
366    /// in that configuration and `rows_deleted` is always 0.
367    pub retention_days_applied: Option<u32>,
368}
369
370/// Aggregate result of a full sweep run (returned by
371/// [`WriterHandle::sweep`] after looping over batches).
372#[derive(Debug, Clone)]
373pub struct SweepResult {
374    /// Total rows deleted across all batches.
375    pub rows_deleted: i64,
376    /// Number of batches issued.
377    pub batches: u64,
378    /// Wall-clock duration of the full sweep, in milliseconds.
379    pub duration_ms: u64,
380    /// Cutoff days actually applied, or `None` when the writer's
381    /// `retention_days` is `None` (sweep is a no-op).
382    pub retention_days_applied: Option<u32>,
383}
384
385/// Cheap handle to the writer task. Clones share the same underlying
386/// mpsc channel; a drop of the last clone is a silent shutdown signal
387/// to the writer task (receiver closes). For a clean shutdown that
388/// releases the lease row, call [`WriterHandle::shutdown`].
389#[derive(Debug, Clone)]
390pub struct WriterHandle {
391    tx: mpsc::Sender<WriteCommand>,
392    broadcast_tx: broadcast::Sender<LabelEvent>,
393    /// Flipped to `true` when the writer task starts its shutdown path.
394    /// Exposed via [`WriterHandle::shutdown_signal`] so peer components
395    /// (the subscribeLabels handler, future maintenance tasks) can close
396    /// their own resources cleanly. Using a watch channel rather than the
397    /// broadcast-channel close signal because clones of `WriterHandle`
398    /// keep `broadcast_tx` alive — receivers would otherwise never see
399    /// `RecvError::Closed`.
400    shutdown_rx: watch::Receiver<bool>,
401}
402
403impl WriterHandle {
404    /// Submit an apply-label request. Resolves when the writer has
405    /// committed the transaction and broadcast the event, or returns an
406    /// `Err` describing why the write was rejected.
407    pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
408        let (reply_tx, reply_rx) = oneshot::channel();
409        self.tx
410            .send(WriteCommand::Apply(req, reply_tx))
411            .await
412            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
413        reply_rx
414            .await
415            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
416    }
417
418    /// Submit a negate-label request. Returns [`Error::LabelNotFound`] if
419    /// no applied label currently exists for `(service_did, uri, val)`.
420    pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
421        let (reply_tx, reply_rx) = oneshot::channel();
422        self.tx
423            .send(WriteCommand::Negate(req, reply_tx))
424            .await
425            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
426        reply_rx
427            .await
428            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
429    }
430
431    /// Resolve a report, optionally applying a label in the same
432    /// atomic transaction. The label INSERT + audit row + report
433    /// UPDATE + resolution audit row all commit together or not at
434    /// all; the label broadcast fires post-commit inside the writer
435    /// task (§F5 + §F12 atomicity contract).
436    ///
437    /// Errors:
438    /// - [`Error::ReportNotFound`] — `report_id` doesn't exist.
439    /// - [`Error::ReportAlreadyResolved`] — report is not in the
440    ///   `pending` state. Handler maps to a generic `InvalidRequest`.
441    pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
442        let (reply_tx, reply_rx) = oneshot::channel();
443        self.tx
444            .send(WriteCommand::ResolveReport(req, reply_tx))
445            .await
446            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
447        reply_rx
448            .await
449            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
450    }
451
452    /// Trigger a §F4 retention sweep through the writer task. Loops
453    /// over single-batch dispatches until a batch returns
454    /// `has_more = false`, aggregating rows + duration.
455    /// Other writer commands interleave between batches (single-writer
456    /// invariant + bounded latency).
457    ///
458    /// Returns `rows_deleted = 0, batches = 0` when the writer was
459    /// spawned with `retention_days = None` (sweep configured off):
460    /// the first batch returns immediately and the loop exits with
461    /// the no-op result.
462    pub async fn sweep(&self, _req: SweepRequest) -> Result<SweepResult> {
463        let start = std::time::Instant::now();
464        let mut total_rows: i64 = 0;
465        let mut batches: u64 = 0;
466
467        loop {
468            let (reply_tx, reply_rx) = oneshot::channel();
469            self.tx
470                .send(WriteCommand::Sweep(SweepRequest, reply_tx))
471                .await
472                .map_err(|_| Error::Signing("writer task is shut down".into()))?;
473            let batch = reply_rx
474                .await
475                .map_err(|_| Error::Signing("writer dropped reply channel".into()))??;
476
477            total_rows += batch.rows_deleted;
478            // Count the no-op early-exit batch too — it represents
479            // the round-trip the caller paid for. Useful for tracing.
480            batches += 1;
481
482            if !batch.has_more {
483                return Ok(SweepResult {
484                    rows_deleted: total_rows,
485                    batches,
486                    duration_ms: start.elapsed().as_millis() as u64,
487                    retention_days_applied: batch.retention_days_applied,
488                });
489            }
490        }
491    }
492
493    /// Append a hash-chained audit row through the writer task (#39).
494    /// Used by in-process callers that don't already hold a transaction
495    /// — e.g., `retentionSweep`'s post-sweep audit row. Callers that
496    /// have an open transaction (writer-internal handlers,
497    /// `flag_reporter`) use [`crate::audit::append::append_in_tx`]
498    /// directly so the audit row commits atomically with the rest of
499    /// their work. Cross-process CLIs (publish/unpublish-service-
500    /// record) use [`crate::audit::append::append_via_pool`].
501    ///
502    /// Returns the inserted `audit_log.id`.
503    pub async fn append_audit(&self, row: crate::audit::append::AuditRowForAppend) -> Result<i64> {
504        let (reply_tx, reply_rx) = oneshot::channel();
505        self.tx
506            .send(WriteCommand::AppendAudit(row, reply_tx))
507            .await
508            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
509        reply_rx
510            .await
511            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
512    }
513
514    /// Subscribe to committed events. The returned receiver lags past
515    /// the internal broadcast buffer; the consumer (subscribeLabels, #7)
516    /// turns `RecvError::Lagged` into a connection close.
517    pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
518        self.broadcast_tx.subscribe()
519    }
520
521    /// Count of live broadcast receivers. Exposed for test sync points
522    /// (the WS handler subscribes asynchronously after upgrade; tests
523    /// that want to emit an event "after the subscriber is ready" poll
524    /// this until it reaches the expected count). Not a production hot
525    /// path — `broadcast::Sender::receiver_count` walks an atomic chain.
526    pub fn receiver_count(&self) -> usize {
527        self.broadcast_tx.receiver_count()
528    }
529
530    /// Observe writer lifecycle. The returned receiver starts at `false`
531    /// and flips to `true` once exactly, when the writer has accepted a
532    /// shutdown command and is about to release the lease. Holders close
533    /// downstream connections gracefully on the `true` transition.
534    pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
535        self.shutdown_rx.clone()
536    }
537
538    /// Explicit shutdown. Drains in-flight writes, releases the lease
539    /// row, stops the heartbeat, and returns. Idempotent-ish: calling on
540    /// an already-shut-down writer returns a channel-closed error, not a
541    /// panic.
542    pub async fn shutdown(&self) -> Result<()> {
543        let (reply_tx, reply_rx) = oneshot::channel();
544        self.tx
545            .send(WriteCommand::Shutdown(reply_tx))
546            .await
547            .map_err(|_| Error::Signing("writer task is already shut down".into()))?;
548        reply_rx
549            .await
550            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
551    }
552}
553
554/// Start the writer task. Acquires the instance lease, bootstraps the
555/// `signing_keys` row if empty, and spawns the event loop + heartbeat.
556///
557/// `service_did` stamps `labels.src` on every emitted event. `key` must
558/// be the private key whose public form is recorded (or will be recorded)
559/// in `signing_keys`.
560///
561/// `retention_days` is the §F4 retention cutoff in days. `None` disables
562/// the retention sweep entirely (`WriteCommand::Sweep` becomes a no-op
563/// returning `rows_deleted = 0`); `Some(N)` lets the sweep delete labels
564/// whose `created_at` is older than `now - N days`. Source of truth is
565/// [`crate::SubscribeConfig::retention_days`] — pass it through verbatim.
566///
567/// `retention` is the sweep execution policy (schedule + batching);
568/// distinct from `retention_days` per the [§F4 split](crate::RetentionConfig).
569pub async fn spawn(
570    pool: Pool<Sqlite>,
571    key: SigningKey,
572    service_did: String,
573    retention_days: Option<u32>,
574    retention: RetentionConfig,
575) -> Result<WriterHandle> {
576    let instance_id = acquire_lease(&pool).await?;
577    let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
578
579    let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
580    let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
581    let (shutdown_tx, shutdown_rx) = watch::channel(false);
582
583    let writer = Writer {
584        pool,
585        key,
586        service_did,
587        signing_key_id,
588        instance_id,
589        rx,
590        broadcast_tx: broadcast_tx.clone(),
591        shutdown_tx,
592        retention_days,
593        retention,
594    };
595
596    tokio::spawn(writer.run());
597
598    Ok(WriterHandle {
599        tx,
600        broadcast_tx,
601        shutdown_rx,
602    })
603}
604
605// ---------- Writer task ----------
606
607struct Writer {
608    pool: Pool<Sqlite>,
609    key: SigningKey,
610    service_did: String,
611    signing_key_id: i64,
612    instance_id: String,
613    rx: mpsc::Receiver<WriteCommand>,
614    broadcast_tx: broadcast::Sender<LabelEvent>,
615    shutdown_tx: watch::Sender<bool>,
616    /// Retention cutoff in days. `None` makes the sweep a no-op.
617    /// Source of truth is [`crate::SubscribeConfig::retention_days`].
618    retention_days: Option<u32>,
619    /// Sweep execution policy (schedule + batching).
620    retention: RetentionConfig,
621}
622
623/// Internal accumulator for an in-flight scheduled sweep. Lives only
624/// while [`Writer::run`] is mid-sweep; absent between sweeps.
625struct SweepRunState {
626    started_at: Instant,
627    rows: i64,
628    batches: u64,
629}
630
631impl Writer {
632    /// Compute the next absolute [`Instant`] at which the scheduled
633    /// sweep should fire, or `None` when the sweep is disabled.
634    ///
635    /// "Disabled" = `sweep_enabled = false` OR `retention_days = None`.
636    /// In the second case the sweep would be a no-op anyway; skipping
637    /// the timer entirely keeps the run loop quiet.
638    ///
639    /// Targets the next occurrence of `sweep_run_at_utc_hour:00:00`
640    /// UTC. If we're already past that hour today, target tomorrow at
641    /// the same hour.
642    fn compute_next_sweep_fire(&self) -> Option<Instant> {
643        if !self.retention.sweep_enabled || self.retention_days.is_none() {
644            return None;
645        }
646        let now_utc = OffsetDateTime::now_utc();
647        let target_hour = self.retention.sweep_run_at_utc_hour;
648        let target_today_time = time::Time::from_hms(target_hour, 0, 0)
649            .expect("hour validated < 24 by Config::validate");
650        let target_today = now_utc.replace_time(target_today_time);
651        let target_dt = if target_today <= now_utc {
652            target_today + time::Duration::days(1)
653        } else {
654            target_today
655        };
656        let wait = target_dt - now_utc;
657        let wait_secs = wait.whole_seconds().max(0) as u64;
658        Some(Instant::now() + Duration::from_secs(wait_secs))
659    }
660
661    async fn run(mut self) {
662        let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
663        heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
664        // First tick fires immediately; skip it — we just acquired the lease
665        // with a fresh timestamp, so touching it again is redundant.
666        heartbeat_timer.tick().await;
667
668        // Scheduled-sweep wiring (§F4). The check timer wakes the run
669        // loop once a minute to evaluate "is it time to start a sweep?";
670        // when a sweep starts, `sweep_state` is populated and the
671        // immediate-batch arm runs one batch per loop iteration until
672        // `has_more = false`. Inter-batch yielding to incoming commands
673        // is automatic — the biased select prefers `rx.recv()` and
674        // `heartbeat_timer.tick()` over running another batch.
675        let mut sweep_check_timer = interval(SWEEP_CHECK_INTERVAL);
676        sweep_check_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
677        sweep_check_timer.tick().await;
678
679        let mut next_scheduled_fire: Option<Instant> = self.compute_next_sweep_fire();
680        let mut sweep_state: Option<SweepRunState> = None;
681
682        loop {
683            let sweep_in_progress = sweep_state.is_some();
684            tokio::select! {
685                biased;
686                // Prefer draining write commands over heartbeats so a
687                // steady-state of inbound work does not starve itself
688                // behind a background housekeeping task.
689                cmd = self.rx.recv() => {
690                    match cmd {
691                        Some(WriteCommand::Apply(req, reply)) => {
692                            let res = self.handle_apply(req).await;
693                            // Caller may have cancelled; dropping the reply is fine.
694                            let _ = reply.send(res);
695                        }
696                        Some(WriteCommand::Negate(req, reply)) => {
697                            let res = self.handle_negate(req).await;
698                            let _ = reply.send(res);
699                        }
700                        Some(WriteCommand::ResolveReport(req, reply)) => {
701                            let res = self.handle_resolve_report(req).await;
702                            let _ = reply.send(res);
703                        }
704                        Some(WriteCommand::Sweep(req, reply)) => {
705                            // Single-batch dispatch — the caller loops
706                            // until has_more=false. Per-batch return
707                            // lets the writer's biased select pick up
708                            // pending Apply / Negate / ResolveReport
709                            // commands between batches (§F4 + §F5).
710                            let res = self.handle_sweep(req).await;
711                            let _ = reply.send(res);
712                        }
713                        Some(WriteCommand::AppendAudit(row, reply)) => {
714                            let res = self.handle_append_audit(row).await;
715                            let _ = reply.send(res);
716                        }
717                        Some(WriteCommand::Shutdown(reply)) => {
718                            // Flip the shutdown watch *before* releasing the
719                            // lease so subscriber tasks see the signal while
720                            // the DB is still accessible for their close-
721                            // frame sends.
722                            let _ = self.shutdown_tx.send(true);
723                            let res = self.release_lease().await;
724                            let _ = reply.send(res);
725                            return;
726                        }
727                        None => {
728                            // All handles dropped without explicit shutdown.
729                            // Best-effort lease release so a same-process
730                            // restart doesn't trip the 60s wait.
731                            let _ = self.shutdown_tx.send(true);
732                            if let Err(e) = self.release_lease().await {
733                                tracing::error!("lease release on handle drop: {e}");
734                            }
735                            return;
736                        }
737                    }
738                }
739                _ = heartbeat_timer.tick() => {
740                    if let Err(e) = self.heartbeat().await {
741                        // Transient DB errors are logged, not fatal. A
742                        // sustained failure is caught at the next-instance
743                        // startup check — not at this writer's expense.
744                        tracing::error!("lease heartbeat failed: {e}");
745                    }
746                }
747                _ = sweep_check_timer.tick() => {
748                    if sweep_state.is_none()
749                        && let Some(fire_at) = next_scheduled_fire
750                        && Instant::now() >= fire_at
751                    {
752                        sweep_state = Some(SweepRunState {
753                            started_at: Instant::now(),
754                            rows: 0,
755                            batches: 0,
756                        });
757                        next_scheduled_fire = Some(fire_at + Duration::from_secs(86_400));
758                        tracing::info!(
759                            retention_days = ?self.retention_days,
760                            sweep_batch_size = self.retention.sweep_batch_size,
761                            "scheduled retention sweep starting"
762                        );
763                    }
764                }
765                // Always-ready arm gated on sweep_in_progress. With biased
766                // order this ranks below rx.recv() + the two timers, so
767                // normal commands and heartbeats interleave naturally
768                // between batches (§F4 inter-batch yield, §F5 single-
769                // writer invariant).
770                _ = std::future::ready(()), if sweep_in_progress => {
771                    match self.handle_sweep(SweepRequest).await {
772                        Ok(batch) => {
773                            let state = sweep_state
774                                .as_mut()
775                                .expect("sweep_in_progress => sweep_state Some");
776                            state.rows += batch.rows_deleted;
777                            state.batches += 1;
778                            if !batch.has_more {
779                                let final_state = sweep_state.take().expect("just set");
780                                tracing::info!(
781                                    rows_deleted = final_state.rows,
782                                    batches = final_state.batches,
783                                    duration_ms = final_state.started_at.elapsed().as_millis() as u64,
784                                    retention_days_applied = ?batch.retention_days_applied,
785                                    "scheduled retention sweep complete"
786                                );
787                            }
788                        }
789                        Err(e) => {
790                            let final_state = sweep_state.take();
791                            tracing::error!(
792                                error = %e,
793                                batches_done = final_state.as_ref().map(|s| s.batches).unwrap_or(0),
794                                rows_so_far = final_state.as_ref().map(|s| s.rows).unwrap_or(0),
795                                "scheduled retention sweep batch failed; aborting run (will retry on next schedule)"
796                            );
797                        }
798                    }
799                }
800            }
801        }
802    }
803
804    /// Process one batch of the §F4 retention sweep. Single-batch
805    /// dispatch — see [`WriteCommand::Sweep`] for the loop ordering.
806    ///
807    /// Returns immediately with `rows_deleted = 0, has_more = false`
808    /// when `retention_days` is `None`. Otherwise issues one
809    /// `DELETE FROM labels WHERE created_at < cutoff LIMIT N` in its
810    /// own transaction; sets `has_more = true` iff the batch hit the
811    /// `sweep_batch_size` limit (suggesting more rows may match).
812    ///
813    /// Errors are propagated as `Err(_)` rather than swallowed: a
814    /// transient DB failure should surface to the operator on a
815    /// manual sweep, and to the schedule-loop logger on a scheduled
816    /// sweep. Idempotency (Q5) means the next sweep retries cleanly.
817    async fn handle_sweep(&self, _req: SweepRequest) -> Result<SweepBatchResult> {
818        let Some(days) = self.retention_days else {
819            return Ok(SweepBatchResult {
820                rows_deleted: 0,
821                has_more: false,
822                retention_days_applied: None,
823            });
824        };
825
826        let cutoff_ms = epoch_ms_now() - (days as i64) * 86_400_000;
827        let limit = self.retention.sweep_batch_size;
828
829        let mut tx = self.pool.begin().await?;
830        // SQLite's DELETE doesn't support LIMIT without the
831        // SQLITE_ENABLE_UPDATE_DELETE_LIMIT compile flag (off in
832        // bundled builds). The rowid IN (SELECT ... LIMIT) trick is
833        // the canonical portable workaround.
834        let result = sqlx::query!(
835            "DELETE FROM labels WHERE rowid IN (
836               SELECT rowid FROM labels WHERE created_at < ?1 LIMIT ?2
837             )",
838            cutoff_ms,
839            limit,
840        )
841        .execute(&mut *tx)
842        .await?;
843        tx.commit().await?;
844
845        let rows = result.rows_affected() as i64;
846        Ok(SweepBatchResult {
847            rows_deleted: rows,
848            has_more: rows >= limit,
849            retention_days_applied: Some(days),
850        })
851    }
852
853    /// Handler for [`WriteCommand::AppendAudit`]. Opens its own
854    /// transaction (BEGIN DEFERRED — fine because the writer task is
855    /// the only in-process audit appender during a `cairn serve`
856    /// session, and cross-process appenders use BEGIN IMMEDIATE on
857    /// their side), inserts the audit row with a freshly-computed
858    /// hash, commits.
859    async fn handle_append_audit(
860        &self,
861        row: crate::audit::append::AuditRowForAppend,
862    ) -> Result<i64> {
863        let mut tx = self.pool.begin().await?;
864        let id = crate::audit::append::append_in_tx(&mut tx, &row).await?;
865        tx.commit().await?;
866        Ok(id)
867    }
868
869    async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
870        let mut tx = self.pool.begin().await?;
871        let created_at = epoch_ms_now();
872        let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
873        tx.commit().await?;
874        // No-receivers is not a write failure (§plan point G).
875        let _ = self.broadcast_tx.send(event.clone());
876        Ok(event)
877    }
878
879    /// Inner label-application pipeline: reserve seq → clamp cts →
880    /// sign → INSERT label → INSERT audit. Does NOT commit the
881    /// transaction and does NOT broadcast — caller handles both.
882    ///
883    /// Extracted from `handle_apply` so `handle_resolve_report` can
884    /// reuse it inside the same transaction as the report update,
885    /// preserving §F5's single-writer-owns-label-emission invariant
886    /// (this helper only runs on the writer task; no other code path
887    /// can access seq allocation or cts clamping).
888    async fn apply_label_inner(
889        &self,
890        tx: &mut sqlx::Transaction<'_, Sqlite>,
891        req: &ApplyLabelRequest,
892        created_at_ms: i64,
893    ) -> Result<LabelEvent> {
894        let seq = reserve_seq(tx).await?;
895        let prev_cts: Option<String> = sqlx::query_scalar!(
896            // sqlx type override — MAX() strips column origin metadata so
897            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
898            // forces the nullable wrapper (Option<String>).
899            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
900               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
901            self.service_did,
902            req.uri,
903            req.val,
904        )
905        .fetch_one(&mut **tx)
906        .await?;
907
908        let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
909
910        let mut label = Label {
911            ver: 1,
912            src: self.service_did.clone(),
913            uri: req.uri.clone(),
914            cid: req.cid.clone(),
915            val: req.val.clone(),
916            neg: false,
917            cts: cts.clone(),
918            exp: req.exp.clone(),
919            sig: None,
920        };
921        label.sig = Some(sign_label(&self.key, &label)?);
922        let sig_bytes = label.sig.expect("just set").to_vec();
923
924        let neg_int: i64 = 0;
925        sqlx::query!(
926            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
927             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
928            seq,
929            label.ver,
930            label.src,
931            label.uri,
932            label.cid,
933            label.val,
934            neg_int,
935            label.cts,
936            label.exp,
937            sig_bytes,
938            self.signing_key_id,
939            created_at_ms,
940        )
941        .execute(&mut **tx)
942        .await?;
943
944        let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
945        crate::audit::append::append_in_tx(
946            tx,
947            &crate::audit::append::AuditRowForAppend {
948                created_at: created_at_ms,
949                action: "label_applied".into(),
950                actor_did: req.actor_did.clone(),
951                target: Some(req.uri.clone()),
952                target_cid: req.cid.clone(),
953                outcome: "success".into(),
954                reason: Some(audit_reason),
955            },
956        )
957        .await?;
958
959        Ok(LabelEvent { seq, label })
960    }
961
962    async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
963        let mut tx = self.pool.begin().await?;
964
965        // Most-recent event for the tuple. Not-found OR latest-is-already-
966        // a-negation both surface as LabelNotFound (§F6).
967        let latest = sqlx::query!(
968            "SELECT neg, cid FROM labels
969             WHERE src = ?1 AND uri = ?2 AND val = ?3
970             ORDER BY seq DESC LIMIT 1",
971            self.service_did,
972            req.uri,
973            req.val,
974        )
975        .fetch_optional(&mut *tx)
976        .await?;
977
978        let cid = match latest {
979            Some(row) if row.neg == 0 => row.cid,
980            _ => {
981                return Err(Error::LabelNotFound {
982                    src: self.service_did.clone(),
983                    uri: req.uri,
984                    val: req.val,
985                });
986            }
987        };
988
989        let seq = reserve_seq(&mut tx).await?;
990        // Prev cts lookup is the same query as apply; tuple uniqueness is
991        // on (src, uri, val) regardless of neg.
992        let prev_cts: Option<String> = sqlx::query_scalar!(
993            // sqlx type override — MAX() strips column origin metadata so
994            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
995            // forces the nullable wrapper (Option<String>).
996            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
997               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
998            self.service_did,
999            req.uri,
1000            req.val,
1001        )
1002        .fetch_one(&mut *tx)
1003        .await?;
1004
1005        let wall_now_ms = epoch_ms_now();
1006        let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
1007
1008        let mut label = Label {
1009            ver: 1,
1010            src: self.service_did.clone(),
1011            uri: req.uri.clone(),
1012            cid: cid.clone(),
1013            val: req.val.clone(),
1014            neg: true,
1015            cts: cts.clone(),
1016            exp: None, // Negations never carry an expiry.
1017            sig: None,
1018        };
1019        label.sig = Some(sign_label(&self.key, &label)?);
1020        let sig_bytes = label.sig.expect("just set").to_vec();
1021
1022        let created_at = wall_now_ms;
1023        let neg_int: i64 = 1;
1024        sqlx::query!(
1025            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
1026             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1027            seq,
1028            label.ver,
1029            label.src,
1030            label.uri,
1031            label.cid,
1032            label.val,
1033            neg_int,
1034            label.cts,
1035            label.exp,
1036            sig_bytes,
1037            self.signing_key_id,
1038            created_at,
1039        )
1040        .execute(&mut *tx)
1041        .await?;
1042
1043        let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
1044        crate::audit::append::append_in_tx(
1045            &mut tx,
1046            &crate::audit::append::AuditRowForAppend {
1047                created_at,
1048                action: "label_negated".into(),
1049                actor_did: req.actor_did.clone(),
1050                target: Some(req.uri.clone()),
1051                target_cid: cid.clone(),
1052                outcome: "success".into(),
1053                reason: Some(audit_reason),
1054            },
1055        )
1056        .await?;
1057
1058        tx.commit().await?;
1059
1060        let event = LabelEvent { seq, label };
1061        let _ = self.broadcast_tx.send(event.clone());
1062        Ok(event)
1063    }
1064
1065    /// Atomic resolveReport flow — §F12 requires label INSERT +
1066    /// audit(label_applied) + report UPDATE + audit(report_resolved)
1067    /// all commit together or not at all. Everything happens inside
1068    /// one SQLite transaction on the writer task, preserving §F5's
1069    /// single-writer invariant for the optional label emission.
1070    ///
1071    /// Returns [`Error::ReportNotFound`] if the report doesn't
1072    /// exist, [`Error::ReportAlreadyResolved`] if it is not in
1073    /// `status='pending'`. Label-value validation is a handler
1074    /// concern (the handler consults `AdminConfig.label_values`
1075    /// before sending the command here — saves a writer round-trip
1076    /// on invalid input and keeps the anti-leak message local).
1077    async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
1078        use crate::report::ReportStatus;
1079
1080        let mut tx = self.pool.begin().await?;
1081
1082        // 1. Load the report; confirm pending status.
1083        let current = sqlx::query_as!(
1084            crate::report::Report,
1085            r#"SELECT
1086                 id                 AS "id!: i64",
1087                 created_at         AS "created_at!: String",
1088                 reported_by        AS "reported_by!: String",
1089                 reason_type        AS "reason_type!: String",
1090                 reason,
1091                 subject_type       AS "subject_type!: String",
1092                 subject_did        AS "subject_did!: String",
1093                 subject_uri,
1094                 subject_cid,
1095                 status             AS "status!: ReportStatus",
1096                 resolved_at,
1097                 resolved_by,
1098                 resolution_label,
1099                 resolution_reason
1100               FROM reports WHERE id = ?1"#,
1101            req.report_id,
1102        )
1103        .fetch_optional(&mut *tx)
1104        .await?;
1105
1106        let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
1107        if report.status != ReportStatus::Pending {
1108            return Err(Error::ReportAlreadyResolved { id: req.report_id });
1109        }
1110
1111        let created_at = epoch_ms_now();
1112
1113        // 2. Optional label-apply BEFORE the report UPDATE so audit
1114        // row insertion order reflects logical sequence
1115        // (label_applied, then report_resolved — §G per #15 criteria).
1116        let label_event = if let Some(apply) = req.action.as_apply() {
1117            let apply_req = ApplyLabelRequest {
1118                actor_did: req.actor_did.clone(),
1119                uri: apply.uri.clone(),
1120                cid: apply.cid.clone(),
1121                val: apply.val.clone(),
1122                exp: apply.exp.clone(),
1123                // The label's own audit_reason JSON captures {val,
1124                // neg, moderator_reason}; the resolution-level reason
1125                // goes on the report_resolved audit row below.
1126                moderator_reason: None,
1127            };
1128            Some(
1129                self.apply_label_inner(&mut tx, &apply_req, created_at)
1130                    .await?,
1131            )
1132        } else {
1133            None
1134        };
1135
1136        // 3. UPDATE reports.
1137        let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
1138        let resolution_label = req.action.as_apply().map(|a| a.val.clone());
1139        let resolved_status = ReportStatus::Resolved;
1140        sqlx::query!(
1141            "UPDATE reports SET
1142                status = ?1,
1143                resolved_at = ?2,
1144                resolved_by = ?3,
1145                resolution_label = ?4,
1146                resolution_reason = ?5
1147             WHERE id = ?6",
1148            resolved_status,
1149            resolved_at_rfc,
1150            req.actor_did,
1151            resolution_label,
1152            req.resolution_reason,
1153            req.report_id,
1154        )
1155        .execute(&mut *tx)
1156        .await?;
1157
1158        // 4. Audit: report_resolved.
1159        let audit_reason = build_resolve_audit_reason(
1160            req.action.as_apply().map(|a| a.val.as_str()),
1161            req.resolution_reason.as_deref(),
1162        );
1163        crate::audit::append::append_in_tx(
1164            &mut tx,
1165            &crate::audit::append::AuditRowForAppend {
1166                created_at,
1167                action: "report_resolved".into(),
1168                actor_did: req.actor_did.clone(),
1169                target: Some(req.report_id.to_string()),
1170                target_cid: None,
1171                outcome: "success".into(),
1172                reason: Some(audit_reason),
1173            },
1174        )
1175        .await?;
1176
1177        tx.commit().await?;
1178
1179        // 5. Broadcast (post-commit, per §F5 broadcast-after-commit rule).
1180        if let Some(event) = &label_event {
1181            let _ = self.broadcast_tx.send(event.clone());
1182        }
1183
1184        // 6. Mutate the loaded report struct to reflect the committed
1185        // state and return it. Avoids a re-fetch round-trip since we
1186        // know exactly what changed.
1187        report.status = ReportStatus::Resolved;
1188        report.resolved_at = Some(resolved_at_rfc);
1189        report.resolved_by = Some(req.actor_did);
1190        report.resolution_label = resolution_label;
1191        report.resolution_reason = req.resolution_reason;
1192
1193        Ok(ResolvedReport {
1194            report,
1195            label_event,
1196        })
1197    }
1198
1199    async fn heartbeat(&self) -> Result<()> {
1200        let now_ms = epoch_ms_now();
1201        sqlx::query!(
1202            "UPDATE server_instance_lease SET last_heartbeat = ?1
1203             WHERE id = 1 AND instance_id = ?2",
1204            now_ms,
1205            self.instance_id,
1206        )
1207        .execute(&self.pool)
1208        .await?;
1209        Ok(())
1210    }
1211
1212    async fn release_lease(&self) -> Result<()> {
1213        sqlx::query!(
1214            "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
1215            self.instance_id,
1216        )
1217        .execute(&self.pool)
1218        .await?;
1219        Ok(())
1220    }
1221}
1222
1223// ---------- Lease + key bootstrap ----------
1224
1225/// Acquire the §F5 single-instance lease.
1226///
1227/// Returns the new instance id on success. Returns
1228/// [`Error::LeaseHeld`] when an existing row's `last_heartbeat` is
1229/// still within [`LEASE_STALE_MS`] — i.e., another writer is live.
1230///
1231/// Exposed as `pub(crate)` so one-shot CLI tools that mutate
1232/// `audit_log` (`cairn audit-rebuild`, #40) can assert the same
1233/// "no live writer" invariant the long-running serve process holds.
1234/// The CLI tool releases via [`release_lease_by_id`] when done.
1235pub(crate) async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
1236    let now_ms = epoch_ms_now();
1237
1238    let existing =
1239        sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
1240            .fetch_optional(pool)
1241            .await?;
1242
1243    if let Some(row) = existing {
1244        let age_ms = (now_ms - row.last_heartbeat).max(0);
1245        if age_ms < LEASE_STALE_MS {
1246            return Err(Error::LeaseHeld {
1247                instance_id: row.instance_id,
1248                age_secs: (age_ms / 1000) as u64,
1249            });
1250        }
1251    }
1252
1253    let new_id = Uuid::new_v4().to_string();
1254    // INSERT OR REPLACE on id=1: takes over a stale lease row if present,
1255    // or creates the row on first startup. Legal because the staleness
1256    // check above confirms no live peer.
1257    sqlx::query!(
1258        "INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
1259         VALUES (1, ?1, ?2, ?2)
1260         ON CONFLICT(id) DO UPDATE SET
1261             instance_id = excluded.instance_id,
1262             acquired_at = excluded.acquired_at,
1263             last_heartbeat = excluded.last_heartbeat",
1264        new_id,
1265        now_ms,
1266    )
1267    .execute(pool)
1268    .await?;
1269
1270    Ok(new_id)
1271}
1272
1273/// Release a lease acquired via [`acquire_lease`] by its instance id.
1274/// Free-function counterpart to [`Writer::release_lease`] for callers
1275/// that don't own a `Writer` (one-shot CLI tools that just hold the
1276/// lease for the duration of a data migration).
1277pub(crate) async fn release_lease_by_id(pool: &Pool<Sqlite>, instance_id: &str) -> Result<()> {
1278    sqlx::query!(
1279        "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
1280        instance_id,
1281    )
1282    .execute(pool)
1283    .await?;
1284    Ok(())
1285}
1286
1287async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
1288    let kp = K256Keypair::from_private_key(key.expose_secret())?;
1289    let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
1290
1291    let existing =
1292        sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
1293            .fetch_optional(pool)
1294            .await?;
1295
1296    if let Some(row) = existing {
1297        if row.public_key_multibase != my_multibase {
1298            return Err(Error::Signing(format!(
1299                "signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
1300                row.public_key_multibase, my_multibase
1301            )));
1302        }
1303        return Ok(row.id);
1304    }
1305
1306    let now_ms = epoch_ms_now();
1307    let valid_from = rfc3339_from_epoch_ms(now_ms)?;
1308    let id = sqlx::query_scalar!(
1309        "INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
1310         VALUES (?1, ?2, NULL, ?3)
1311         RETURNING id",
1312        my_multibase,
1313        valid_from,
1314        now_ms,
1315    )
1316    .fetch_one(pool)
1317    .await?;
1318
1319    Ok(id)
1320}
1321
1322// ---------- Pure helpers (unit-testable without a DB) ----------
1323
1324async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
1325    // label_sequence has a single auto-increment column; `DEFAULT VALUES`
1326    // triggers a fresh allocation. Under AUTOINCREMENT the assigned seq
1327    // is never reused even across rolled-back transactions — this is what
1328    // gives us "strictly monotonic, no gaps from writer's perspective."
1329    let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
1330        .fetch_one(tx)
1331        .await?;
1332    Ok(seq)
1333}
1334
1335/// §6.1 monotonicity clamp. `prev_cts_str`, when present, is expected to
1336/// be in the writer's canonical `CTS_FORMAT`.
1337fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
1338    let effective_ms = match prev_cts_str {
1339        Some(s) => {
1340            let prev_ms = parse_rfc3339_ms(s)?;
1341            wall_now_ms.max(prev_ms + 1)
1342        }
1343        None => wall_now_ms,
1344    };
1345    rfc3339_from_epoch_ms(effective_ms)
1346}
1347
1348/// Epoch-ms → RFC-3339 Z with millisecond precision. `pub(crate)` so
1349/// peer modules (admin/audit_view, future retention sweep) share the
1350/// single formatter used on the writer's timestamp boundary.
1351pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
1352    let nanos: i128 = (ms as i128) * 1_000_000;
1353    let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
1354        .map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
1355    let formatted = dt
1356        .format(&CTS_FORMAT)
1357        .map_err(|e| Error::Signing(format!("format cts: {e}")))?;
1358    Ok(format!("{formatted}Z"))
1359}
1360
1361/// RFC-3339 Z with millisecond precision → epoch-ms. `pub(crate)` so
1362/// admin handlers can validate `since` / `until` query params against
1363/// the same parser the writer uses on its input boundary.
1364pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
1365    let stripped = s
1366        .strip_suffix('Z')
1367        .ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
1368    let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
1369        .map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
1370    let nanos = pdt.assume_utc().unix_timestamp_nanos();
1371    Ok((nanos / 1_000_000) as i64)
1372}
1373
1374/// Current wall-clock time as Unix epoch milliseconds. `pub(crate)` so
1375/// peer modules (server, future retention sweep) share one implementation
1376/// rather than each inlining `SystemTime::now()` with its own error
1377/// handling.
1378pub(crate) fn epoch_ms_now() -> i64 {
1379    SystemTime::now()
1380        .duration_since(UNIX_EPOCH)
1381        .expect("system clock before unix epoch")
1382        .as_millis() as i64
1383}
1384
1385fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
1386    let body = serde_json::json!({
1387        "val": val,
1388        "neg": neg,
1389        "moderator_reason": moderator_reason,
1390    });
1391    body.to_string()
1392}
1393
1394/// `report_resolved` audit reason JSON — see [`AUDIT_REASON_RESOLVE_REPORT`]
1395/// for the schema.
1396fn build_resolve_audit_reason(
1397    applied_label_val: Option<&str>,
1398    resolution_reason: Option<&str>,
1399) -> String {
1400    serde_json::json!({
1401        "applied_label_val": applied_label_val,
1402        "resolution_reason": resolution_reason,
1403    })
1404    .to_string()
1405}
1406
1407/// `retention_sweep` audit reason JSON — see
1408/// [`AUDIT_REASON_RETENTION_SWEEP`] for the schema. Shared with the
1409/// retentionSweep admin handler (audited only on the operator-
1410/// initiated path per Q6/D2).
1411pub(crate) fn build_retention_sweep_audit_reason(result: &SweepResult) -> String {
1412    serde_json::json!({
1413        "rows_deleted": result.rows_deleted,
1414        "batches": result.batches,
1415        "duration_ms": result.duration_ms,
1416        "retention_days_applied": result.retention_days_applied,
1417    })
1418    .to_string()
1419}
1420
1421/// `reporter_flagged` / `reporter_unflagged` audit reason JSON —
1422/// see [`AUDIT_REASON_FLAG_REPORTER`]. Shared with the flagReporter
1423/// handler (not the writer — flagReporter is a direct handler txn
1424/// per the #15 criteria).
1425pub(crate) fn build_flag_reporter_audit_reason(
1426    did: &str,
1427    suppressed: bool,
1428    moderator_reason: Option<&str>,
1429) -> String {
1430    serde_json::json!({
1431        "did": did,
1432        "suppressed": suppressed,
1433        "moderator_reason": moderator_reason,
1434    })
1435    .to_string()
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440    use super::*;
1441
1442    // 1_715_000_000_000 ms since epoch = 2024-05-06T12:53:20.000Z UTC.
1443    // The four tests below pin every branch of the clamp: no prior,
1444    // wall-ahead, wall-behind (use prev+1ms), wall-equal (advance anyway).
1445
1446    #[test]
1447    fn clamp_cts_uses_wall_clock_when_no_prior() {
1448        let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
1449        assert_eq!(out, "2024-05-06T12:53:20.000Z");
1450    }
1451
1452    #[test]
1453    fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
1454        // wall_clock one second behind prev: result is prev + 1ms, not wall.
1455        let prev = "2024-05-06T12:53:20.500Z";
1456        let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
1457        assert_eq!(out, "2024-05-06T12:53:20.501Z");
1458    }
1459
1460    #[test]
1461    fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
1462        let prev = "2024-05-06T12:53:20.500Z";
1463        let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
1464        assert_eq!(out, "2024-05-06T12:53:22.000Z");
1465    }
1466
1467    #[test]
1468    fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
1469        // wall_clock == prev exactly: must advance by 1ms (strictly greater).
1470        let prev = "2024-05-06T12:53:20.500Z";
1471        let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
1472        assert_eq!(out, "2024-05-06T12:53:20.501Z");
1473    }
1474
1475    #[test]
1476    fn build_audit_reason_shape_matches_documented_schema() {
1477        let json = build_audit_reason("spam", false, Some("user reported"));
1478        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
1479        assert_eq!(v["val"], "spam");
1480        assert_eq!(v["neg"], false);
1481        assert_eq!(v["moderator_reason"], "user reported");
1482    }
1483
1484    #[test]
1485    fn build_audit_reason_null_when_moderator_reason_absent() {
1486        let json = build_audit_reason("spam", true, None);
1487        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
1488        assert!(v["moderator_reason"].is_null());
1489    }
1490
1491    #[test]
1492    fn rfc3339_roundtrip() {
1493        let s = "2026-04-22T12:00:00.938Z";
1494        let ms = parse_rfc3339_ms(s).expect("parse");
1495        let back = rfc3339_from_epoch_ms(ms).expect("format");
1496        assert_eq!(back, s);
1497    }
1498}