Skip to main content

systemprompt_analytics/feedback/
backfill.rs

1//! Replayable backfill pages enqueue changes and advance checkpoints in one
2//! transaction.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use super::{BackfillProgress, FeedbackFactsRepository, validation};
8use crate::Result;
9use systemprompt_identifiers::{TaskId, UserId};
10use systemprompt_models::feedback::ContentDigest;
11
12impl FeedbackFactsRepository {
13    pub async fn begin_backfill(
14        &self,
15        owner: &UserId,
16        job: &TaskId,
17        source: &str,
18    ) -> Result<BackfillProgress> {
19        if source.is_empty() || source.len() > 128 || source.chars().any(char::is_control) {
20            return Err(validation::invalid());
21        }
22        let mut tx = self.pool.begin().await?;
23        sqlx::query!(
24            "INSERT INTO analytics_fact_checkpoints(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
25            owner.as_str()
26        )
27        .execute(&mut *tx)
28        .await?;
29        sqlx::query!(
30            "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
31            owner.as_str()
32        )
33        .fetch_one(&mut *tx)
34        .await?;
35        sqlx::query!("INSERT INTO analytics_fact_backfills(owner_id,job_id,source) VALUES($1,$2,$3) ON CONFLICT DO NOTHING", owner.as_str(), job.as_str(), source).execute(&mut *tx).await?;
36        tx.commit().await?;
37        let progress = self.backfill(owner, job).await?;
38        if progress.source != source {
39            return Err(validation::invalid());
40        }
41        Ok(progress)
42    }
43
44    pub async fn backfill(&self, owner: &UserId, job: &TaskId) -> Result<BackfillProgress> {
45        let row = sqlx::query!("SELECT source,cursor,generation,pages,facts,complete FROM analytics_fact_backfills WHERE owner_id=$1 AND job_id=$2", owner.as_str(), job.as_str()).fetch_optional(&self.pool).await?.ok_or_else(validation::invalid)?;
46        Ok(BackfillProgress {
47            job_id: job.clone(),
48            source: row.source,
49            cursor: row.cursor,
50            generation: row.generation,
51            pages: row.pages,
52            facts: row.facts,
53            complete: row.complete,
54        })
55    }
56
57    pub async fn append_backfill_page(
58        &self,
59        owner: &UserId,
60        job: &TaskId,
61        page: &super::BackfillPage,
62    ) -> Result<BackfillProgress> {
63        if page.changes.len() > 256
64            || page.next_cursor.len() > 1024
65            || page.expected_generation < 0
66            || (page.changes.is_empty() && !page.complete)
67        {
68            return Err(validation::invalid());
69        }
70        let digest = ContentDigest::of(&serde_json::to_vec(page)?);
71        let mut tx = self.pool.begin().await?;
72        sqlx::query!(
73            "INSERT INTO analytics_fact_checkpoints(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
74            owner.as_str()
75        )
76        .execute(&mut *tx)
77        .await?;
78        sqlx::query!(
79            "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
80            owner.as_str()
81        )
82        .fetch_one(&mut *tx)
83        .await?;
84
85        let progress = sqlx::query!("SELECT source,generation,complete FROM analytics_fact_backfills WHERE owner_id=$1 AND job_id=$2 FOR UPDATE", owner.as_str(), job.as_str()).fetch_optional(&mut *tx).await?.ok_or_else(validation::invalid)?;
86        let previous = sqlx::query_scalar!("SELECT digest FROM analytics_fact_backfill_pages WHERE owner_id=$1 AND job_id=$2 AND page_generation=$3", owner.as_str(), job.as_str(), page.expected_generation).fetch_optional(&mut *tx).await?;
87        if let Some(previous) = previous {
88            if previous != digest.as_str() {
89                return Err(validation::invalid());
90            }
91            tx.commit().await?;
92            return self.backfill(owner, job).await;
93        }
94        if progress.complete
95            || progress.generation != page.expected_generation
96            || page
97                .changes
98                .iter()
99                .any(|change| change.key.source != progress.source)
100        {
101            return Err(validation::invalid());
102        }
103        for change in &page.changes {
104            Self::submit_in(&mut tx, owner, change).await?;
105        }
106        let count = i64::try_from(page.changes.len()).map_err(|_error| validation::invalid())?;
107        sqlx::query!("INSERT INTO analytics_fact_backfill_pages(owner_id,job_id,page_generation,digest) VALUES($1,$2,$3,$4)", owner.as_str(), job.as_str(), page.expected_generation, digest.as_str()).execute(&mut *tx).await?;
108        sqlx::query!("UPDATE analytics_fact_backfills SET cursor=$3,generation=generation+1,pages=pages+1,facts=facts+$4,complete=$5,updated_at=clock_timestamp() WHERE owner_id=$1 AND job_id=$2", owner.as_str(), job.as_str(), &page.next_cursor, count, page.complete).execute(&mut *tx).await?;
109        tx.commit().await?;
110        self.backfill(owner, job).await
111    }
112}