systemprompt_analytics/feedback/
deltas.rs1use super::{DeltaClaim, DeltaLease, FactDelta, FeedbackFactsRepository, validation};
8use crate::Result;
9use systemprompt_identifiers::{AnalyticsWorkerId, UserId};
10
11impl FeedbackFactsRepository {
12 pub async fn claim_deltas(
13 &self,
14 owner: &UserId,
15 consumer: &str,
16 worker: &AnalyticsWorkerId,
17 claim: DeltaClaim,
18 ) -> Result<Option<DeltaLease>> {
19 let DeltaClaim {
20 limit,
21 lease_seconds,
22 } = claim;
23 if consumer.is_empty()
24 || consumer.len() > 128
25 || consumer.chars().any(char::is_control)
26 || !(1..=256).contains(&limit)
27 || !(1..=300).contains(&lease_seconds)
28 {
29 return Err(validation::invalid());
30 }
31 let mut tx = self.pool.begin().await?;
32 sqlx::query!(
33 "INSERT INTO analytics_fact_checkpoints(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
34 owner.as_str()
35 )
36 .execute(&mut *tx)
37 .await?;
38 sqlx::query!(
39 "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
40 owner.as_str()
41 )
42 .fetch_one(&mut *tx)
43 .await?;
44
45 sqlx::query!("INSERT INTO analytics_fact_consumers(owner_id,consumer) VALUES($1,$2) ON CONFLICT DO NOTHING", owner.as_str(), consumer).execute(&mut *tx).await?;
46 let state = sqlx::query!("SELECT generation FROM analytics_fact_consumers WHERE owner_id=$1 AND consumer=$2 AND (lease_until IS NULL OR lease_until<=clock_timestamp()) FOR UPDATE SKIP LOCKED", owner.as_str(), consumer).fetch_optional(&mut *tx).await?;
47 let Some(state) = state else {
48 tx.commit().await?;
49 return Ok(None);
50 };
51 let limit = i64::from(limit);
52 let through = sqlx::query_scalar!("SELECT MAX(generation) FROM (SELECT generation FROM analytics_fact_deltas WHERE owner_id=$1 AND generation>$2 ORDER BY generation LIMIT $3) batch", owner.as_str(), state.generation, limit).fetch_one(&mut *tx).await?;
53 let Some(through) = through else {
54 tx.commit().await?;
55 return Ok(None);
56 };
57 let seconds = f64::from(lease_seconds);
58 let updated = sqlx::query!(r#"UPDATE analytics_fact_consumers SET lease_worker=$3,lease_epoch=lease_epoch+1,lease_until=clock_timestamp()+make_interval(secs=>$4),lease_through=$5 WHERE owner_id=$1 AND consumer=$2 RETURNING lease_epoch,lease_until AS "lease_until!""#, owner.as_str(), consumer, worker.as_str(), seconds, through).fetch_one(&mut *tx).await?;
59 tx.commit().await?;
60 Ok(Some(DeltaLease {
61 consumer: consumer.to_owned(),
62 worker_id: worker.clone(),
63 epoch: updated.lease_epoch,
64 after_generation: state.generation,
65 through_generation: through,
66 expires_at: updated.lease_until,
67 }))
68 }
69
70 pub async fn delta_batch(&self, owner: &UserId, lease: &DeltaLease) -> Result<Vec<FactDelta>> {
71 let mut tx = self.pool.begin().await?;
72 Self::lock_delta_lease(&mut tx, owner, lease).await?;
73 let rows = sqlx::query!("SELECT generation,fact_kind,source,fact_id,before_fact,after_fact FROM analytics_fact_deltas WHERE owner_id=$1 AND generation>$2 AND generation<=$3 ORDER BY generation LIMIT 256", owner.as_str(), lease.after_generation, lease.through_generation).fetch_all(&mut *tx).await?;
74 let deltas = rows
75 .into_iter()
76 .map(|row| {
77 Ok(FactDelta {
78 generation: row.generation,
79 key: systemprompt_models::feedback::analytics::AnalyticsFactKey {
80 kind: validation::parse_kind(&row.fact_kind)?,
81 source: row.source,
82 id: systemprompt_identifiers::AnalyticsFactId::new(row.fact_id),
83 },
84 before: row.before_fact.map(serde_json::from_value).transpose()?,
85 after: row.after_fact.map(serde_json::from_value).transpose()?,
86 })
87 })
88 .collect::<Result<Vec<_>>>()?;
89 tx.commit().await?;
90 Ok(deltas)
91 }
92
93 pub async fn lock_delta_lease(
94 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
95 owner: &UserId,
96 lease: &DeltaLease,
97 ) -> Result<()> {
98 let valid = sqlx::query_scalar!("SELECT generation FROM analytics_fact_consumers WHERE owner_id=$1 AND consumer=$2 AND lease_worker=$3 AND lease_epoch=$4 AND generation=$5 AND lease_through=$6 AND lease_until>clock_timestamp() FOR UPDATE", owner.as_str(), &lease.consumer, lease.worker_id.as_str(), lease.epoch, lease.after_generation, lease.through_generation).fetch_optional(&mut **tx).await?;
99 if valid.is_none() {
100 return Err(validation::invalid());
101 }
102 Ok(())
103 }
104
105 pub async fn complete_delta_batch(
106 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
107 owner: &UserId,
108 lease: &DeltaLease,
109 ) -> Result<()> {
110 Self::lock_delta_lease(tx, owner, lease).await?;
111 let updated = sqlx::query!("UPDATE analytics_fact_consumers SET generation=$6,lease_worker=NULL,lease_until=NULL,lease_through=NULL WHERE owner_id=$1 AND consumer=$2 AND lease_worker=$3 AND lease_epoch=$4 AND generation=$5 AND lease_through=$6 AND lease_until>clock_timestamp()", owner.as_str(), &lease.consumer, lease.worker_id.as_str(), lease.epoch, lease.after_generation, lease.through_generation).execute(&mut **tx).await?;
112 if updated.rows_affected() != 1 {
113 return Err(validation::invalid());
114 }
115 sqlx::query!("UPDATE analytics_fact_deltas SET consumed_at=clock_timestamp() WHERE owner_id=$1 AND consumed_at IS NULL AND generation<=(SELECT MIN(generation) FROM analytics_fact_consumers WHERE owner_id=$1)", owner.as_str()).execute(&mut **tx).await?;
116 Ok(())
117 }
118}