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