systemprompt_runtime/reporting/
rebuild.rs1use 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
35const MAX_SUPERSEDED_RETRIES: u32 = 3;
41const SUPERSEDED_BACKOFF: Duration = Duration::from_millis(250);
42
43#[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}