Skip to main content

systemprompt_analytics/feedback/
queue.rs

1//! Leases fence every completion; committed pending rows remain discoverable
2//! after restart.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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}