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, 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::signing::sign_label;
55use crate::signing_key::SigningKey;
56
57/// Lease freshness threshold (§F5). A lease younger than this is held by
58/// a live peer; younger than 10s would be flaky under a single missed
59/// heartbeat, older than ~2× the value risks a genuine zombie blocking a
60/// legitimate restart for too long.
61const LEASE_STALE_MS: i64 = 60_000;
62
63/// Heartbeat interval. 10s × 6 = 60s staleness budget per the threshold.
64const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
65
66/// Upper bound on queued write commands. Moderators rarely drive more
67/// than single-digit req/s even on busy operators; 64 is a comfortable
68/// overshoot that still makes backpressure visible quickly under a bug.
69const COMMAND_BUFFER: usize = 64;
70
71/// Broadcast buffer for [`LabelEvent`] fan-out. Slow subscribers hit
72/// `RecvError::Lagged` past this; the subscribeLabels endpoint (#7) turns
73/// that into a client-connection close.
74const BROADCAST_BUFFER: usize = 1024;
75
76/// Audit-log `reason` JSON schema for `label_applied` / `label_negated`.
77/// Centralized as a doc constant so future actions (signing_key_added,
78/// report_resolved, etc.) have an obvious place to register their shape.
79///
80/// ```json
81/// {
82///   "val": "<label value>",
83///   "neg": true | false,
84///   "moderator_reason": "<free text>" | null
85/// }
86/// ```
87#[doc(alias = "audit_log.reason")]
88pub const AUDIT_REASON_SCHEMA: &str =
89    "label_applied / label_negated: { val, neg, moderator_reason }";
90
91/// Audit-log `reason` JSON schema for `report_resolved` (§F12 resolveReport
92/// atomicity — two rows land in one transaction: the inner label_applied
93/// per [`AUDIT_REASON_SCHEMA`] plus this one).
94///
95/// ```json
96/// {
97///   "applied_label_val": "<val>" | null,
98///   "resolution_reason": "<free text>" | null
99/// }
100/// ```
101#[doc(alias = "audit_log.reason.report_resolved")]
102pub const AUDIT_REASON_RESOLVE_REPORT: &str =
103    "report_resolved: { applied_label_val, resolution_reason }";
104
105/// Audit-log `reason` JSON schema for `reporter_flagged` /
106/// `reporter_unflagged` (§F12 flagReporter).
107///
108/// ```json
109/// {
110///   "did": "<flagged DID>",
111///   "suppressed": true | false,
112///   "moderator_reason": "<free text>" | null
113/// }
114/// ```
115#[doc(alias = "audit_log.reason.flag_reporter")]
116pub const AUDIT_REASON_FLAG_REPORTER: &str =
117    "reporter_flagged / reporter_unflagged: { did, suppressed, moderator_reason }";
118
119/// Closed set of `audit_log.action` values emitted by Cairn write paths.
120///
121/// §F10: audit rows only for moderation decisions, not for input/operational
122/// events — `createReport` (input) intentionally does NOT audit. The
123/// listAuditLog handler validates its `action` query param against this
124/// exact set; unknown values return `InvalidRequest` rather than silently
125/// matching zero rows.
126///
127/// Adding a new moderation action requires: (1) the handler writes an
128/// `INSERT INTO audit_log (action, ...)` row, (2) the value is added here,
129/// (3) the `reason` schema for the new action is documented as a const
130/// alongside [`AUDIT_REASON_SCHEMA`] and friends.
131pub const AUDIT_ACTION_VALUES: &[&str] = &[
132    "label_applied",
133    "label_negated",
134    "report_resolved",
135    "reporter_flagged",
136    "reporter_unflagged",
137];
138
139/// Closed set of `audit_log.outcome` values matching the SQL `CHECK`
140/// constraint in `migrations/0001_init.sql`. listAuditLog validates the
141/// `outcome` query param against this set; the SQL constraint itself is
142/// the durable source of truth — this slice mirrors it so the handler
143/// can reject invalid values pre-query without round-tripping to SQLite.
144pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
145
146/// RFC-3339 with millisecond precision. `Z` is appended by
147/// [`rfc3339_from_epoch_ms`] and stripped by [`parse_rfc3339_ms`] — kept
148/// out of the format description because `time::OffsetDateTime::parse`
149/// can't infer a UTC offset from the literal character `Z`, and
150/// symmetric Z-handling at the string boundary is simpler than switching
151/// parsing-only to `well_known::Rfc3339`.
152const CTS_FORMAT: &[FormatItem<'_>] =
153    format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
154
155/// Moderator request to apply a label (positive event).
156#[derive(Debug, Clone)]
157pub struct ApplyLabelRequest {
158    /// The DID of the moderator issuing the request. Becomes
159    /// `audit_log.actor_did`. `src` on the label itself is the writer's
160    /// service DID, not this.
161    pub actor_did: String,
162    /// AT-URI or DID of the subject being labeled.
163    pub uri: String,
164    /// Optional record-version pin. Present means "this specific version";
165    /// absent means "all versions / the account" per §6.1.
166    pub cid: Option<String>,
167    /// Label value, ≤128 bytes (validated on insert by schema CHECK and
168    /// by caller-side input validation on the XRPC boundary).
169    pub val: String,
170    /// Optional expiration timestamp (RFC-3339 Z). Stored only; expiry
171    /// enforcement is v1.1 (§F7).
172    pub exp: Option<String>,
173    /// Free-text reason recorded to `audit_log` only. Not signed, not
174    /// included in the label record.
175    pub moderator_reason: Option<String>,
176}
177
178/// Moderator request to negate (withdraw) a previously-applied label.
179///
180/// Uniqueness is at the `(src, uri, val)` tuple. The negation copies the
181/// most-recent applied event's `cid` so the negation pins the same record
182/// version — callers don't supply it.
183#[derive(Debug, Clone)]
184pub struct NegateLabelRequest {
185    /// Moderator DID issuing the negation. Becomes
186    /// `audit_log.actor_did`.
187    pub actor_did: String,
188    /// AT-URI or DID of the subject whose label is being withdrawn.
189    /// The `cid` pinning (if any) is copied from the prior apply.
190    pub uri: String,
191    /// Label value being negated. Must match an existing applied
192    /// label on `(src, uri, val)` or the call errors with
193    /// `LabelNotFound`.
194    pub val: String,
195    /// Free-text reason recorded to `audit_log` only.
196    pub moderator_reason: Option<String>,
197}
198
199/// A committed label event. Returned by `apply_label` / `negate_label` and
200/// broadcast to subscribers (subscribeLabels consumers in #7). Wire
201/// serialization is that consumer's concern — the LabelEvent itself
202/// carries no serde derive because the canonical encoding for the wire is
203/// the DAG-CBOR path in `crate::signing`, not serde JSON.
204#[derive(Debug, Clone)]
205pub struct LabelEvent {
206    /// Frame sequence number from `label_sequence`. Strictly monotonic.
207    pub seq: i64,
208    /// The full signed label (`sig` populated).
209    pub label: Label,
210}
211
212/// Internal write command. One variant per public method.
213enum WriteCommand {
214    Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
215    Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
216    ResolveReport(
217        ResolveReportRequest,
218        oneshot::Sender<Result<ResolvedReport>>,
219    ),
220    Shutdown(oneshot::Sender<Result<()>>),
221}
222
223/// Inline label-application sub-object for [`ResolveReportRequest`].
224/// Structurally aligned with [`ApplyLabelRequest`] minus the
225/// moderator-reason (that field on the outer request already captures
226/// the resolution rationale; the label's own audit reason is derived
227/// from the resolution flow).
228#[derive(Debug, Clone)]
229pub struct ApplyLabelInline {
230    /// Subject the label applies to (AT-URI or DID).
231    pub uri: String,
232    /// Optional record-version pin; `None` targets the account / all
233    /// versions of the record.
234    pub cid: Option<String>,
235    /// Label value to apply (≤128 bytes).
236    pub val: String,
237    /// Optional expiration timestamp (RFC-3339 Z). Stored only;
238    /// enforcement is v1.1.
239    pub exp: Option<String>,
240}
241
242/// Request to resolve a report (§F12 `resolveReport`). The optional
243/// `apply_label` is applied **in the same transaction** as the report
244/// status update and both audit rows — §F5 single-writer invariant
245/// plus §F12 atomicity requirement documented in the `resolveReport`
246/// lexicon.
247#[derive(Debug, Clone)]
248pub struct ResolveReportRequest {
249    /// Moderator DID issuing the resolution. Becomes
250    /// `audit_log.actor_did` on both the label-applied (if any) and
251    /// report_resolved rows.
252    pub actor_did: String,
253    /// Primary key of the report being resolved.
254    pub report_id: i64,
255    /// Optional label to emit atomically with the resolution. When
256    /// `None`, the resolution only updates report state + writes
257    /// the audit row.
258    pub apply_label: Option<ApplyLabelInline>,
259    /// Free-text resolution rationale recorded to audit_log.
260    pub resolution_reason: Option<String>,
261}
262
263/// Result of a successful resolve. `label_event` is `Some(..)` iff
264/// the request carried `apply_label`; the broadcast has already
265/// happened inside the writer task post-commit.
266#[derive(Debug, Clone)]
267pub struct ResolvedReport {
268    /// The updated report row (status now `resolved`).
269    pub report: crate::report::Report,
270    /// The emitted label event, if the resolution included an
271    /// `apply_label`. Signed, broadcast, and committed as part of
272    /// the same transaction as the report UPDATE.
273    pub label_event: Option<LabelEvent>,
274}
275
276/// Cheap handle to the writer task. Clones share the same underlying
277/// mpsc channel; a drop of the last clone is a silent shutdown signal
278/// to the writer task (receiver closes). For a clean shutdown that
279/// releases the lease row, call [`WriterHandle::shutdown`].
280#[derive(Debug, Clone)]
281pub struct WriterHandle {
282    tx: mpsc::Sender<WriteCommand>,
283    broadcast_tx: broadcast::Sender<LabelEvent>,
284    /// Flipped to `true` when the writer task starts its shutdown path.
285    /// Exposed via [`WriterHandle::shutdown_signal`] so peer components
286    /// (the subscribeLabels handler, future maintenance tasks) can close
287    /// their own resources cleanly. Using a watch channel rather than the
288    /// broadcast-channel close signal because clones of `WriterHandle`
289    /// keep `broadcast_tx` alive — receivers would otherwise never see
290    /// `RecvError::Closed`.
291    shutdown_rx: watch::Receiver<bool>,
292}
293
294impl WriterHandle {
295    /// Submit an apply-label request. Resolves when the writer has
296    /// committed the transaction and broadcast the event, or returns an
297    /// `Err` describing why the write was rejected.
298    pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
299        let (reply_tx, reply_rx) = oneshot::channel();
300        self.tx
301            .send(WriteCommand::Apply(req, reply_tx))
302            .await
303            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
304        reply_rx
305            .await
306            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
307    }
308
309    /// Submit a negate-label request. Returns [`Error::LabelNotFound`] if
310    /// no applied label currently exists for `(service_did, uri, val)`.
311    pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
312        let (reply_tx, reply_rx) = oneshot::channel();
313        self.tx
314            .send(WriteCommand::Negate(req, reply_tx))
315            .await
316            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
317        reply_rx
318            .await
319            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
320    }
321
322    /// Resolve a report, optionally applying a label in the same
323    /// atomic transaction. The label INSERT + audit row + report
324    /// UPDATE + resolution audit row all commit together or not at
325    /// all; the label broadcast fires post-commit inside the writer
326    /// task (§F5 + §F12 atomicity contract).
327    ///
328    /// Errors:
329    /// - [`Error::ReportNotFound`] — `report_id` doesn't exist.
330    /// - [`Error::ReportAlreadyResolved`] — report is not in the
331    ///   `pending` state. Handler maps to a generic `InvalidRequest`.
332    pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
333        let (reply_tx, reply_rx) = oneshot::channel();
334        self.tx
335            .send(WriteCommand::ResolveReport(req, reply_tx))
336            .await
337            .map_err(|_| Error::Signing("writer task is shut down".into()))?;
338        reply_rx
339            .await
340            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
341    }
342
343    /// Subscribe to committed events. The returned receiver lags past
344    /// the internal broadcast buffer; the consumer (subscribeLabels, #7)
345    /// turns `RecvError::Lagged` into a connection close.
346    pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
347        self.broadcast_tx.subscribe()
348    }
349
350    /// Count of live broadcast receivers. Exposed for test sync points
351    /// (the WS handler subscribes asynchronously after upgrade; tests
352    /// that want to emit an event "after the subscriber is ready" poll
353    /// this until it reaches the expected count). Not a production hot
354    /// path — `broadcast::Sender::receiver_count` walks an atomic chain.
355    pub fn receiver_count(&self) -> usize {
356        self.broadcast_tx.receiver_count()
357    }
358
359    /// Observe writer lifecycle. The returned receiver starts at `false`
360    /// and flips to `true` once exactly, when the writer has accepted a
361    /// shutdown command and is about to release the lease. Holders close
362    /// downstream connections gracefully on the `true` transition.
363    pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
364        self.shutdown_rx.clone()
365    }
366
367    /// Explicit shutdown. Drains in-flight writes, releases the lease
368    /// row, stops the heartbeat, and returns. Idempotent-ish: calling on
369    /// an already-shut-down writer returns a channel-closed error, not a
370    /// panic.
371    pub async fn shutdown(&self) -> Result<()> {
372        let (reply_tx, reply_rx) = oneshot::channel();
373        self.tx
374            .send(WriteCommand::Shutdown(reply_tx))
375            .await
376            .map_err(|_| Error::Signing("writer task is already shut down".into()))?;
377        reply_rx
378            .await
379            .map_err(|_| Error::Signing("writer dropped reply channel".into()))?
380    }
381}
382
383/// Start the writer task. Acquires the instance lease, bootstraps the
384/// `signing_keys` row if empty, and spawns the event loop + heartbeat.
385///
386/// `service_did` stamps `labels.src` on every emitted event. `key` must
387/// be the private key whose public form is recorded (or will be recorded)
388/// in `signing_keys`.
389pub async fn spawn(
390    pool: Pool<Sqlite>,
391    key: SigningKey,
392    service_did: String,
393) -> Result<WriterHandle> {
394    let instance_id = acquire_lease(&pool).await?;
395    let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
396
397    let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
398    let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
399    let (shutdown_tx, shutdown_rx) = watch::channel(false);
400
401    let writer = Writer {
402        pool,
403        key,
404        service_did,
405        signing_key_id,
406        instance_id,
407        rx,
408        broadcast_tx: broadcast_tx.clone(),
409        shutdown_tx,
410    };
411
412    tokio::spawn(writer.run());
413
414    Ok(WriterHandle {
415        tx,
416        broadcast_tx,
417        shutdown_rx,
418    })
419}
420
421// ---------- Writer task ----------
422
423struct Writer {
424    pool: Pool<Sqlite>,
425    key: SigningKey,
426    service_did: String,
427    signing_key_id: i64,
428    instance_id: String,
429    rx: mpsc::Receiver<WriteCommand>,
430    broadcast_tx: broadcast::Sender<LabelEvent>,
431    shutdown_tx: watch::Sender<bool>,
432}
433
434impl Writer {
435    async fn run(mut self) {
436        let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
437        heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
438        // First tick fires immediately; skip it — we just acquired the lease
439        // with a fresh timestamp, so touching it again is redundant.
440        heartbeat_timer.tick().await;
441
442        loop {
443            tokio::select! {
444                biased;
445                // Prefer draining write commands over heartbeats so a
446                // steady-state of inbound work does not starve itself
447                // behind a background housekeeping task.
448                cmd = self.rx.recv() => {
449                    match cmd {
450                        Some(WriteCommand::Apply(req, reply)) => {
451                            let res = self.handle_apply(req).await;
452                            // Caller may have cancelled; dropping the reply is fine.
453                            let _ = reply.send(res);
454                        }
455                        Some(WriteCommand::Negate(req, reply)) => {
456                            let res = self.handle_negate(req).await;
457                            let _ = reply.send(res);
458                        }
459                        Some(WriteCommand::ResolveReport(req, reply)) => {
460                            let res = self.handle_resolve_report(req).await;
461                            let _ = reply.send(res);
462                        }
463                        Some(WriteCommand::Shutdown(reply)) => {
464                            // Flip the shutdown watch *before* releasing the
465                            // lease so subscriber tasks see the signal while
466                            // the DB is still accessible for their close-
467                            // frame sends.
468                            let _ = self.shutdown_tx.send(true);
469                            let res = self.release_lease().await;
470                            let _ = reply.send(res);
471                            return;
472                        }
473                        None => {
474                            // All handles dropped without explicit shutdown.
475                            // Best-effort lease release so a same-process
476                            // restart doesn't trip the 60s wait.
477                            let _ = self.shutdown_tx.send(true);
478                            if let Err(e) = self.release_lease().await {
479                                tracing::error!("lease release on handle drop: {e}");
480                            }
481                            return;
482                        }
483                    }
484                }
485                _ = heartbeat_timer.tick() => {
486                    if let Err(e) = self.heartbeat().await {
487                        // Transient DB errors are logged, not fatal. A
488                        // sustained failure is caught at the next-instance
489                        // startup check — not at this writer's expense.
490                        tracing::error!("lease heartbeat failed: {e}");
491                    }
492                }
493            }
494        }
495    }
496
497    async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
498        let mut tx = self.pool.begin().await?;
499        let created_at = epoch_ms_now();
500        let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
501        tx.commit().await?;
502        // No-receivers is not a write failure (§plan point G).
503        let _ = self.broadcast_tx.send(event.clone());
504        Ok(event)
505    }
506
507    /// Inner label-application pipeline: reserve seq → clamp cts →
508    /// sign → INSERT label → INSERT audit. Does NOT commit the
509    /// transaction and does NOT broadcast — caller handles both.
510    ///
511    /// Extracted from `handle_apply` so `handle_resolve_report` can
512    /// reuse it inside the same transaction as the report update,
513    /// preserving §F5's single-writer-owns-label-emission invariant
514    /// (this helper only runs on the writer task; no other code path
515    /// can access seq allocation or cts clamping).
516    async fn apply_label_inner(
517        &self,
518        tx: &mut sqlx::Transaction<'_, Sqlite>,
519        req: &ApplyLabelRequest,
520        created_at_ms: i64,
521    ) -> Result<LabelEvent> {
522        let seq = reserve_seq(tx).await?;
523        let prev_cts: Option<String> = sqlx::query_scalar!(
524            // sqlx type override — MAX() strips column origin metadata so
525            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
526            // forces the nullable wrapper (Option<String>).
527            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
528               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
529            self.service_did,
530            req.uri,
531            req.val,
532        )
533        .fetch_one(&mut **tx)
534        .await?;
535
536        let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
537
538        let mut label = Label {
539            ver: 1,
540            src: self.service_did.clone(),
541            uri: req.uri.clone(),
542            cid: req.cid.clone(),
543            val: req.val.clone(),
544            neg: false,
545            cts: cts.clone(),
546            exp: req.exp.clone(),
547            sig: None,
548        };
549        label.sig = Some(sign_label(&self.key, &label)?);
550        let sig_bytes = label.sig.expect("just set").to_vec();
551
552        let neg_int: i64 = 0;
553        sqlx::query!(
554            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
555             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
556            seq,
557            label.ver,
558            label.src,
559            label.uri,
560            label.cid,
561            label.val,
562            neg_int,
563            label.cts,
564            label.exp,
565            sig_bytes,
566            self.signing_key_id,
567            created_at_ms,
568        )
569        .execute(&mut **tx)
570        .await?;
571
572        let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
573        let action = "label_applied";
574        let outcome = "success";
575        sqlx::query!(
576            "INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
577             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
578            created_at_ms,
579            action,
580            req.actor_did,
581            req.uri,
582            req.cid,
583            outcome,
584            audit_reason,
585        )
586        .execute(&mut **tx)
587        .await?;
588
589        Ok(LabelEvent { seq, label })
590    }
591
592    async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
593        let mut tx = self.pool.begin().await?;
594
595        // Most-recent event for the tuple. Not-found OR latest-is-already-
596        // a-negation both surface as LabelNotFound (§F6).
597        let latest = sqlx::query!(
598            "SELECT neg, cid FROM labels
599             WHERE src = ?1 AND uri = ?2 AND val = ?3
600             ORDER BY seq DESC LIMIT 1",
601            self.service_did,
602            req.uri,
603            req.val,
604        )
605        .fetch_optional(&mut *tx)
606        .await?;
607
608        let cid = match latest {
609            Some(row) if row.neg == 0 => row.cid,
610            _ => {
611                return Err(Error::LabelNotFound {
612                    src: self.service_did.clone(),
613                    uri: req.uri,
614                    val: req.val,
615                });
616            }
617        };
618
619        let seq = reserve_seq(&mut tx).await?;
620        // Prev cts lookup is the same query as apply; tuple uniqueness is
621        // on (src, uri, val) regardless of neg.
622        let prev_cts: Option<String> = sqlx::query_scalar!(
623            // sqlx type override — MAX() strips column origin metadata so
624            // sqlx can't infer the TEXT+nullable shape on its own. `?:`
625            // forces the nullable wrapper (Option<String>).
626            r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
627               WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
628            self.service_did,
629            req.uri,
630            req.val,
631        )
632        .fetch_one(&mut *tx)
633        .await?;
634
635        let wall_now_ms = epoch_ms_now();
636        let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
637
638        let mut label = Label {
639            ver: 1,
640            src: self.service_did.clone(),
641            uri: req.uri.clone(),
642            cid: cid.clone(),
643            val: req.val.clone(),
644            neg: true,
645            cts: cts.clone(),
646            exp: None, // Negations never carry an expiry.
647            sig: None,
648        };
649        label.sig = Some(sign_label(&self.key, &label)?);
650        let sig_bytes = label.sig.expect("just set").to_vec();
651
652        let created_at = wall_now_ms;
653        let neg_int: i64 = 1;
654        sqlx::query!(
655            "INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
656             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
657            seq,
658            label.ver,
659            label.src,
660            label.uri,
661            label.cid,
662            label.val,
663            neg_int,
664            label.cts,
665            label.exp,
666            sig_bytes,
667            self.signing_key_id,
668            created_at,
669        )
670        .execute(&mut *tx)
671        .await?;
672
673        let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
674        let action = "label_negated";
675        let outcome = "success";
676        sqlx::query!(
677            "INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
678             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
679            created_at,
680            action,
681            req.actor_did,
682            req.uri,
683            cid,
684            outcome,
685            audit_reason,
686        )
687        .execute(&mut *tx)
688        .await?;
689
690        tx.commit().await?;
691
692        let event = LabelEvent { seq, label };
693        let _ = self.broadcast_tx.send(event.clone());
694        Ok(event)
695    }
696
697    /// Atomic resolveReport flow — §F12 requires label INSERT +
698    /// audit(label_applied) + report UPDATE + audit(report_resolved)
699    /// all commit together or not at all. Everything happens inside
700    /// one SQLite transaction on the writer task, preserving §F5's
701    /// single-writer invariant for the optional label emission.
702    ///
703    /// Returns [`Error::ReportNotFound`] if the report doesn't
704    /// exist, [`Error::ReportAlreadyResolved`] if it is not in
705    /// `status='pending'`. Label-value validation is a handler
706    /// concern (the handler consults `AdminConfig.label_values`
707    /// before sending the command here — saves a writer round-trip
708    /// on invalid input and keeps the anti-leak message local).
709    async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
710        let mut tx = self.pool.begin().await?;
711
712        // 1. Load the report; confirm pending status.
713        let current = sqlx::query_as!(
714            crate::report::Report,
715            r#"SELECT
716                 id                 AS "id!: i64",
717                 created_at         AS "created_at!: String",
718                 reported_by        AS "reported_by!: String",
719                 reason_type        AS "reason_type!: String",
720                 reason,
721                 subject_type       AS "subject_type!: String",
722                 subject_did        AS "subject_did!: String",
723                 subject_uri,
724                 subject_cid,
725                 status             AS "status!: String",
726                 resolved_at,
727                 resolved_by,
728                 resolution_label,
729                 resolution_reason
730               FROM reports WHERE id = ?1"#,
731            req.report_id,
732        )
733        .fetch_optional(&mut *tx)
734        .await?;
735
736        let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
737        if report.status != "pending" {
738            return Err(Error::ReportAlreadyResolved { id: req.report_id });
739        }
740
741        let created_at = epoch_ms_now();
742
743        // 2. Optional label-apply BEFORE the report UPDATE so audit
744        // row insertion order reflects logical sequence
745        // (label_applied, then report_resolved — §G per #15 criteria).
746        let label_event = if let Some(apply) = &req.apply_label {
747            let apply_req = ApplyLabelRequest {
748                actor_did: req.actor_did.clone(),
749                uri: apply.uri.clone(),
750                cid: apply.cid.clone(),
751                val: apply.val.clone(),
752                exp: apply.exp.clone(),
753                // The label's own audit_reason JSON captures {val,
754                // neg, moderator_reason}; the resolution-level reason
755                // goes on the report_resolved audit row below.
756                moderator_reason: None,
757            };
758            Some(
759                self.apply_label_inner(&mut tx, &apply_req, created_at)
760                    .await?,
761            )
762        } else {
763            None
764        };
765
766        // 3. UPDATE reports.
767        let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
768        let resolution_label = req.apply_label.as_ref().map(|a| a.val.clone());
769        sqlx::query!(
770            "UPDATE reports SET
771                status = 'resolved',
772                resolved_at = ?1,
773                resolved_by = ?2,
774                resolution_label = ?3,
775                resolution_reason = ?4
776             WHERE id = ?5",
777            resolved_at_rfc,
778            req.actor_did,
779            resolution_label,
780            req.resolution_reason,
781            req.report_id,
782        )
783        .execute(&mut *tx)
784        .await?;
785
786        // 4. Audit: report_resolved.
787        let audit_reason = build_resolve_audit_reason(
788            req.apply_label.as_ref().map(|a| a.val.as_str()),
789            req.resolution_reason.as_deref(),
790        );
791        let report_id_str = req.report_id.to_string();
792        let action = "report_resolved";
793        let outcome = "success";
794        sqlx::query!(
795            "INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
796             VALUES (?1, ?2, ?3, ?4, NULL, ?5, ?6)",
797            created_at,
798            action,
799            req.actor_did,
800            report_id_str,
801            outcome,
802            audit_reason,
803        )
804        .execute(&mut *tx)
805        .await?;
806
807        tx.commit().await?;
808
809        // 5. Broadcast (post-commit, per §F5 broadcast-after-commit rule).
810        if let Some(event) = &label_event {
811            let _ = self.broadcast_tx.send(event.clone());
812        }
813
814        // 6. Mutate the loaded report struct to reflect the committed
815        // state and return it. Avoids a re-fetch round-trip since we
816        // know exactly what changed.
817        report.status = "resolved".to_string();
818        report.resolved_at = Some(resolved_at_rfc);
819        report.resolved_by = Some(req.actor_did);
820        report.resolution_label = resolution_label;
821        report.resolution_reason = req.resolution_reason;
822
823        Ok(ResolvedReport {
824            report,
825            label_event,
826        })
827    }
828
829    async fn heartbeat(&self) -> Result<()> {
830        let now_ms = epoch_ms_now();
831        sqlx::query!(
832            "UPDATE server_instance_lease SET last_heartbeat = ?1
833             WHERE id = 1 AND instance_id = ?2",
834            now_ms,
835            self.instance_id,
836        )
837        .execute(&self.pool)
838        .await?;
839        Ok(())
840    }
841
842    async fn release_lease(&self) -> Result<()> {
843        sqlx::query!(
844            "DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
845            self.instance_id,
846        )
847        .execute(&self.pool)
848        .await?;
849        Ok(())
850    }
851}
852
853// ---------- Lease + key bootstrap ----------
854
855async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
856    let now_ms = epoch_ms_now();
857
858    let existing =
859        sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
860            .fetch_optional(pool)
861            .await?;
862
863    if let Some(row) = existing {
864        let age_ms = (now_ms - row.last_heartbeat).max(0);
865        if age_ms < LEASE_STALE_MS {
866            return Err(Error::LeaseHeld {
867                instance_id: row.instance_id,
868                age_secs: (age_ms / 1000) as u64,
869            });
870        }
871    }
872
873    let new_id = Uuid::new_v4().to_string();
874    // INSERT OR REPLACE on id=1: takes over a stale lease row if present,
875    // or creates the row on first startup. Legal because the staleness
876    // check above confirms no live peer.
877    sqlx::query!(
878        "INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
879         VALUES (1, ?1, ?2, ?2)
880         ON CONFLICT(id) DO UPDATE SET
881             instance_id = excluded.instance_id,
882             acquired_at = excluded.acquired_at,
883             last_heartbeat = excluded.last_heartbeat",
884        new_id,
885        now_ms,
886    )
887    .execute(pool)
888    .await?;
889
890    Ok(new_id)
891}
892
893async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
894    let kp = K256Keypair::from_private_key(key.expose_secret())?;
895    let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
896
897    let existing =
898        sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
899            .fetch_optional(pool)
900            .await?;
901
902    if let Some(row) = existing {
903        if row.public_key_multibase != my_multibase {
904            return Err(Error::Signing(format!(
905                "signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
906                row.public_key_multibase, my_multibase
907            )));
908        }
909        return Ok(row.id);
910    }
911
912    let now_ms = epoch_ms_now();
913    let valid_from = rfc3339_from_epoch_ms(now_ms)?;
914    let id = sqlx::query_scalar!(
915        "INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
916         VALUES (?1, ?2, NULL, ?3)
917         RETURNING id",
918        my_multibase,
919        valid_from,
920        now_ms,
921    )
922    .fetch_one(pool)
923    .await?;
924
925    Ok(id)
926}
927
928// ---------- Pure helpers (unit-testable without a DB) ----------
929
930async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
931    // label_sequence has a single auto-increment column; `DEFAULT VALUES`
932    // triggers a fresh allocation. Under AUTOINCREMENT the assigned seq
933    // is never reused even across rolled-back transactions — this is what
934    // gives us "strictly monotonic, no gaps from writer's perspective."
935    let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
936        .fetch_one(tx)
937        .await?;
938    Ok(seq)
939}
940
941/// §6.1 monotonicity clamp. `prev_cts_str`, when present, is expected to
942/// be in the writer's canonical `CTS_FORMAT`.
943fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
944    let effective_ms = match prev_cts_str {
945        Some(s) => {
946            let prev_ms = parse_rfc3339_ms(s)?;
947            wall_now_ms.max(prev_ms + 1)
948        }
949        None => wall_now_ms,
950    };
951    rfc3339_from_epoch_ms(effective_ms)
952}
953
954/// Epoch-ms → RFC-3339 Z with millisecond precision. `pub(crate)` so
955/// peer modules (admin/audit_view, future retention sweep) share the
956/// single formatter used on the writer's timestamp boundary.
957pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
958    let nanos: i128 = (ms as i128) * 1_000_000;
959    let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
960        .map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
961    let formatted = dt
962        .format(&CTS_FORMAT)
963        .map_err(|e| Error::Signing(format!("format cts: {e}")))?;
964    Ok(format!("{formatted}Z"))
965}
966
967/// RFC-3339 Z with millisecond precision → epoch-ms. `pub(crate)` so
968/// admin handlers can validate `since` / `until` query params against
969/// the same parser the writer uses on its input boundary.
970pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
971    let stripped = s
972        .strip_suffix('Z')
973        .ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
974    let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
975        .map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
976    let nanos = pdt.assume_utc().unix_timestamp_nanos();
977    Ok((nanos / 1_000_000) as i64)
978}
979
980/// Current wall-clock time as Unix epoch milliseconds. `pub(crate)` so
981/// peer modules (server, future retention sweep) share one implementation
982/// rather than each inlining `SystemTime::now()` with its own error
983/// handling.
984pub(crate) fn epoch_ms_now() -> i64 {
985    SystemTime::now()
986        .duration_since(UNIX_EPOCH)
987        .expect("system clock before unix epoch")
988        .as_millis() as i64
989}
990
991fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
992    let body = serde_json::json!({
993        "val": val,
994        "neg": neg,
995        "moderator_reason": moderator_reason,
996    });
997    body.to_string()
998}
999
1000/// `report_resolved` audit reason JSON — see [`AUDIT_REASON_RESOLVE_REPORT`]
1001/// for the schema.
1002fn build_resolve_audit_reason(
1003    applied_label_val: Option<&str>,
1004    resolution_reason: Option<&str>,
1005) -> String {
1006    serde_json::json!({
1007        "applied_label_val": applied_label_val,
1008        "resolution_reason": resolution_reason,
1009    })
1010    .to_string()
1011}
1012
1013/// `reporter_flagged` / `reporter_unflagged` audit reason JSON —
1014/// see [`AUDIT_REASON_FLAG_REPORTER`]. Shared with the flagReporter
1015/// handler (not the writer — flagReporter is a direct handler txn
1016/// per the #15 criteria).
1017pub(crate) fn build_flag_reporter_audit_reason(
1018    did: &str,
1019    suppressed: bool,
1020    moderator_reason: Option<&str>,
1021) -> String {
1022    serde_json::json!({
1023        "did": did,
1024        "suppressed": suppressed,
1025        "moderator_reason": moderator_reason,
1026    })
1027    .to_string()
1028}
1029
1030#[cfg(test)]
1031mod tests {
1032    use super::*;
1033
1034    // 1_715_000_000_000 ms since epoch = 2024-05-06T12:53:20.000Z UTC.
1035    // The four tests below pin every branch of the clamp: no prior,
1036    // wall-ahead, wall-behind (use prev+1ms), wall-equal (advance anyway).
1037
1038    #[test]
1039    fn clamp_cts_uses_wall_clock_when_no_prior() {
1040        let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
1041        assert_eq!(out, "2024-05-06T12:53:20.000Z");
1042    }
1043
1044    #[test]
1045    fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
1046        // wall_clock one second behind prev: result is prev + 1ms, not wall.
1047        let prev = "2024-05-06T12:53:20.500Z";
1048        let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
1049        assert_eq!(out, "2024-05-06T12:53:20.501Z");
1050    }
1051
1052    #[test]
1053    fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
1054        let prev = "2024-05-06T12:53:20.500Z";
1055        let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
1056        assert_eq!(out, "2024-05-06T12:53:22.000Z");
1057    }
1058
1059    #[test]
1060    fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
1061        // wall_clock == prev exactly: must advance by 1ms (strictly greater).
1062        let prev = "2024-05-06T12:53:20.500Z";
1063        let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
1064        assert_eq!(out, "2024-05-06T12:53:20.501Z");
1065    }
1066
1067    #[test]
1068    fn build_audit_reason_shape_matches_documented_schema() {
1069        let json = build_audit_reason("spam", false, Some("user reported"));
1070        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
1071        assert_eq!(v["val"], "spam");
1072        assert_eq!(v["neg"], false);
1073        assert_eq!(v["moderator_reason"], "user reported");
1074    }
1075
1076    #[test]
1077    fn build_audit_reason_null_when_moderator_reason_absent() {
1078        let json = build_audit_reason("spam", true, None);
1079        let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
1080        assert!(v["moderator_reason"].is_null());
1081    }
1082
1083    #[test]
1084    fn rfc3339_roundtrip() {
1085        let s = "2026-04-22T12:00:00.938Z";
1086        let ms = parse_rfc3339_ms(s).expect("parse");
1087        let back = rfc3339_from_epoch_ms(ms).expect("format");
1088        assert_eq!(back, s);
1089    }
1090}