systemprompt_analytics/snapshots/
reads.rs1use super::ranges::Range;
8use super::{FeedbackSnapshot, FeedbackSnapshotsRepository, SnapshotHealth, invalid};
9use chrono::{DateTime, Utc};
10use systemprompt_identifiers::{ManagedResourceId, UserId};
11
12impl FeedbackSnapshotsRepository {
13 pub async fn refresh(
14 &self,
15 owner: &UserId,
16 resources: &[ManagedResourceId],
17 now: DateTime<Utc>,
18 ) -> crate::Result<i64> {
19 if resources.len() > 10000 {
20 return Err(invalid("Too many snapshot inventory scopes"));
21 }
22 let mut tx = self.pool.begin().await?;
23 let generation = Self::refresh_in(&mut tx, owner, resources, now).await?;
24 tx.commit().await?;
25 Ok(generation)
26 }
27 pub(super) async fn refresh_in(
28 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
29 owner: &UserId,
30 resources: &[ManagedResourceId],
31 now: DateTime<Utc>,
32 ) -> crate::Result<i64> {
33 sqlx::query!(
34 "INSERT INTO analytics_snapshot_state(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
35 owner.as_str()
36 )
37 .execute(&mut **tx)
38 .await?;
39 let state=sqlx::query!("SELECT generation,fact_generation,day FROM analytics_snapshot_state WHERE owner_id=$1 FOR UPDATE",owner.as_str()).fetch_one(&mut **tx).await?;
40 let generation = state
41 .generation
42 .checked_add(1)
43 .ok_or_else(|| invalid("Snapshot generation overflow"))?;
44 let existing:std::collections::BTreeSet<String>=sqlx::query_scalar!("SELECT DISTINCT scope FROM analytics_feedback_snapshots WHERE owner_id=$1 ORDER BY scope LIMIT 10002",owner.as_str()).fetch_all(&mut **tx).await?.into_iter().collect();
45 let mut scopes: std::collections::BTreeSet<String> = resources
46 .iter()
47 .map(|id| id.as_str().to_owned())
48 .filter(|scope| !existing.contains(scope))
49 .collect();
50 if state.day != Some(now.date_naive()) {
51 scopes.extend(existing);
52 scopes.insert(String::new());
53 }
54 scopes.extend(sqlx::query_scalar!("SELECT scope FROM analytics_snapshot_dirty WHERE owner_id=$1 ORDER BY scope LIMIT 10002",owner.as_str()).fetch_all(&mut **tx).await?);
55 if scopes.is_empty() {
56 return Ok(state.generation);
57 }
58 if scopes.len() > 10001 {
59 return Err(invalid("Snapshot scope count exceeds limit"));
60 }
61 let to = now
62 .date_naive()
63 .succ_opt()
64 .ok_or_else(|| invalid("Invalid snapshot date"))?;
65 for scope in scopes {
66 for window in [1i32, 7, 30, 90, 365] {
67 let from = to - chrono::Duration::days(i64::from(window));
68 let snapshot = Self::assemble(
69 tx,
70 owner,
71 &Range {
72 scope: &scope,
73 from,
74 to,
75 generation,
76 fact_generation: state.fact_generation,
77 now,
78 },
79 )
80 .await?;
81 let body = serde_json::to_value(snapshot)?;
82 sqlx::query!("INSERT INTO analytics_feedback_snapshots(owner_id,scope,window_days,generation,from_day,to_day,body) VALUES($1,$2,$3,$4,$5,$6,$7) ON CONFLICT(owner_id,scope,window_days) DO UPDATE SET generation=EXCLUDED.generation,from_day=EXCLUDED.from_day,to_day=EXCLUDED.to_day,body=EXCLUDED.body",owner.as_str(),&scope,window,generation,from,to,body).execute(&mut **tx).await?;
83 }
84 }
85 sqlx::query!(
86 "DELETE FROM analytics_snapshot_dirty WHERE owner_id=$1",
87 owner.as_str()
88 )
89 .execute(&mut **tx)
90 .await?;
91 sqlx::query!("UPDATE analytics_snapshot_state SET generation=$2,generated_at=$3,day=$4,retained_from=$5,last_error=NULL WHERE owner_id=$1",owner.as_str(),generation,now,now.date_naive(),to-chrono::Duration::days(365)).execute(&mut **tx).await?;
92 sqlx::query!(
93 "SELECT pg_notify('feedback_snapshots',$1)",
94 format!("{owner}:{generation}")
95 )
96 .execute(&mut **tx)
97 .await?;
98 Ok(generation)
99 }
100 pub async fn snapshot(
101 &self,
102 owner: &UserId,
103 resource: Option<&ManagedResourceId>,
104 window: u32,
105 ) -> crate::Result<Option<FeedbackSnapshot>> {
106 if ![1, 7, 30, 90, 365].contains(&window) {
107 return Err(invalid("Unsupported snapshot window"));
108 }
109 let scope = resource.map_or("", ManagedResourceId::as_str);
110 let window = i32::try_from(window).map_err(|_error| invalid("Invalid snapshot window"))?;
111 let value=sqlx::query_scalar!("SELECT body FROM analytics_feedback_snapshots WHERE owner_id=$1 AND scope=$2 AND window_days=$3",owner.as_str(),scope,window).fetch_optional(&self.pool).await?;
112 value
113 .map(serde_json::from_value)
114 .transpose()
115 .map_err(Into::into)
116 }
117 pub async fn snapshots(
118 &self,
119 owner: &UserId,
120 after: Option<&ManagedResourceId>,
121 window: u32,
122 limit: u32,
123 ) -> crate::Result<Vec<FeedbackSnapshot>> {
124 if ![1, 7, 30, 90, 365].contains(&window) || !(1..=100).contains(&limit) {
125 return Err(invalid("Invalid snapshot page"));
126 }
127 let after = after.map_or("", ManagedResourceId::as_str);
128 let window = i32::try_from(window).map_err(|_error| invalid("Invalid snapshot window"))?;
129 let limit = i64::from(limit);
130 let rows=sqlx::query_scalar!("SELECT body FROM analytics_feedback_snapshots WHERE owner_id=$1 AND scope>$2 AND window_days=$3 ORDER BY scope LIMIT $4",owner.as_str(),after,window,limit).fetch_all(&self.pool).await?;
131 rows.into_iter()
132 .map(|value| serde_json::from_value(value).map_err(Into::into))
133 .collect()
134 }
135 pub async fn health(&self, owner: &UserId) -> crate::Result<SnapshotHealth> {
136 let row=sqlx::query!(r#"SELECT s.generation,s.fact_generation,s.generated_at,s.last_error,
137 (SELECT COUNT(*) FROM analytics_fact_changes WHERE owner_id=$1 AND state IN('pending','leased')) AS "pending!",
138 (SELECT COALESCE(SUM(pending_count),0)::bigint FROM analytics_ingestion_producers) AS "producers!",
139 COALESCE((SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1),0) AS "facts!",
140 (SELECT COUNT(*) FROM analytics_snapshot_jobs WHERE owner_id=$1 AND state IN('pending','leased')) AS "jobs!"
141 FROM analytics_snapshot_state s WHERE s.owner_id=$1"#,owner.as_str()).fetch_optional(&self.pool).await?;
142 Ok(row.map_or_else(
143 || SnapshotHealth {
144 generation: 0,
145 fact_generation: 0,
146 generated_at: None,
147 last_error: Some("Snapshots have not been initialized".to_owned()),
148 pending_changes: 0,
149 pending_producer_changes: 0,
150 facts_generation: 0,
151 pending_jobs: 0,
152 },
153 |row| SnapshotHealth {
154 generation: row.generation,
155 fact_generation: row.fact_generation,
156 generated_at: row.generated_at,
157 last_error: row.last_error,
158 pending_changes: row.pending,
159 pending_producer_changes: row.producers,
160 facts_generation: row.facts,
161 pending_jobs: row.jobs,
162 },
163 ))
164 }
165}