Skip to main content

systemprompt_analytics/snapshots/
reads.rs

1//! Dashboard reads use retained snapshots; workers refresh bounded standard
2//! windows.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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}