systemprompt_analytics/projection/
state.rs1use 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#[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}