use std::time::{Duration, SystemTime, UNIX_EPOCH};
use proto_blue_crypto::{K256Keypair, Keypair as _, format_multikey};
use sqlx::{Pool, Sqlite};
use time::format_description::FormatItem;
use time::macros::format_description;
use time::{OffsetDateTime, PrimitiveDateTime};
use tokio::sync::{broadcast, mpsc, oneshot, watch};
use tokio::time::{MissedTickBehavior, interval};
use uuid::Uuid;
use crate::error::{Error, Result};
use crate::label::Label;
use crate::signing::sign_label;
use crate::signing_key::SigningKey;
const LEASE_STALE_MS: i64 = 60_000;
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
const COMMAND_BUFFER: usize = 64;
const BROADCAST_BUFFER: usize = 1024;
#[doc(alias = "audit_log.reason")]
pub const AUDIT_REASON_SCHEMA: &str =
"label_applied / label_negated: { val, neg, moderator_reason }";
#[doc(alias = "audit_log.reason.report_resolved")]
pub const AUDIT_REASON_RESOLVE_REPORT: &str =
"report_resolved: { applied_label_val, resolution_reason }";
#[doc(alias = "audit_log.reason.flag_reporter")]
pub const AUDIT_REASON_FLAG_REPORTER: &str =
"reporter_flagged / reporter_unflagged: { did, suppressed, moderator_reason }";
pub const AUDIT_ACTION_VALUES: &[&str] = &[
"label_applied",
"label_negated",
"report_resolved",
"reporter_flagged",
"reporter_unflagged",
];
pub const AUDIT_OUTCOME_VALUES: &[&str] = &["success", "failure"];
const CTS_FORMAT: &[FormatItem<'_>] =
format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
#[derive(Debug, Clone)]
pub struct ApplyLabelRequest {
pub actor_did: String,
pub uri: String,
pub cid: Option<String>,
pub val: String,
pub exp: Option<String>,
pub moderator_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct NegateLabelRequest {
pub actor_did: String,
pub uri: String,
pub val: String,
pub moderator_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct LabelEvent {
pub seq: i64,
pub label: Label,
}
enum WriteCommand {
Apply(ApplyLabelRequest, oneshot::Sender<Result<LabelEvent>>),
Negate(NegateLabelRequest, oneshot::Sender<Result<LabelEvent>>),
ResolveReport(
ResolveReportRequest,
oneshot::Sender<Result<ResolvedReport>>,
),
Shutdown(oneshot::Sender<Result<()>>),
}
#[derive(Debug, Clone)]
pub struct ApplyLabelInline {
pub uri: String,
pub cid: Option<String>,
pub val: String,
pub exp: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ResolveReportRequest {
pub actor_did: String,
pub report_id: i64,
pub apply_label: Option<ApplyLabelInline>,
pub resolution_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ResolvedReport {
pub report: crate::report::Report,
pub label_event: Option<LabelEvent>,
}
#[derive(Debug, Clone)]
pub struct WriterHandle {
tx: mpsc::Sender<WriteCommand>,
broadcast_tx: broadcast::Sender<LabelEvent>,
shutdown_rx: watch::Receiver<bool>,
}
impl WriterHandle {
pub async fn apply_label(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Apply(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn negate_label(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Negate(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub async fn resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::ResolveReport(req, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
pub fn subscribe(&self) -> broadcast::Receiver<LabelEvent> {
self.broadcast_tx.subscribe()
}
pub fn receiver_count(&self) -> usize {
self.broadcast_tx.receiver_count()
}
pub fn shutdown_signal(&self) -> watch::Receiver<bool> {
self.shutdown_rx.clone()
}
pub async fn shutdown(&self) -> Result<()> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Shutdown(reply_tx))
.await
.map_err(|_| Error::Signing("writer task is already shut down".into()))?;
reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))?
}
}
pub async fn spawn(
pool: Pool<Sqlite>,
key: SigningKey,
service_did: String,
) -> Result<WriterHandle> {
let instance_id = acquire_lease(&pool).await?;
let signing_key_id = ensure_signing_key_row(&pool, &key).await?;
let (tx, rx) = mpsc::channel(COMMAND_BUFFER);
let (broadcast_tx, _first_rx) = broadcast::channel(BROADCAST_BUFFER);
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let writer = Writer {
pool,
key,
service_did,
signing_key_id,
instance_id,
rx,
broadcast_tx: broadcast_tx.clone(),
shutdown_tx,
};
tokio::spawn(writer.run());
Ok(WriterHandle {
tx,
broadcast_tx,
shutdown_rx,
})
}
struct Writer {
pool: Pool<Sqlite>,
key: SigningKey,
service_did: String,
signing_key_id: i64,
instance_id: String,
rx: mpsc::Receiver<WriteCommand>,
broadcast_tx: broadcast::Sender<LabelEvent>,
shutdown_tx: watch::Sender<bool>,
}
impl Writer {
async fn run(mut self) {
let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
heartbeat_timer.tick().await;
loop {
tokio::select! {
biased;
cmd = self.rx.recv() => {
match cmd {
Some(WriteCommand::Apply(req, reply)) => {
let res = self.handle_apply(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::Negate(req, reply)) => {
let res = self.handle_negate(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::ResolveReport(req, reply)) => {
let res = self.handle_resolve_report(req).await;
let _ = reply.send(res);
}
Some(WriteCommand::Shutdown(reply)) => {
let _ = self.shutdown_tx.send(true);
let res = self.release_lease().await;
let _ = reply.send(res);
return;
}
None => {
let _ = self.shutdown_tx.send(true);
if let Err(e) = self.release_lease().await {
tracing::error!("lease release on handle drop: {e}");
}
return;
}
}
}
_ = heartbeat_timer.tick() => {
if let Err(e) = self.heartbeat().await {
tracing::error!("lease heartbeat failed: {e}");
}
}
}
}
}
async fn handle_apply(&self, req: ApplyLabelRequest) -> Result<LabelEvent> {
let mut tx = self.pool.begin().await?;
let created_at = epoch_ms_now();
let event = self.apply_label_inner(&mut tx, &req, created_at).await?;
tx.commit().await?;
let _ = self.broadcast_tx.send(event.clone());
Ok(event)
}
async fn apply_label_inner(
&self,
tx: &mut sqlx::Transaction<'_, Sqlite>,
req: &ApplyLabelRequest,
created_at_ms: i64,
) -> Result<LabelEvent> {
let seq = reserve_seq(tx).await?;
let prev_cts: Option<String> = sqlx::query_scalar!(
r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
self.service_did,
req.uri,
req.val,
)
.fetch_one(&mut **tx)
.await?;
let cts = clamp_cts(created_at_ms, prev_cts.as_deref())?;
let mut label = Label {
ver: 1,
src: self.service_did.clone(),
uri: req.uri.clone(),
cid: req.cid.clone(),
val: req.val.clone(),
neg: false,
cts: cts.clone(),
exp: req.exp.clone(),
sig: None,
};
label.sig = Some(sign_label(&self.key, &label)?);
let sig_bytes = label.sig.expect("just set").to_vec();
let neg_int: i64 = 0;
sqlx::query!(
"INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
seq,
label.ver,
label.src,
label.uri,
label.cid,
label.val,
neg_int,
label.cts,
label.exp,
sig_bytes,
self.signing_key_id,
created_at_ms,
)
.execute(&mut **tx)
.await?;
let audit_reason = build_audit_reason(&req.val, false, req.moderator_reason.as_deref());
let action = "label_applied";
let outcome = "success";
sqlx::query!(
"INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
created_at_ms,
action,
req.actor_did,
req.uri,
req.cid,
outcome,
audit_reason,
)
.execute(&mut **tx)
.await?;
Ok(LabelEvent { seq, label })
}
async fn handle_negate(&self, req: NegateLabelRequest) -> Result<LabelEvent> {
let mut tx = self.pool.begin().await?;
let latest = sqlx::query!(
"SELECT neg, cid FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3
ORDER BY seq DESC LIMIT 1",
self.service_did,
req.uri,
req.val,
)
.fetch_optional(&mut *tx)
.await?;
let cid = match latest {
Some(row) if row.neg == 0 => row.cid,
_ => {
return Err(Error::LabelNotFound {
src: self.service_did.clone(),
uri: req.uri,
val: req.val,
});
}
};
let seq = reserve_seq(&mut tx).await?;
let prev_cts: Option<String> = sqlx::query_scalar!(
r#"SELECT MAX(cts) AS "max_cts?: String" FROM labels
WHERE src = ?1 AND uri = ?2 AND val = ?3"#,
self.service_did,
req.uri,
req.val,
)
.fetch_one(&mut *tx)
.await?;
let wall_now_ms = epoch_ms_now();
let cts = clamp_cts(wall_now_ms, prev_cts.as_deref())?;
let mut label = Label {
ver: 1,
src: self.service_did.clone(),
uri: req.uri.clone(),
cid: cid.clone(),
val: req.val.clone(),
neg: true,
cts: cts.clone(),
exp: None, sig: None,
};
label.sig = Some(sign_label(&self.key, &label)?);
let sig_bytes = label.sig.expect("just set").to_vec();
let created_at = wall_now_ms;
let neg_int: i64 = 1;
sqlx::query!(
"INSERT INTO labels (seq, ver, src, uri, cid, val, neg, cts, exp, sig, signing_key_id, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
seq,
label.ver,
label.src,
label.uri,
label.cid,
label.val,
neg_int,
label.cts,
label.exp,
sig_bytes,
self.signing_key_id,
created_at,
)
.execute(&mut *tx)
.await?;
let audit_reason = build_audit_reason(&req.val, true, req.moderator_reason.as_deref());
let action = "label_negated";
let outcome = "success";
sqlx::query!(
"INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
created_at,
action,
req.actor_did,
req.uri,
cid,
outcome,
audit_reason,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
let event = LabelEvent { seq, label };
let _ = self.broadcast_tx.send(event.clone());
Ok(event)
}
async fn handle_resolve_report(&self, req: ResolveReportRequest) -> Result<ResolvedReport> {
let mut tx = self.pool.begin().await?;
let current = sqlx::query_as!(
crate::report::Report,
r#"SELECT
id AS "id!: i64",
created_at AS "created_at!: String",
reported_by AS "reported_by!: String",
reason_type AS "reason_type!: String",
reason,
subject_type AS "subject_type!: String",
subject_did AS "subject_did!: String",
subject_uri,
subject_cid,
status AS "status!: String",
resolved_at,
resolved_by,
resolution_label,
resolution_reason
FROM reports WHERE id = ?1"#,
req.report_id,
)
.fetch_optional(&mut *tx)
.await?;
let mut report = current.ok_or(Error::ReportNotFound { id: req.report_id })?;
if report.status != "pending" {
return Err(Error::ReportAlreadyResolved { id: req.report_id });
}
let created_at = epoch_ms_now();
let label_event = if let Some(apply) = &req.apply_label {
let apply_req = ApplyLabelRequest {
actor_did: req.actor_did.clone(),
uri: apply.uri.clone(),
cid: apply.cid.clone(),
val: apply.val.clone(),
exp: apply.exp.clone(),
moderator_reason: None,
};
Some(
self.apply_label_inner(&mut tx, &apply_req, created_at)
.await?,
)
} else {
None
};
let resolved_at_rfc = rfc3339_from_epoch_ms(created_at)?;
let resolution_label = req.apply_label.as_ref().map(|a| a.val.clone());
sqlx::query!(
"UPDATE reports SET
status = 'resolved',
resolved_at = ?1,
resolved_by = ?2,
resolution_label = ?3,
resolution_reason = ?4
WHERE id = ?5",
resolved_at_rfc,
req.actor_did,
resolution_label,
req.resolution_reason,
req.report_id,
)
.execute(&mut *tx)
.await?;
let audit_reason = build_resolve_audit_reason(
req.apply_label.as_ref().map(|a| a.val.as_str()),
req.resolution_reason.as_deref(),
);
let report_id_str = req.report_id.to_string();
let action = "report_resolved";
let outcome = "success";
sqlx::query!(
"INSERT INTO audit_log (created_at, action, actor_did, target, target_cid, outcome, reason)
VALUES (?1, ?2, ?3, ?4, NULL, ?5, ?6)",
created_at,
action,
req.actor_did,
report_id_str,
outcome,
audit_reason,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
if let Some(event) = &label_event {
let _ = self.broadcast_tx.send(event.clone());
}
report.status = "resolved".to_string();
report.resolved_at = Some(resolved_at_rfc);
report.resolved_by = Some(req.actor_did);
report.resolution_label = resolution_label;
report.resolution_reason = req.resolution_reason;
Ok(ResolvedReport {
report,
label_event,
})
}
async fn heartbeat(&self) -> Result<()> {
let now_ms = epoch_ms_now();
sqlx::query!(
"UPDATE server_instance_lease SET last_heartbeat = ?1
WHERE id = 1 AND instance_id = ?2",
now_ms,
self.instance_id,
)
.execute(&self.pool)
.await?;
Ok(())
}
async fn release_lease(&self) -> Result<()> {
sqlx::query!(
"DELETE FROM server_instance_lease WHERE id = 1 AND instance_id = ?1",
self.instance_id,
)
.execute(&self.pool)
.await?;
Ok(())
}
}
async fn acquire_lease(pool: &Pool<Sqlite>) -> Result<String> {
let now_ms = epoch_ms_now();
let existing =
sqlx::query!("SELECT instance_id, last_heartbeat FROM server_instance_lease WHERE id = 1")
.fetch_optional(pool)
.await?;
if let Some(row) = existing {
let age_ms = (now_ms - row.last_heartbeat).max(0);
if age_ms < LEASE_STALE_MS {
return Err(Error::LeaseHeld {
instance_id: row.instance_id,
age_secs: (age_ms / 1000) as u64,
});
}
}
let new_id = Uuid::new_v4().to_string();
sqlx::query!(
"INSERT INTO server_instance_lease (id, instance_id, acquired_at, last_heartbeat)
VALUES (1, ?1, ?2, ?2)
ON CONFLICT(id) DO UPDATE SET
instance_id = excluded.instance_id,
acquired_at = excluded.acquired_at,
last_heartbeat = excluded.last_heartbeat",
new_id,
now_ms,
)
.execute(pool)
.await?;
Ok(new_id)
}
async fn ensure_signing_key_row(pool: &Pool<Sqlite>, key: &SigningKey) -> Result<i64> {
let kp = K256Keypair::from_private_key(key.expose_secret())?;
let my_multibase = format_multikey("ES256K", &kp.public_key_compressed());
let existing =
sqlx::query!("SELECT id, public_key_multibase FROM signing_keys ORDER BY id LIMIT 1")
.fetch_optional(pool)
.await?;
if let Some(row) = existing {
if row.public_key_multibase != my_multibase {
return Err(Error::Signing(format!(
"signing_keys.public_key_multibase ({}) does not match the loaded signing key's derived public key ({}) — key rotation is v1.1 scope",
row.public_key_multibase, my_multibase
)));
}
return Ok(row.id);
}
let now_ms = epoch_ms_now();
let valid_from = rfc3339_from_epoch_ms(now_ms)?;
let id = sqlx::query_scalar!(
"INSERT INTO signing_keys (public_key_multibase, valid_from, valid_to, created_at)
VALUES (?1, ?2, NULL, ?3)
RETURNING id",
my_multibase,
valid_from,
now_ms,
)
.fetch_one(pool)
.await?;
Ok(id)
}
async fn reserve_seq(tx: &mut sqlx::SqliteConnection) -> Result<i64> {
let seq = sqlx::query_scalar!("INSERT INTO label_sequence DEFAULT VALUES RETURNING seq")
.fetch_one(tx)
.await?;
Ok(seq)
}
fn clamp_cts(wall_now_ms: i64, prev_cts_str: Option<&str>) -> Result<String> {
let effective_ms = match prev_cts_str {
Some(s) => {
let prev_ms = parse_rfc3339_ms(s)?;
wall_now_ms.max(prev_ms + 1)
}
None => wall_now_ms,
};
rfc3339_from_epoch_ms(effective_ms)
}
pub(crate) fn rfc3339_from_epoch_ms(ms: i64) -> Result<String> {
let nanos: i128 = (ms as i128) * 1_000_000;
let dt = OffsetDateTime::from_unix_timestamp_nanos(nanos)
.map_err(|e| Error::Signing(format!("epoch ms {ms} out of range: {e}")))?;
let formatted = dt
.format(&CTS_FORMAT)
.map_err(|e| Error::Signing(format!("format cts: {e}")))?;
Ok(format!("{formatted}Z"))
}
pub(crate) fn parse_rfc3339_ms(s: &str) -> Result<i64> {
let stripped = s
.strip_suffix('Z')
.ok_or_else(|| Error::Signing(format!("cts {s:?} missing trailing Z")))?;
let pdt = PrimitiveDateTime::parse(stripped, &CTS_FORMAT)
.map_err(|e| Error::Signing(format!("parse cts {s:?}: {e}")))?;
let nanos = pdt.assume_utc().unix_timestamp_nanos();
Ok((nanos / 1_000_000) as i64)
}
pub(crate) fn epoch_ms_now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system clock before unix epoch")
.as_millis() as i64
}
fn build_audit_reason(val: &str, neg: bool, moderator_reason: Option<&str>) -> String {
let body = serde_json::json!({
"val": val,
"neg": neg,
"moderator_reason": moderator_reason,
});
body.to_string()
}
fn build_resolve_audit_reason(
applied_label_val: Option<&str>,
resolution_reason: Option<&str>,
) -> String {
serde_json::json!({
"applied_label_val": applied_label_val,
"resolution_reason": resolution_reason,
})
.to_string()
}
pub(crate) fn build_flag_reporter_audit_reason(
did: &str,
suppressed: bool,
moderator_reason: Option<&str>,
) -> String {
serde_json::json!({
"did": did,
"suppressed": suppressed,
"moderator_reason": moderator_reason,
})
.to_string()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn clamp_cts_uses_wall_clock_when_no_prior() {
let out = clamp_cts(1_715_000_000_000, None).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.000Z");
}
#[test]
fn clamp_cts_advances_one_ms_past_prior_when_wall_clock_lags() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_000 - 1000, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.501Z");
}
#[test]
fn clamp_cts_uses_wall_clock_when_ahead_of_prior() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_000 + 2000, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:22.000Z");
}
#[test]
fn clamp_cts_handles_equal_wall_and_prior_by_advancing() {
let prev = "2024-05-06T12:53:20.500Z";
let out = clamp_cts(1_715_000_000_500, Some(prev)).expect("clamp");
assert_eq!(out, "2024-05-06T12:53:20.501Z");
}
#[test]
fn build_audit_reason_shape_matches_documented_schema() {
let json = build_audit_reason("spam", false, Some("user reported"));
let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
assert_eq!(v["val"], "spam");
assert_eq!(v["neg"], false);
assert_eq!(v["moderator_reason"], "user reported");
}
#[test]
fn build_audit_reason_null_when_moderator_reason_absent() {
let json = build_audit_reason("spam", true, None);
let v: serde_json::Value = serde_json::from_str(&json).expect("parse");
assert!(v["moderator_reason"].is_null());
}
#[test]
fn rfc3339_roundtrip() {
let s = "2026-04-22T12:00:00.938Z";
let ms = parse_rfc3339_ms(s).expect("parse");
let back = rfc3339_from_epoch_ms(ms).expect("format");
assert_eq!(back, s);
}
}