Skip to main content

systemprompt_analytics/projection/
state.rs

1//! Projection bookkeeping: the singleton state row, the projector lock and
2//! the durable-queue lag behind it.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use chrono::{DateTime, Utc};
8use serde::Serialize;
9use sqlx::{PgConnection, PgPool};
10
11use crate::Result;
12
13const PROJECTOR_LOCK: i64 = 0x5350_414e_414c_5954;
14
15/// Baseline state plus the pending durable facts still owed to the projection.
16#[derive(Debug, Clone, Copy, Serialize)]
17pub struct ProjectionStatus {
18    pub initialized: bool,
19    pub generation: i64,
20    pub rebuilt_at: Option<DateTime<Utc>>,
21    pub pending_count: i64,
22    pub oldest_pending_at: Option<DateTime<Utc>>,
23    pub last_processed_at: Option<DateTime<Utc>>,
24}
25
26pub async fn lock_projector(connection: &mut PgConnection) -> Result<()> {
27    sqlx::query!("SELECT pg_advisory_xact_lock($1)", PROJECTOR_LOCK)
28        .fetch_one(connection)
29        .await?;
30    Ok(())
31}
32
33pub async fn lock_user_deletion(connection: &mut PgConnection) -> Result<()> {
34    sqlx::query_scalar!(r#"SELECT public.lock_user_deletion_for_retention() AS "locked!""#)
35        .fetch_one(connection)
36        .await?;
37    Ok(())
38}
39
40pub async fn is_initialized(connection: &mut PgConnection) -> Result<bool> {
41    Ok(sqlx::query_scalar!(
42        r#"SELECT initialized AS "initialized!" FROM analytics_projection_state WHERE singleton"#
43    )
44    .fetch_one(connection)
45    .await?)
46}
47
48pub async fn next_cutoff_revision(connection: &mut PgConnection) -> Result<i64> {
49    Ok(
50        sqlx::query_scalar!(r#"SELECT nextval('event_outbox_reporting_revision') AS "cutoff!""#)
51            .fetch_one(connection)
52            .await?,
53    )
54}
55
56pub async fn status(pool: &PgPool, consumer: &str) -> Result<ProjectionStatus> {
57    Ok(sqlx::query_as!(
58        ProjectionStatus,
59        r#"SELECT initialized, generation, rebuilt_at,
60            (SELECT COUNT(*) FROM event_outbox WHERE consumer = $1 AND processed_at IS NULL) AS "pending_count!",
61            (SELECT MIN(created_at) FROM event_outbox WHERE consumer = $1 AND processed_at IS NULL) AS oldest_pending_at,
62            (SELECT MAX(processed_at) FROM event_outbox WHERE consumer = $1) AS last_processed_at
63         FROM analytics_projection_state WHERE singleton"#,
64        consumer
65    )
66    .fetch_one(pool)
67    .await?)
68}