systemprompt_analytics/feedback/
queue.rs1use super::{FactLease, FactsHealth, FeedbackFactsRepository, validation};
8use crate::Result;
9use systemprompt_identifiers::{AnalyticsChangeId, AnalyticsWorkerId, UserId};
10
11impl FeedbackFactsRepository {
12 pub async fn claim(
13 &self,
14 owner: &UserId,
15 worker: &AnalyticsWorkerId,
16 limit: u32,
17 lease_seconds: u32,
18 ) -> Result<Vec<FactLease>> {
19 if !(1..=256).contains(&limit) || !(1..=300).contains(&lease_seconds) {
20 return Err(validation::invalid());
21 }
22 let limit = i64::from(limit);
23 let seconds = f64::from(lease_seconds);
24 let rows = sqlx::query!(r#"WITH ready AS (SELECT change_id FROM analytics_fact_changes WHERE owner_id=$1 AND next_attempt_at<=now() AND (state='pending' OR (state='leased' AND lease_until<=now())) ORDER BY recorded_at,change_id LIMIT $3 FOR UPDATE SKIP LOCKED) UPDATE analytics_fact_changes c SET state='leased',lease_worker=$2,lease_epoch=c.lease_epoch+1,lease_until=now()+make_interval(secs=>$4),attempts=c.attempts+1 FROM ready WHERE c.owner_id=$1 AND c.change_id=ready.change_id RETURNING c.change_id,c.lease_epoch,c.lease_until AS "lease_until!""#, owner.as_str(), worker.as_str(), limit, seconds).fetch_all(&self.pool).await?;
25 Ok(rows
26 .into_iter()
27 .map(|row| FactLease {
28 change_id: AnalyticsChangeId::new(row.change_id),
29 worker_id: worker.clone(),
30 epoch: row.lease_epoch,
31 expires_at: row.lease_until,
32 })
33 .collect())
34 }
35
36 pub async fn retry(&self, owner: &UserId, lease: &FactLease) -> Result<()> {
37 let result = sqlx::query!("UPDATE analytics_fact_changes SET state='pending',lease_worker=NULL,lease_until=NULL,last_error='Fact processing failed; retry scheduled',next_attempt_at=now()+make_interval(secs=>LEAST(attempts,60)::double precision) WHERE owner_id=$1 AND change_id=$2 AND state='leased' AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>now()", owner.as_str(), lease.change_id.as_str(), lease.worker_id.as_str(), lease.epoch).execute(&self.pool).await?;
38 if result.rows_affected() != 1 {
39 return Err(validation::invalid());
40 }
41 Ok(())
42 }
43
44 pub async fn health(&self, owner: &UserId) -> Result<FactsHealth> {
45 let checkpoint = sqlx::query!("SELECT generation,last_applied_at,last_recorded_at FROM analytics_fact_checkpoints WHERE owner_id=$1", owner.as_str()).fetch_optional(&self.pool).await?;
46 let queue = sqlx::query!(r#"SELECT COUNT(*) FILTER(WHERE state='pending') AS "pending!", COUNT(*) FILTER(WHERE state='leased') AS "leased!", COUNT(*) FILTER(WHERE last_error IS NOT NULL AND state IN ('pending','leased')) AS "retries!", MIN(recorded_at) FILTER(WHERE state IN ('pending','leased')) AS oldest_pending_at FROM analytics_fact_changes WHERE owner_id=$1"#, owner.as_str()).fetch_one(&self.pool).await?;
47 Ok(FactsHealth {
48 generation: checkpoint.as_ref().map_or(0, |row| row.generation),
49 pending: queue.pending,
50 leased: queue.leased,
51 retries: queue.retries,
52 oldest_pending_at: queue.oldest_pending_at,
53 last_applied_at: checkpoint.as_ref().and_then(|row| row.last_applied_at),
54 last_recorded_at: checkpoint.and_then(|row| row.last_recorded_at),
55 })
56 }
57}