systemprompt_analytics/snapshots/
jobs.rs1use super::ranges::Range;
7use super::{
8 FeedbackSnapshotsRepository, SnapshotJobLease, SnapshotJobState, SnapshotRangeJob,
9 SnapshotRangeRequest, invalid,
10};
11use chrono::{DateTime, Utc};
12use sqlx::PgTransaction;
13use systemprompt_identifiers::{AnalyticsSnapshotJobId, AnalyticsWorkerId, UserId};
14
15struct LeasedRange {
16 scope: String,
17 from_day: chrono::NaiveDate,
18 to_day: chrono::NaiveDate,
19 generation: i64,
20 fact_generation: i64,
21}
22
23impl FeedbackSnapshotsRepository {
24 pub async fn request_range(
25 &self,
26 owner: &UserId,
27 request: &SnapshotRangeRequest,
28 now: DateTime<Utc>,
29 ) -> crate::Result<SnapshotRangeJob> {
30 if request.operation_id.as_str().is_empty()
31 || request.operation_id.as_str().len() > 180
32 || request
33 .resource_id
34 .as_ref()
35 .is_some_and(|id| id.as_str().is_empty() || id.as_str().len() > 180)
36 {
37 return Err(invalid("Invalid range operation identity"));
38 }
39 let days = (request.to_day - request.from_day).num_days();
40 if !(1..=365).contains(&days)
41 || request.from_day < now.date_naive() - chrono::Duration::days(364)
42 || request.to_day > now.date_naive() + chrono::Duration::days(1)
43 {
44 return Err(invalid(
45 "Range must contain 1–365 UTC days within retention",
46 ));
47 }
48 let scope = request.resource_id.as_ref().map_or("", |id| id.as_str());
49 let mut tx = self.pool.begin().await?;
50 sqlx::query!(
51 "INSERT INTO analytics_snapshot_state(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
52 owner.as_str()
53 )
54 .execute(&mut *tx)
55 .await?;
56 sqlx::query!(
57 "SELECT generation FROM analytics_snapshot_state WHERE owner_id=$1 FOR UPDATE",
58 owner.as_str()
59 )
60 .fetch_one(&mut *tx)
61 .await?;
62 let existing=sqlx::query!("SELECT scope,from_day,to_day FROM analytics_snapshot_jobs WHERE owner_id=$1 AND job_id=$2",owner.as_str(),request.operation_id.as_str()).fetch_optional(&mut *tx).await?;
63 if let Some(row) = existing {
64 if row.scope != scope
65 || row.from_day != request.from_day
66 || row.to_day != request.to_day
67 {
68 return Err(invalid("Conflicting range operation retry"));
69 }
70 } else {
71 let count=sqlx::query_scalar!(r#"SELECT COUNT(*) AS "count!" FROM analytics_snapshot_jobs WHERE owner_id=$1 AND state IN('pending','leased')"#,owner.as_str()).fetch_one(&mut *tx).await?;
72 if count >= 100 {
73 return Err(invalid("Pending range job limit reached"));
74 }
75 sqlx::query!("INSERT INTO analytics_snapshot_jobs(owner_id,job_id,scope,from_day,to_day) VALUES($1,$2,$3,$4,$5)",owner.as_str(),request.operation_id.as_str(),scope,request.from_day,request.to_day).execute(&mut *tx).await?;
76 }
77 tx.commit().await?;
78 self.range_job(owner, &request.operation_id)
79 .await?
80 .ok_or_else(|| invalid("Range operation unavailable"))
81 }
82 pub async fn range_job(
83 &self,
84 owner: &UserId,
85 job: &AnalyticsSnapshotJobId,
86 ) -> crate::Result<Option<SnapshotRangeJob>> {
87 let row=sqlx::query!("SELECT state,result,last_error FROM analytics_snapshot_jobs WHERE owner_id=$1 AND job_id=$2",owner.as_str(),job.as_str()).fetch_optional(&self.pool).await?;
88 row.map(|row| {
89 Ok(SnapshotRangeJob {
90 operation_id: job.clone(),
91 state: SnapshotJobState::parse(&row.state)?,
92 result: row.result.map(serde_json::from_value).transpose()?,
93 diagnostic: row.last_error,
94 })
95 })
96 .transpose()
97 }
98 pub async fn claim_range(
99 &self,
100 owner: &UserId,
101 worker: &AnalyticsWorkerId,
102 ) -> crate::Result<Option<SnapshotJobLease>> {
103 let mut tx = self.pool.begin().await?;
104 let row=sqlx::query!("SELECT job_id FROM analytics_snapshot_jobs WHERE owner_id=$1 AND (state='pending' OR (state='leased' AND lease_until<=clock_timestamp())) ORDER BY created_at,job_id FOR UPDATE SKIP LOCKED LIMIT 1",owner.as_str()).fetch_optional(&mut *tx).await?;
105 let Some(row) = row else {
106 tx.commit().await?;
107 return Ok(None);
108 };
109 let epoch=sqlx::query_scalar!("UPDATE analytics_snapshot_jobs SET state='leased',lease_worker=$3,lease_epoch=lease_epoch+1,lease_until=clock_timestamp()+interval '5 minutes' WHERE owner_id=$1 AND job_id=$2 RETURNING lease_epoch",owner.as_str(),&row.job_id,worker.as_str()).fetch_one(&mut *tx).await?;
110 tx.commit().await?;
111 Ok(Some(SnapshotJobLease {
112 operation_id: AnalyticsSnapshotJobId::new(row.job_id),
113 worker_id: worker.clone(),
114 epoch,
115 }))
116 }
117 pub async fn complete_range(
118 &self,
119 owner: &UserId,
120 lease: &SnapshotJobLease,
121 now: DateTime<Utc>,
122 ) -> crate::Result<()> {
123 let mut tx = self.pool.begin().await?;
124 let leased = Self::lock_leased_range(&mut tx, owner, lease).await?;
125 let range = Range {
126 scope: &leased.scope,
127 from: leased.from_day,
128 to: leased.to_day,
129 generation: leased.generation,
130 fact_generation: leased.fact_generation,
131 now,
132 };
133 let result = match Self::assemble(&mut tx, owner, &range).await {
134 Ok(snapshot) => match serde_json::to_value(snapshot) {
135 Ok(result) => result,
136 Err(error) => {
137 drop(tx);
138 self.fail_range(owner, lease, &error.to_string()).await?;
139 return Err(error.into());
140 },
141 },
142 Err(error) => {
143 drop(tx);
144 self.fail_range(owner, lease, &error.to_string()).await?;
145 return Err(error);
146 },
147 };
148 let changed=sqlx::query!("UPDATE analytics_snapshot_jobs SET state='ready',result=$5,completed_at=clock_timestamp(),lease_worker=NULL,lease_until=NULL WHERE owner_id=$1 AND job_id=$2 AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>clock_timestamp()",owner.as_str(),lease.operation_id.as_str(),lease.worker_id.as_str(),lease.epoch,result).execute(&mut *tx).await?;
149 if changed.rows_affected() != 1 {
150 return Err(invalid("Expired range job lease"));
151 }
152 tx.commit().await?;
153 Ok(())
154 }
155 pub async fn fail_range(
158 &self,
159 owner: &UserId,
160 lease: &SnapshotJobLease,
161 diagnostic: &str,
162 ) -> crate::Result<()> {
163 let diagnostic = systemprompt_models::text::truncate_with_ellipsis(diagnostic, 2000);
164 let changed=sqlx::query!("UPDATE analytics_snapshot_jobs SET state='failed',last_error=$5,completed_at=clock_timestamp(),lease_worker=NULL,lease_until=NULL WHERE owner_id=$1 AND job_id=$2 AND state='leased' AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>clock_timestamp()",owner.as_str(),lease.operation_id.as_str(),lease.worker_id.as_str(),lease.epoch,diagnostic).execute(&self.pool).await?;
165 if changed.rows_affected() != 1 {
166 return Err(invalid("Expired range job lease"));
167 }
168 Ok(())
169 }
170 async fn lock_leased_range(
171 tx: &mut PgTransaction<'_>,
172 owner: &UserId,
173 lease: &SnapshotJobLease,
174 ) -> crate::Result<LeasedRange> {
175 let state=sqlx::query!("SELECT generation,fact_generation FROM analytics_snapshot_state WHERE owner_id=$1 FOR UPDATE",owner.as_str()).fetch_one(&mut **tx).await?;
176 let row=sqlx::query!("SELECT scope,from_day,to_day FROM analytics_snapshot_jobs WHERE owner_id=$1 AND job_id=$2 AND state='leased' AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>clock_timestamp() FOR UPDATE",owner.as_str(),lease.operation_id.as_str(),lease.worker_id.as_str(),lease.epoch).fetch_optional(&mut **tx).await?.ok_or_else(||invalid("Stale range job lease"))?;
177 Ok(LeasedRange {
178 scope: row.scope,
179 from_day: row.from_day,
180 to_day: row.to_day,
181 generation: state.generation,
182 fact_generation: state.fact_generation,
183 })
184 }
185}