systemprompt_analytics/feedback/
backfill.rs1use 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}