Skip to main content

systemprompt_analytics/snapshots/
processing.rs

1//! Fenced delta processing updates shadows, daily totals and durable consumer
2//! checkpoints atomically.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use super::daily::FactReference;
8use super::{FeedbackSnapshotsRepository, invalid};
9use crate::feedback::{DeltaClaim, DeltaLease, FeedbackFactsRepository};
10use chrono::{DateTime, NaiveDate, Utc};
11use std::collections::BTreeSet;
12use systemprompt_identifiers::{AnalyticsWorkerId, UserId};
13
14#[derive(Debug, Clone, Copy)]
15struct ShadowWindow {
16    cutoff: NaiveDate,
17    generation: i64,
18}
19
20impl FeedbackSnapshotsRepository {
21    pub async fn process(
22        &self,
23        owner: &UserId,
24        worker: &AnalyticsWorkerId,
25        now: DateTime<Utc>,
26    ) -> crate::Result<u64> {
27        let Some(lease) = self
28            .facts
29            .claim_deltas(
30                owner,
31                "snapshots-v1",
32                worker,
33                DeltaClaim {
34                    limit: 256,
35                    lease_seconds: 300,
36                },
37            )
38            .await?
39        else {
40            return Ok(0);
41        };
42        let result = self.apply_batch(owner, &lease, now).await;
43        if result.is_err() {
44            sqlx::query!("INSERT INTO analytics_snapshot_state(owner_id,last_error) VALUES($1,'Snapshot batch failed; leased work remains retryable') ON CONFLICT(owner_id) DO UPDATE SET last_error=EXCLUDED.last_error",owner.as_str()).execute(&self.pool).await?;
45        }
46        result
47    }
48    pub async fn apply_batch(
49        &self,
50        owner: &UserId,
51        lease: &DeltaLease,
52        now: DateTime<Utc>,
53    ) -> crate::Result<u64> {
54        if lease.consumer != "snapshots-v1" {
55            return Err(invalid("Invalid snapshot consumer"));
56        }
57        let mut tx = self.pool.begin().await?;
58        sqlx::query!(
59            "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
60            owner.as_str()
61        )
62        .fetch_one(&mut *tx)
63        .await?;
64        sqlx::query!(
65            "INSERT INTO analytics_snapshot_state(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
66            owner.as_str()
67        )
68        .execute(&mut *tx)
69        .await?;
70        let state = sqlx::query!(
71            "SELECT compacted_before FROM analytics_snapshot_state WHERE owner_id=$1 FOR UPDATE",
72            owner.as_str()
73        )
74        .fetch_one(&mut *tx)
75        .await?;
76        FeedbackFactsRepository::lock_delta_lease(&mut tx, owner, lease).await?;
77        let rows=sqlx::query!(r#"SELECT generation,fact_kind,source,fact_id,before_fact,after_fact,occurred_at,(SELECT revision FROM analytics_normalized_facts f WHERE f.owner_id=d.owner_id AND f.fact_kind=d.fact_kind AND f.source=d.source AND f.fact_id=d.fact_id) AS "revision?" FROM analytics_fact_deltas d WHERE owner_id=$1 AND generation>$2 AND generation<=$3 ORDER BY generation LIMIT 256"#,owner.as_str(),lease.after_generation,lease.through_generation).fetch_all(&mut *tx).await?;
78        let count =
79            u64::try_from(rows.len()).map_err(|_error| invalid("Snapshot batch overflow"))?;
80        let cutoff = state
81            .compacted_before
82            .unwrap_or_else(|| now.date_naive() - chrono::Duration::days(365));
83        let mut days = BTreeSet::new();
84        for row in rows {
85            if state.compacted_before.is_some()
86                && row.before_fact.is_none()
87                && row.revision.is_some_and(|revision| revision > 1)
88            {
89                Self::suppress_sealed_history(&mut tx, owner, cutoff).await?;
90            }
91            for fact in [&row.before_fact, &row.after_fact] {
92                let reference = FactReference {
93                    kind: &row.fact_kind,
94                    source: &row.source,
95                    id: &row.fact_id,
96                    fact: fact.as_ref(),
97                };
98                days.extend(Self::affected_days(&mut tx, owner, &reference).await?);
99            }
100            days.insert(row.occurred_at.date_naive());
101            if row.occurred_at.date_naive() < cutoff {
102                Self::suppress_day(&mut tx, owner, row.occurred_at.date_naive(), row.generation)
103                    .await?;
104            }
105            let reference = FactReference {
106                kind: &row.fact_kind,
107                source: &row.source,
108                id: &row.fact_id,
109                fact: row.after_fact.as_ref(),
110            };
111            let shadow = ShadowWindow {
112                cutoff,
113                generation: row.generation,
114            };
115            if let Some(day) = Self::replace_shadow(&mut tx, owner, &reference, shadow).await? {
116                days.insert(day);
117            }
118            days.extend(Self::affected_days(&mut tx, owner, &reference).await?);
119        }
120        Self::finish_batch(&mut tx, owner, lease, &days, cutoff).await?;
121        tx.commit().await?;
122        Ok(count)
123    }
124
125    async fn suppress_sealed_history(
126        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
127        owner: &UserId,
128        cutoff: NaiveDate,
129    ) -> crate::Result<()> {
130        sqlx::query!("INSERT INTO analytics_snapshot_dirty(owner_id,scope) SELECT owner_id,scope FROM analytics_snapshot_daily WHERE owner_id=$1 AND day<$2 ON CONFLICT DO NOTHING",owner.as_str(),cutoff).execute(&mut **tx).await?;
131        sqlx::query!("UPDATE analytics_snapshot_daily SET metrics='{}'::jsonb,spend='{}'::jsonb,histogram='{}'::jsonb,cohort=0,suppressed=true WHERE owner_id=$1 AND day<$2",owner.as_str(),cutoff).execute(&mut **tx).await?;
132        sqlx::query!("UPDATE analytics_snapshot_jobs SET state='pending',result=NULL,lease_worker=NULL,lease_until=NULL,lease_epoch=lease_epoch+1,last_error='Sealed history suppressed after a correction without retained provenance' WHERE owner_id=$1 AND from_day<$2",owner.as_str(),cutoff).execute(&mut **tx).await?;
133        Ok(())
134    }
135
136    async fn suppress_day(
137        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
138        owner: &UserId,
139        day: NaiveDate,
140        generation: i64,
141    ) -> crate::Result<()> {
142        sqlx::query!("INSERT INTO analytics_snapshot_dirty(owner_id,scope) SELECT owner_id,scope FROM analytics_snapshot_daily WHERE owner_id=$1 AND day=$2 ON CONFLICT DO NOTHING",owner.as_str(),day).execute(&mut **tx).await?;
143        sqlx::query!("UPDATE analytics_snapshot_daily SET metrics='{}'::jsonb,spend='{}'::jsonb,histogram='{}'::jsonb,cohort=0,suppressed=true,generation=$3 WHERE owner_id=$1 AND day=$2",owner.as_str(),day,generation).execute(&mut **tx).await?;
144        Ok(())
145    }
146
147    async fn replace_shadow(
148        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
149        owner: &UserId,
150        reference: &FactReference<'_>,
151        shadow: ShadowWindow,
152    ) -> crate::Result<Option<NaiveDate>> {
153        sqlx::query!("DELETE FROM analytics_snapshot_shadow WHERE owner_id=$1 AND fact_kind=$2 AND source=$3 AND fact_id=$4",owner.as_str(),reference.kind,reference.source,reference.id).execute(&mut **tx).await?;
154        let Some(fact) = reference.fact else {
155            return Ok(None);
156        };
157        let occurred = fact
158            .get("value")
159            .and_then(|value| value.get("occurred_at"))
160            .cloned()
161            .ok_or_else(|| invalid("Missing event timestamp"))?;
162        let occurred: DateTime<Utc> = serde_json::from_value(occurred)?;
163        if occurred.date_naive() >= shadow.cutoff {
164            sqlx::query!("INSERT INTO analytics_snapshot_shadow(owner_id,fact_kind,source,fact_id,fact,occurred_at,generation) VALUES($1,$2,$3,$4,$5,$6,$7)",owner.as_str(),reference.kind,reference.source,reference.id,fact,occurred,shadow.generation).execute(&mut **tx).await?;
165        }
166        Ok(Some(occurred.date_naive()))
167    }
168
169    async fn finish_batch(
170        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
171        owner: &UserId,
172        lease: &DeltaLease,
173        days: &BTreeSet<NaiveDate>,
174        cutoff: NaiveDate,
175    ) -> crate::Result<()> {
176        let all_days: Vec<_> = days.iter().copied().collect();
177        sqlx::query!("UPDATE analytics_snapshot_jobs SET state='pending',result=NULL,lease_worker=NULL,lease_until=NULL,lease_epoch=lease_epoch+1,last_error='Range invalidated by corrected evidence' WHERE owner_id=$1 AND EXISTS(SELECT 1 FROM unnest($2::date[]) d WHERE d>=from_day AND d<to_day)",owner.as_str(),&all_days).execute(&mut **tx).await?;
178        let days: Vec<_> = days.iter().copied().filter(|day| *day >= cutoff).collect();
179        Self::rebuild_days(tx, owner, &days, lease.through_generation).await?;
180        FeedbackFactsRepository::complete_delta_batch(tx, owner, lease).await?;
181        sqlx::query!("UPDATE analytics_snapshot_state SET fact_generation=$2,last_error=NULL WHERE owner_id=$1",owner.as_str(),lease.through_generation).execute(&mut **tx).await?;
182        Ok(())
183    }
184}