systemprompt_analytics/snapshots/
retention.rs1use 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}