Skip to main content

systemprompt_analytics/snapshots/
retention.rs

1//! Retention erases evidence only behind committed producer, fact and snapshot
2//! barriers.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use super::{FeedbackSnapshotsRepository, RetentionOutcome, invalid};
8use chrono::{DateTime, Utc};
9use systemprompt_identifiers::UserId;
10
11impl FeedbackSnapshotsRepository {
12    pub async fn compact_in(
13        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
14        owner: &UserId,
15        now: DateTime<Utc>,
16    ) -> crate::Result<RetentionOutcome> {
17        if now > Utc::now() {
18            return Err(invalid("Retention clock cannot be in the future"));
19        }
20        Self::lock_retention_barriers(tx, owner).await?;
21        let cutoff = now - chrono::Duration::days(90);
22        let day = cutoff.date_naive();
23        let oldest = now.date_naive() - chrono::Duration::days(364);
24        let previous = sqlx::query_scalar!(
25            "SELECT evidence_cutoff FROM analytics_snapshot_state WHERE owner_id=$1",
26            owner.as_str()
27        )
28        .fetch_one(&mut **tx)
29        .await?;
30        if previous.is_some_and(|previous| cutoff < previous) {
31            return Err(invalid("Retention cutoff cannot move backwards"));
32        }
33        let outcome = Self::purge_expired(tx, owner, cutoff, oldest).await?;
34        let days: Vec<_> = (0..90)
35            .map(|offset| now.date_naive() - chrono::Duration::days(offset))
36            .collect();
37        let generation = sqlx::query_scalar!(
38            "SELECT fact_generation FROM analytics_snapshot_state WHERE owner_id=$1",
39            owner.as_str()
40        )
41        .fetch_one(&mut **tx)
42        .await?;
43        Self::rebuild_days(tx, owner, &days, generation).await?;
44        Self::refresh_in(tx, owner, &[], now).await?;
45        sqlx::query!(
46            "SELECT public.finish_reporting_compaction($1) AS processed",
47            now - chrono::Duration::days(90)
48        )
49        .fetch_one(&mut **tx)
50        .await?;
51        Ok(RetentionOutcome {
52            compacted_before: day,
53            ..outcome
54        })
55    }
56    async fn purge_expired(
57        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
58        owner: &UserId,
59        cutoff: DateTime<Utc>,
60        oldest: chrono::NaiveDate,
61    ) -> crate::Result<RetentionOutcome> {
62        let day = cutoff.date_naive();
63        sqlx::query!("INSERT INTO analytics_snapshot_dirty(owner_id,scope) SELECT owner_id,scope FROM analytics_feedback_snapshots WHERE owner_id=$1 ON CONFLICT DO NOTHING",owner.as_str()).execute(&mut **tx).await?;
64        sqlx::query!("UPDATE analytics_snapshot_daily SET metrics='{}'::jsonb,spend='{}'::jsonb,histogram='{}'::jsonb,cohort=0,suppressed=true WHERE owner_id=$1 AND ((day<$2 AND cohort<5) OR day=$2)",owner.as_str(),day).execute(&mut **tx).await?;
65        sqlx::query!(
66            "DELETE FROM analytics_snapshot_identities WHERE owner_id=$1 AND day<=$2",
67            owner.as_str(),
68            day
69        )
70        .execute(&mut **tx)
71        .await?;
72        let removed = Self::erase_expired(tx, owner, cutoff).await?;
73        sqlx::query!(
74            "DELETE FROM analytics_fact_changes WHERE owner_id=$1 AND occurred_at<$2",
75            owner.as_str(),
76            cutoff
77        )
78        .execute(&mut **tx)
79        .await?;
80        sqlx::query!("DELETE FROM analytics_fact_deltas WHERE owner_id=$1 AND (occurred_at<$2 OR (before_fact->'value'->>'occurred_at')::timestamptz<$2 OR (after_fact->'value'->>'occurred_at')::timestamptz<$2)",owner.as_str(),cutoff).execute(&mut **tx).await?;
81        sqlx::query!(
82            "DELETE FROM analytics_fact_backfill_pages WHERE owner_id=$1",
83            owner.as_str()
84        )
85        .execute(&mut **tx)
86        .await?;
87        sqlx::query!(
88            "DELETE FROM analytics_fact_backfills WHERE owner_id=$1",
89            owner.as_str()
90        )
91        .execute(&mut **tx)
92        .await?;
93        let removed_daily = sqlx::query!(
94            "DELETE FROM analytics_snapshot_daily WHERE owner_id=$1 AND day<$2",
95            owner.as_str(),
96            oldest
97        )
98        .execute(&mut **tx)
99        .await?
100        .rows_affected();
101        sqlx::query!(
102            "DELETE FROM analytics_snapshot_jobs WHERE owner_id=$1 AND created_at<$2",
103            owner.as_str(),
104            cutoff
105        )
106        .execute(&mut **tx)
107        .await?;
108        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 privacy retention' WHERE owner_id=$1",owner.as_str()).execute(&mut **tx).await?;
109        sqlx::query!("UPDATE analytics_snapshot_state SET compacted_before=$2,evidence_cutoff=$3 WHERE owner_id=$1",owner.as_str(),day,cutoff).execute(&mut **tx).await?;
110        Ok(RetentionOutcome {
111            compacted_before: day,
112            removed_facts: removed,
113            removed_daily,
114        })
115    }
116    async fn lock_retention_barriers(
117        tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
118        owner: &UserId,
119    ) -> crate::Result<()> {
120        sqlx::query!("SELECT public.prepare_reporting_privacy() AS locked")
121            .fetch_one(&mut **tx)
122            .await?;
123        sqlx::query!("LOCK TABLE analytics_ingestion_producers,analytics_fact_checkpoints,analytics_fact_backfills,analytics_fact_changes,analytics_fact_consumers,analytics_fact_deltas IN EXCLUSIVE MODE").execute(&mut **tx).await?;
124        let producers=sqlx::query!("SELECT producer,pending_count FROM analytics_ingestion_producers ORDER BY producer FOR UPDATE").fetch_all(&mut **tx).await?;
125        if producers.iter().any(|row| row.pending_count != 0) {
126            return Err(invalid(
127                "Retention waits for pending producer evidence and privacy corrections",
128            ));
129        }
130        let facts = sqlx::query_scalar!(
131            "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
132            owner.as_str()
133        )
134        .fetch_optional(&mut **tx)
135        .await?
136        .ok_or_else(|| invalid("Fact checkpoint is not initialized"))?;
137        let state = sqlx::query!(
138            "SELECT fact_generation FROM analytics_snapshot_state WHERE owner_id=$1 FOR UPDATE",
139            owner.as_str()
140        )
141        .fetch_optional(&mut **tx)
142        .await?
143        .ok_or_else(|| invalid("Snapshots are not initialized"))?;
144        let consumers=sqlx::query!("SELECT consumer,generation FROM analytics_fact_consumers WHERE owner_id=$1 ORDER BY consumer FOR UPDATE",owner.as_str()).fetch_all(&mut **tx).await?;
145        let pending=sqlx::query_scalar!(r#"SELECT EXISTS(SELECT 1 FROM analytics_fact_changes WHERE owner_id=$1 AND state IN('pending','leased')) OR EXISTS(SELECT 1 FROM analytics_fact_deltas WHERE owner_id=$1 AND consumed_at IS NULL) OR EXISTS(SELECT 1 FROM analytics_fact_backfills WHERE owner_id=$1 AND NOT complete) AS "pending!""#,owner.as_str()).fetch_one(&mut **tx).await?;
146        if pending
147            || state.fact_generation != facts
148            || !consumers.iter().any(|row| row.consumer == "snapshots-v1")
149            || consumers.iter().any(|row| row.generation != facts)
150        {
151            return Err(invalid(
152                "Retention waits for committed fact and aggregate checkpoints",
153            ));
154        }
155        Ok(())
156    }
157}