Skip to main content

systemprompt_analytics/feedback/
deltas.rs

1//! Fenced downstream batches support aggregate writes and checkpoint commit in
2//! one transaction.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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}