cairn_mod/moderation/cache.rs
1//! Subject-strike-state cache management (#55).
2//!
3//! The cache (`subject_strike_state` SQLite table) is populated by
4//! the recorder (#51) on every `recordAction` / `revokeAction`. It
5//! holds `(subject_did → current_strike_count, last_action_at,
6//! last_recompute_at)`.
7//!
8//! Two roles in v1.4:
9//!
10//! - **Write-through** (already shipped in #51): every action
11//! recorder + revoker writes the cache in the same transaction
12//! as the underlying `subject_actions` mutation, so the cache
13//! stays in sync with the action stream up to whatever decay has
14//! accumulated since the last touch.
15//! - **Lazy recompute-on-read** (this module): a v1.5+ consumer
16//! that wants O(1) "is this user past threshold right now?"
17//! without a full history walk can call
18//! [`get_or_recompute_strike_count`] — fresh cache returns
19//! immediately; stale cache triggers a recompute via the decay
20//! calculator (#50) and a write-back.
21//!
22//! v1.4 has zero in-tree consumers of [`get_or_recompute_strike_count`].
23//! The function ships ahead of v1.5's auto-action-threshold-check
24//! consumer; landing it now keeps cache invariants paired with
25//! their write path (#51) instead of split across releases.
26//!
27//! # Cache-bypass invariant on read endpoints
28//!
29//! `tools.cairn.admin.getSubjectStrikes` (#52) and
30//! `tools.cairn.public.getMyStrikeState` (#54) intentionally DO
31//! NOT consult the cache. They always recompute from
32//! `subject_actions` via `crate::server::strike_state::build_strike_state_view`
33//! (private to the server module). That invariant is preserved
34//! through this module — adding cache management does not change
35//! the read-side correctness story.
36//!
37//! # Freshness semantics
38//!
39//! The cache is fresh when `now - last_recompute_at <
40//! policy.cache_freshness_window_seconds`. Strict less-than:
41//! exact-boundary equality counts as stale. Clock-skew
42//! (`last_recompute_at` in the future relative to `now`) is also
43//! treated as stale — that state is either bug-induced or the
44//! result of an out-of-band clock jump, and the safe answer in
45//! both cases is "recompute."
46//!
47//! # Missing-cache stance
48//!
49//! [`get_or_recompute_strike_count`] returns
50//! [`Error::StrikeCacheMissing`] when the subject_did has no cache
51//! row. Same semantic as the read endpoints' `SubjectNotFound`:
52//! "no actions have ever been recorded against this subject." v1.5
53//! consumers that want optimistic "never-actioned == 0 strikes"
54//! semantics can wrap with `unwrap_or(0)` or branch on the typed
55//! error; the typed branching is friendlier than a magic-zero
56//! return.
57
58use std::time::{Duration, SystemTime, UNIX_EPOCH};
59
60use sqlx::{Pool, Sqlite};
61
62use crate::error::{Error, Result};
63use crate::moderation::decay::calculate_strike_state;
64use crate::moderation::policy::StrikePolicy;
65use crate::server::strike_state::load_action_history;
66
67/// Read-side projection of one `subject_strike_state` row. Stored
68/// timestamps are epoch-ms `i64`; this struct holds the
69/// [`SystemTime`] equivalents so the freshness check works in
70/// `Duration` arithmetic without callers re-doing the conversion.
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct SubjectStrikeStateCache {
73 /// Account DID this row attributes to (PK).
74 pub subject_did: String,
75 /// Cached active strike count. Authoritative-as-of
76 /// `last_recompute_at`; drifts with time as decay accumulates.
77 pub current_strike_count: u32,
78 /// Effective_at of the most recent action recorded against
79 /// this DID. `None` when no actions exist (which shouldn't
80 /// happen in practice — the recorder writes a row on every
81 /// action — but the schema column is nullable for forward
82 /// compatibility with future write paths).
83 pub last_action_at: Option<SystemTime>,
84 /// Wall-clock the `current_strike_count` was last derived. Drives
85 /// the freshness check.
86 pub last_recompute_at: SystemTime,
87}
88
89/// Whether the cache row is fresh enough to read without a
90/// recompute. Pure function: same input always yields the same
91/// output.
92///
93/// Returns `false` for:
94/// - `now - last_recompute_at >= freshness_window` (stale).
95/// - `last_recompute_at > now` (clock skew or corruption — safe
96/// default is "recompute").
97pub fn cache_is_fresh(
98 cache: &SubjectStrikeStateCache,
99 freshness_window: Duration,
100 now: SystemTime,
101) -> bool {
102 match now.duration_since(cache.last_recompute_at) {
103 Ok(elapsed) => elapsed < freshness_window,
104 // last_recompute_at is in the future relative to now —
105 // either clock skew or a corrupted row. Treat as stale.
106 Err(_) => false,
107 }
108}
109
110/// Load the cache row for a DID. Returns `Ok(None)` if no row
111/// exists.
112pub async fn load_cache(
113 pool: &Pool<Sqlite>,
114 subject_did: &str,
115) -> Result<Option<SubjectStrikeStateCache>> {
116 let row = sqlx::query!(
117 "SELECT subject_did, current_strike_count, last_action_at, last_recompute_at
118 FROM subject_strike_state WHERE subject_did = ?1",
119 subject_did,
120 )
121 .fetch_optional(pool)
122 .await?;
123 let Some(r) = row else {
124 return Ok(None);
125 };
126 let current_strike_count = u32::try_from(r.current_strike_count).map_err(|_| {
127 Error::Signing(format!(
128 "subject_strike_state.current_strike_count {} out of u32 range",
129 r.current_strike_count
130 ))
131 })?;
132 Ok(Some(SubjectStrikeStateCache {
133 subject_did: r.subject_did,
134 current_strike_count,
135 last_action_at: r.last_action_at.map(epoch_ms_to_systemtime),
136 last_recompute_at: epoch_ms_to_systemtime(r.last_recompute_at),
137 }))
138}
139
140/// UPSERT a recompute result into the cache. Bumps
141/// `current_strike_count` and `last_recompute_at`; leaves
142/// `last_action_at` untouched on the conflict-update path (this is
143/// a recompute, not a new action). On the insert path
144/// `last_action_at` is `NULL` — the conflict-update branch should
145/// be the only realistic path since the recorder writes the row on
146/// every action; a fresh-INSERT here would only happen if the
147/// caller invokes this on a never-actioned subject (which the
148/// public surface forbids via [`get_or_recompute_strike_count`]'s
149/// missing-cache check).
150pub async fn update_cache(
151 pool: &Pool<Sqlite>,
152 subject_did: &str,
153 current_count: u32,
154 now: SystemTime,
155) -> Result<()> {
156 let now_ms = systemtime_to_epoch_ms(now);
157 let count_i64 = current_count as i64;
158 sqlx::query!(
159 "INSERT INTO subject_strike_state (subject_did, current_strike_count, last_action_at, last_recompute_at)
160 VALUES (?1, ?2, NULL, ?3)
161 ON CONFLICT(subject_did) DO UPDATE SET
162 current_strike_count = excluded.current_strike_count,
163 last_recompute_at = excluded.last_recompute_at",
164 subject_did,
165 count_i64,
166 now_ms,
167 )
168 .execute(pool)
169 .await?;
170 Ok(())
171}
172
173/// Return the subject's currently-active strike count, recomputing
174/// from `subject_actions` history if the cache is stale.
175///
176/// Currently UNUSED by any in-tree code. Ships ahead of v1.5's
177/// auto-action-threshold-check consumer (per the brief that
178/// settled this seam); landing the cache management alongside the
179/// cache write path it pairs with — instead of split across
180/// releases — keeps the invariants reviewable in one place.
181///
182/// Errors:
183/// - [`Error::StrikeCacheMissing`] when no cache row exists for
184/// `subject_did`. The recorder writes the row on every action,
185/// so this means the subject has never been actioned.
186/// - DB / arithmetic / decay-calculator errors propagate.
187///
188/// Cache-write best-effort: when a recompute fires, the new count
189/// is written back via [`update_cache`]. If that write fails
190/// (e.g., a transient DB error), the recomputed count is still
191/// returned to the caller — the cache update is an optimization,
192/// not a correctness gate. The failure is logged via `tracing::warn`.
193pub async fn get_or_recompute_strike_count(
194 pool: &Pool<Sqlite>,
195 subject_did: &str,
196 policy: &StrikePolicy,
197 now: SystemTime,
198) -> Result<u32> {
199 let Some(cache) = load_cache(pool, subject_did).await? else {
200 return Err(Error::StrikeCacheMissing(subject_did.to_string()));
201 };
202
203 let freshness_window = Duration::from_secs(policy.cache_freshness_window_seconds as u64);
204 if cache_is_fresh(&cache, freshness_window, now) {
205 return Ok(cache.current_strike_count);
206 }
207
208 // Stale: recompute from source-of-truth via the shared loader
209 // + decay calculator.
210 let history = load_action_history(pool, subject_did).await?;
211 let state = calculate_strike_state(&history, policy, now);
212 let new_count = state.current_count;
213
214 // Best-effort cache update. A failure here doesn't lose the
215 // recomputed count — the caller gets correct data even if the
216 // cache stays stale for another freshness_window.
217 if let Err(e) = update_cache(pool, subject_did, new_count, now).await {
218 tracing::warn!(
219 subject = %subject_did,
220 error = %e,
221 "strike-state cache update failed; returning recomputed count anyway",
222 );
223 }
224
225 Ok(new_count)
226}
227
228fn epoch_ms_to_systemtime(ms: i64) -> SystemTime {
229 if ms >= 0 {
230 UNIX_EPOCH + Duration::from_millis(ms as u64)
231 } else {
232 UNIX_EPOCH
233 }
234}
235
236fn systemtime_to_epoch_ms(t: SystemTime) -> i64 {
237 t.duration_since(UNIX_EPOCH)
238 .unwrap_or(Duration::ZERO)
239 .as_millis()
240 .try_into()
241 .unwrap_or(i64::MAX)
242}
243
244#[cfg(test)]
245mod tests {
246 use super::*;
247
248 fn t0() -> SystemTime {
249 UNIX_EPOCH + Duration::from_secs(2_000_000_000)
250 }
251
252 fn cache_at(last_recompute: SystemTime) -> SubjectStrikeStateCache {
253 SubjectStrikeStateCache {
254 subject_did: "did:plc:test".to_string(),
255 current_strike_count: 0,
256 last_action_at: None,
257 last_recompute_at: last_recompute,
258 }
259 }
260
261 // ---------- cache_is_fresh ----------
262
263 #[test]
264 fn fresh_within_window() {
265 let now = t0();
266 let c = cache_at(now - Duration::from_secs(1800));
267 assert!(cache_is_fresh(&c, Duration::from_secs(3600), now));
268 }
269
270 #[test]
271 fn stale_past_window() {
272 let now = t0();
273 let c = cache_at(now - Duration::from_secs(7200));
274 assert!(!cache_is_fresh(&c, Duration::from_secs(3600), now));
275 }
276
277 #[test]
278 fn exact_boundary_is_stale_strict_less_than() {
279 // freshness_window = 3600s, last_recompute exactly 3600s
280 // ago → elapsed == window → NOT fresh (the check is strict
281 // less-than).
282 let now = t0();
283 let c = cache_at(now - Duration::from_secs(3600));
284 assert!(!cache_is_fresh(&c, Duration::from_secs(3600), now));
285 }
286
287 #[test]
288 fn last_recompute_in_future_is_stale() {
289 // Clock skew: last_recompute_at is AFTER now. Treat as stale
290 // — the safe answer in either bug-induced or clock-jump
291 // scenarios is to recompute.
292 let now = t0();
293 let c = cache_at(now + Duration::from_secs(60));
294 assert!(!cache_is_fresh(&c, Duration::from_secs(3600), now));
295 }
296
297 #[test]
298 fn zero_freshness_window_makes_everything_stale() {
299 // Defensive: a zero-window policy (which validation
300 // forbids, but defense-in-depth) makes nothing fresh.
301 let now = t0();
302 let c = cache_at(now);
303 assert!(!cache_is_fresh(&c, Duration::ZERO, now));
304 }
305
306 #[test]
307 fn epoch_ms_systemtime_round_trip() {
308 // Pin the conversion contract used by load_cache /
309 // update_cache. Pick a value that doesn't sit on a second
310 // boundary so the ms precision matters.
311 let original_ms: i64 = 1_776_902_400_123;
312 let st = epoch_ms_to_systemtime(original_ms);
313 let back = systemtime_to_epoch_ms(st);
314 assert_eq!(back, original_ms);
315 }
316
317 // The DB-touching paths (load_cache / update_cache /
318 // get_or_recompute_strike_count) are exercised in
319 // tests/strike_cache.rs against a real SQLite database — the
320 // unit tests stay pure to avoid sqlx-prepare churn for what
321 // are otherwise straightforward UPSERT statements.
322}