1use 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}