systemprompt_runtime/reporting/
rebuild.rs1use 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}