Skip to main content

systemprompt_runtime/reporting/
rebuild.rs

1//! Transactional source capture installation and analytics baseline rebuild.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use std::sync::Arc;
7use std::time::Duration;
8
9use sqlx::{PgConnection, PgPool};
10use systemprompt_analytics::AnalyticsError;
11use systemprompt_analytics::projection::{
12    self, ReportingProjector, ReportingRow, SOURCE_DEFINITIONS, SnapshotCursor,
13};
14use systemprompt_database::DbPool;
15use systemprompt_database::resilience::{Outcome, RetryConfig, retry_async};
16
17use crate::RuntimeResult;
18
19pub async fn initialize(db: &DbPool) -> RuntimeResult<()> {
20    configure(db, false).await
21}
22
23pub async fn rebuild(db: &DbPool) -> RuntimeResult<()> {
24    configure(db, true).await
25}
26
27async fn configure(db: &DbPool, force_rebuild: bool) -> RuntimeResult<()> {
28    let pool = db.write_pool_arc()?;
29    let retry = RetryConfig {
30        max_attempts: 4,
31        base_delay: Duration::from_millis(50),
32        max_delay: Duration::from_millis(800),
33        jitter: true,
34    };
35    let classify = |error: &AnalyticsError| match error {
36        AnalyticsError::Repository(repository) if repository.is_serialization_failure() => {
37            Outcome::Transient { retry_after: None }
38        },
39        _ => Outcome::Permanent,
40    };
41    retry_async(&retry, "reporting-rebuild", classify, || {
42        configure_once(&pool, force_rebuild)
43    })
44    .await?;
45    Ok(())
46}
47
48async fn configure_once(pool: &Arc<PgPool>, force_rebuild: bool) -> Result<(), AnalyticsError> {
49    let mut transaction = pool.begin().await.map_err(AnalyticsError::from)?;
50    projection::lock_user_deletion(&mut transaction).await?;
51    SnapshotCursor::lock_sources(&mut transaction).await?;
52    projection::lock_projector(&mut transaction).await?;
53    if force_rebuild || !projection::is_initialized(&mut transaction).await? {
54        rebuild_locked(&mut transaction).await?;
55    }
56    transaction.commit().await.map_err(AnalyticsError::from)?;
57    Ok(())
58}
59
60async fn rebuild_locked(connection: &mut PgConnection) -> Result<(), AnalyticsError> {
61    let cutoff = projection::next_cutoff_revision(&mut *connection).await?;
62    let generation = ReportingProjector::begin_rebuild(connection).await?;
63    for definition in SOURCE_DEFINITIONS {
64        let cursor = SnapshotCursor::open(&mut *connection, definition).await?;
65        loop {
66            let rows = cursor.fetch(&mut *connection).await?;
67            if rows.is_empty() {
68                break;
69            }
70            for row in rows {
71                ReportingProjector::apply_snapshot(
72                    connection,
73                    &ReportingRow {
74                        source: definition.source,
75                        key: row.entity_key,
76                        revision: cutoff,
77                        deleted: false,
78                        row: row.row,
79                    },
80                )
81                .await?;
82            }
83        }
84        cursor.close(&mut *connection).await?;
85    }
86    ReportingProjector::finish_rebuild(connection, generation, cutoff).await?;
87    Ok(())
88}