Skip to main content

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}