Skip to main content

systemprompt_analytics/snapshots/
jobs.rs

1//! Idempotent asynchronous range requests with restart-safe worker fencing.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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    // Why: the failure is recorded under the same lease fence as completion,
156    // so a worker whose lease was stolen cannot mark another worker's job.
157    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}