Skip to main content

systemprompt_runtime/reporting/
rebuild.rs

1//! Analytics baseline rebuild in fenced, committed phases. Phase A mints the
2//! millisecond cutoff and opens the generation under every source lock, so no
3//! writer is mid-flight; phase A2 empties the targets in its own transaction;
4//! phase B snapshots each source as keyset pages, one transaction and one
5//! heartbeat per page, until a page comes back short; phase C flips
6//! `initialized` if this generation is still the live one. Nothing here holds
7//! a source lock or an open transaction for longer than one page.
8//!
9//! `initialize` is synchronous: it returns once the projection is initialized
10//! or another node is known to be building it. `rebuild` is unconditional and
11//! fences any rebuild in flight. A rebuild whose heartbeat is older than
12//! `STALE_HEARTBEAT` is presumed dead and taken over; a page is a bounded
13//! statement, so a live rebuild heartbeats well inside that window.
14//!
15//! Copyright (c) systemprompt.io — Business Source License 1.1.
16//! See <https://systemprompt.io> for licensing details.
17
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use chrono::Utc;
22use sqlx::PgPool;
23use systemprompt_analytics::AnalyticsError;
24use systemprompt_analytics::projection::{
25    self, ReportingProjector, SOURCE_DEFINITIONS, SourceDefinition,
26};
27use systemprompt_database::DbPool;
28use systemprompt_database::resilience::{Outcome, RetryConfig, retry_async};
29
30use crate::RuntimeResult;
31
32const PAGE_ROWS: i64 = 10_000;
33const STALE_HEARTBEAT: chrono::Duration = chrono::Duration::seconds(120);
34
35// Why: each fence opens a generation that supersedes any in flight, so two
36// forced rebuilds racing will supersede each other for as long as both keep
37// retrying. Bounding the retries turns an unbounded livelock into a typed
38// outcome the caller can act on; the backoff gives the winning run room to
39// finish rather than being superseded again immediately.
40const MAX_SUPERSEDED_RETRIES: u32 = 3;
41const SUPERSEDED_BACKOFF: Duration = Duration::from_millis(250);
42
43/// What a rebuild run actually did, which is not always what was asked for.
44///
45/// `InProgressElsewhere` is the one to handle: a forced rebuild is exclusive,
46/// because every fence mints a generation superseding whatever is in flight,
47/// so concurrent forced rebuilds cannot all win and the losers rebuilt
48/// nothing. A caller that needs a baseline containing rows it has just
49/// written must either act on this outcome or serialise its rebuilds; it
50/// cannot get that guarantee by retrying, which is what made an earlier
51/// unbounded retry a livelock.
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub enum RebuildOutcome {
54    Rebuilt,
55    AlreadyInitialized,
56    InProgressElsewhere,
57}
58
59#[derive(Debug, Clone, Copy, PartialEq, Eq)]
60enum Mode {
61    IfNeeded,
62    Force,
63}
64
65#[derive(Debug, Clone, Copy)]
66struct Plan {
67    generation: i64,
68}
69
70pub async fn initialize(db: &DbPool) -> RuntimeResult<RebuildOutcome> {
71    run(db, Mode::IfNeeded).await
72}
73
74pub async fn rebuild(db: &DbPool) -> RuntimeResult<RebuildOutcome> {
75    run(db, Mode::Force).await
76}
77
78async fn run(db: &DbPool, mode: Mode) -> RuntimeResult<RebuildOutcome> {
79    let pool = db.write_pool_arc()?;
80    let mut superseded = 0u32;
81    loop {
82        let Some(plan) = fence(&pool, mode).await? else {
83            return Ok(if in_progress_elsewhere(&pool).await? {
84                RebuildOutcome::InProgressElsewhere
85            } else {
86                RebuildOutcome::AlreadyInitialized
87            });
88        };
89        let started = Instant::now();
90        match build(&pool, plan).await {
91            Ok(()) => {
92                tracing::info!(
93                    generation = plan.generation,
94                    elapsed_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
95                    "Analytics baseline rebuilt"
96                );
97                return Ok(RebuildOutcome::Rebuilt);
98            },
99            Err(error) if error.is_rebuild_superseded() => {
100                superseded += 1;
101                if superseded > MAX_SUPERSEDED_RETRIES {
102                    tracing::warn!(
103                        generation = plan.generation,
104                        attempts = superseded,
105                        "Analytics baseline rebuild superseded on every attempt; \
106                         another rebuild owns the generation"
107                    );
108                    return Ok(RebuildOutcome::InProgressElsewhere);
109                }
110                tracing::warn!(
111                    generation = plan.generation,
112                    attempt = superseded,
113                    "Analytics baseline rebuild superseded; starting over"
114                );
115                tokio::time::sleep(SUPERSEDED_BACKOFF * superseded).await;
116            },
117            Err(error) => return Err(error.into()),
118        }
119    }
120}
121
122async fn build(pool: &Arc<PgPool>, plan: Plan) -> Result<(), AnalyticsError> {
123    clear(pool, plan).await?;
124    for definition in SOURCE_DEFINITIONS {
125        snapshot_source(pool, plan, definition).await?;
126    }
127    finish(pool, plan).await
128}
129
130const fn retry_config() -> RetryConfig {
131    RetryConfig {
132        max_attempts: 4,
133        base_delay: Duration::from_millis(50),
134        max_delay: Duration::from_millis(800),
135        jitter: true,
136    }
137}
138
139fn classify(error: &AnalyticsError) -> Outcome {
140    match error {
141        AnalyticsError::Repository(repository) if repository.is_serialization_failure() => {
142            Outcome::Transient { retry_after: None }
143        },
144        _ => Outcome::Permanent,
145    }
146}
147
148async fn fence(pool: &Arc<PgPool>, mode: Mode) -> Result<Option<Plan>, AnalyticsError> {
149    retry_async(
150        &retry_config(),
151        "reporting-rebuild-fence",
152        classify,
153        || async {
154            let mut tx = pool.begin().await?;
155            projection::lock_user_deletion(&mut tx).await?;
156            projection::lock_sources(&mut tx).await?;
157            projection::lock_projector(&mut tx).await?;
158            let state = projection::rebuild_state(&mut tx).await?;
159            if mode == Mode::IfNeeded && (state.initialized || heartbeat_is_fresh(&state)) {
160                tx.commit().await?;
161                return Ok(None);
162            }
163            let cutoff = projection::next_cutoff_revision(&mut tx).await?;
164            let generation = ReportingProjector::begin_rebuild(&mut tx, cutoff).await?;
165            tx.commit().await?;
166            tracing::info!(generation, cutoff, "Analytics baseline rebuild started");
167            Ok(Some(Plan { generation }))
168        },
169    )
170    .await
171}
172
173fn heartbeat_is_fresh(state: &projection::RebuildState) -> bool {
174    state
175        .rebuild_heartbeat_at
176        .is_some_and(|beat| Utc::now() - beat < STALE_HEARTBEAT)
177}
178
179async fn in_progress_elsewhere(pool: &Arc<PgPool>) -> Result<bool, AnalyticsError> {
180    let mut connection = pool.acquire().await?;
181    let state = projection::rebuild_state(&mut connection).await?;
182    Ok(!state.initialized && heartbeat_is_fresh(&state))
183}
184
185async fn clear(pool: &Arc<PgPool>, plan: Plan) -> Result<(), AnalyticsError> {
186    retry_async(
187        &retry_config(),
188        "reporting-rebuild-clear",
189        classify,
190        || async {
191            let mut tx = pool.begin().await?;
192            projection::lock_projector(&mut tx).await?;
193            ReportingProjector::clear_targets(&mut tx, plan.generation).await?;
194            tx.commit().await?;
195            Ok(())
196        },
197    )
198    .await
199}
200
201async fn snapshot_source(
202    pool: &Arc<PgPool>,
203    plan: Plan,
204    definition: &SourceDefinition,
205) -> Result<(), AnalyticsError> {
206    let started = Instant::now();
207    let mut after: Option<String> = None;
208    let mut written = 0;
209    loop {
210        let page = retry_async(
211            &retry_config(),
212            "reporting-rebuild-page",
213            classify,
214            || async {
215                let mut tx = pool.begin().await?;
216                projection::lock_projector(&mut tx).await?;
217                let page = projection::write_snapshot_page(
218                    &mut tx,
219                    definition,
220                    after.as_deref(),
221                    PAGE_ROWS,
222                )
223                .await?;
224                projection::heartbeat_rebuild(
225                    &mut tx,
226                    plan.generation,
227                    definition.table,
228                    page.written,
229                )
230                .await?;
231                tx.commit().await?;
232                Ok(page)
233            },
234        )
235        .await?;
236        written += page.written;
237        if page.fetched < PAGE_ROWS {
238            break;
239        }
240        after = page.last_key;
241    }
242    tracing::info!(
243        source = definition.table,
244        rows = written,
245        elapsed_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
246        "Analytics baseline source snapshotted"
247    );
248    Ok(())
249}
250
251async fn finish(pool: &Arc<PgPool>, plan: Plan) -> Result<(), AnalyticsError> {
252    retry_async(
253        &retry_config(),
254        "reporting-rebuild-finish",
255        classify,
256        || async {
257            let mut tx = pool.begin().await?;
258            projection::lock_projector(&mut tx).await?;
259            ReportingProjector::finish_rebuild(&mut tx, plan.generation).await?;
260            tx.commit().await?;
261            Ok(())
262        },
263    )
264    .await
265}