use std::time::{Duration, Instant, 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::server::RetentionConfig;
use crate::signing::sign_label;
use crate::signing_key::SigningKey;
pub(crate) const LEASE_STALE_MS: i64 = 60_000;
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
const SWEEP_CHECK_INTERVAL: Duration = Duration::from_secs(60);
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 }";
#[doc(alias = "audit_log.reason.retention_sweep")]
pub const AUDIT_REASON_RETENTION_SWEEP: &str =
"retention_sweep: { rows_deleted, batches, duration_ms, retention_days_applied }";
pub const AUDIT_ACTION_VALUES: &[&str] = &[
"label_applied",
"label_negated",
"report_resolved",
"reporter_flagged",
"reporter_unflagged",
"retention_sweep",
];
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>>,
),
Sweep(SweepRequest, oneshot::Sender<Result<SweepBatchResult>>),
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 enum ResolutionAction {
Dismiss,
ApplyLabel(ApplyLabelInline),
}
impl ResolutionAction {
pub fn as_apply(&self) -> Option<&ApplyLabelInline> {
match self {
ResolutionAction::Dismiss => None,
ResolutionAction::ApplyLabel(a) => Some(a),
}
}
}
#[derive(Debug, Clone)]
pub struct ResolveReportRequest {
pub actor_did: String,
pub report_id: i64,
pub action: ResolutionAction,
pub resolution_reason: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ResolvedReport {
pub report: crate::report::Report,
pub label_event: Option<LabelEvent>,
}
#[derive(Debug, Clone, Default)]
pub struct SweepRequest;
#[derive(Debug, Clone)]
pub struct SweepBatchResult {
pub rows_deleted: i64,
pub has_more: bool,
pub retention_days_applied: Option<u32>,
}
#[derive(Debug, Clone)]
pub struct SweepResult {
pub rows_deleted: i64,
pub batches: u64,
pub duration_ms: u64,
pub retention_days_applied: Option<u32>,
}
#[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 async fn sweep(&self, _req: SweepRequest) -> Result<SweepResult> {
let start = std::time::Instant::now();
let mut total_rows: i64 = 0;
let mut batches: u64 = 0;
loop {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCommand::Sweep(SweepRequest, reply_tx))
.await
.map_err(|_| Error::Signing("writer task is shut down".into()))?;
let batch = reply_rx
.await
.map_err(|_| Error::Signing("writer dropped reply channel".into()))??;
total_rows += batch.rows_deleted;
batches += 1;
if !batch.has_more {
return Ok(SweepResult {
rows_deleted: total_rows,
batches,
duration_ms: start.elapsed().as_millis() as u64,
retention_days_applied: batch.retention_days_applied,
});
}
}
}
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,
retention_days: Option<u32>,
retention: RetentionConfig,
) -> 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,
retention_days,
retention,
};
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>,
retention_days: Option<u32>,
retention: RetentionConfig,
}
struct SweepRunState {
started_at: Instant,
rows: i64,
batches: u64,
}
impl Writer {
fn compute_next_sweep_fire(&self) -> Option<Instant> {
if !self.retention.sweep_enabled || self.retention_days.is_none() {
return None;
}
let now_utc = OffsetDateTime::now_utc();
let target_hour = self.retention.sweep_run_at_utc_hour;
let target_today_time = time::Time::from_hms(target_hour, 0, 0)
.expect("hour validated < 24 by Config::validate");
let target_today = now_utc.replace_time(target_today_time);
let target_dt = if target_today <= now_utc {
target_today + time::Duration::days(1)
} else {
target_today
};
let wait = target_dt - now_utc;
let wait_secs = wait.whole_seconds().max(0) as u64;
Some(Instant::now() + Duration::from_secs(wait_secs))
}
async fn run(mut self) {
let mut heartbeat_timer = interval(HEARTBEAT_INTERVAL);
heartbeat_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
heartbeat_timer.tick().await;
let mut sweep_check_timer = interval(SWEEP_CHECK_INTERVAL);
sweep_check_timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
sweep_check_timer.tick().await;
let mut next_scheduled_fire: Option<Instant> = self.compute_next_sweep_fire();
let mut sweep_state: Option<SweepRunState> = None;
loop {
let sweep_in_progress = sweep_state.is_some();
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::Sweep(req, reply)) => {
let res = self.handle_sweep(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}");
}
}
_ = sweep_check_timer.tick() => {
if sweep_state.is_none()
&& let Some(fire_at) = next_scheduled_fire
&& Instant::now() >= fire_at
{
sweep_state = Some(SweepRunState {
started_at: Instant::now(),
rows: 0,
batches: 0,
});
next_scheduled_fire = Some(fire_at + Duration::from_secs(86_400));
tracing::info!(
retention_days = ?self.retention_days,
sweep_batch_size = self.retention.sweep_batch_size,
"scheduled retention sweep starting"
);
}
}
_ = std::future::ready(()), if sweep_in_progress => {
match self.handle_sweep(SweepRequest).await {
Ok(batch) => {
let state = sweep_state
.as_mut()
.expect("sweep_in_progress => sweep_state Some");
state.rows += batch.rows_deleted;
state.batches += 1;
if !batch.has_more {
let final_state = sweep_state.take().expect("just set");
tracing::info!(
rows_deleted = final_state.rows,
batches = final_state.batches,
duration_ms = final_state.started_at.elapsed().as_millis() as u64,
retention_days_applied = ?batch.retention_days_applied,
"scheduled retention sweep complete"
);
}
}
Err(e) => {
let final_state = sweep_state.take();
tracing::error!(
error = %e,
batches_done = final_state.as_ref().map(|s| s.batches).unwrap_or(0),
rows_so_far = final_state.as_ref().map(|s| s.rows).unwrap_or(0),
"scheduled retention sweep batch failed; aborting run (will retry on next schedule)"
);
}
}
}
}
}
}
async fn handle_sweep(&self, _req: SweepRequest) -> Result<SweepBatchResult> {
let Some(days) = self.retention_days else {
return Ok(SweepBatchResult {
rows_deleted: 0,
has_more: false,
retention_days_applied: None,
});
};
let cutoff_ms = epoch_ms_now() - (days as i64) * 86_400_000;
let limit = self.retention.sweep_batch_size;
let mut tx = self.pool.begin().await?;
let result = sqlx::query!(
"DELETE FROM labels WHERE rowid IN (
SELECT rowid FROM labels WHERE created_at < ?1 LIMIT ?2
)",
cutoff_ms,
limit,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
let rows = result.rows_affected() as i64;
Ok(SweepBatchResult {
rows_deleted: rows,
has_more: rows >= limit,
retention_days_applied: Some(days),
})
}
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> {
use crate::report::ReportStatus;
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!: ReportStatus",
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 != ReportStatus::Pending {
return Err(Error::ReportAlreadyResolved { id: req.report_id });
}
let created_at = epoch_ms_now();
let label_event = if let Some(apply) = req.action.as_apply() {
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.action.as_apply().map(|a| a.val.clone());
let resolved_status = ReportStatus::Resolved;
sqlx::query!(
"UPDATE reports SET
status = ?1,
resolved_at = ?2,
resolved_by = ?3,
resolution_label = ?4,
resolution_reason = ?5
WHERE id = ?6",
resolved_status,
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.action.as_apply().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 = ReportStatus::Resolved;
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_retention_sweep_audit_reason(result: &SweepResult) -> String {
serde_json::json!({
"rows_deleted": result.rows_deleted,
"batches": result.batches,
"duration_ms": result.duration_ms,
"retention_days_applied": result.retention_days_applied,
})
.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);
}
}