systemprompt_runtime/reporting/
worker.rs1use std::time::Duration;
7
8use sqlx::PgPool;
9use systemprompt_analytics::AnalyticsError;
10use systemprompt_analytics::projection::{
11 self, REPORTING_CONSUMER, REPORTING_KIND, REPORTING_VERSION, ReportingProjector, ReportingRow,
12};
13use systemprompt_database::DbPool;
14use systemprompt_events::services::durable::OutboxConsumer;
15use tokio::task::JoinHandle;
16
17use crate::RuntimeResult;
18
19pub fn spawn(db: &DbPool) -> RuntimeResult<JoinHandle<()>> {
20 let pool = db.write_pool_arc()?;
21 Ok(tokio::spawn(async move {
22 let mut interval = tokio::time::interval(Duration::from_secs(1));
23 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
24 let mut failing = false;
25 let mut next_metrics = tokio::time::Instant::now();
26 loop {
27 interval.tick().await;
28 match drain(&pool, 256).await {
29 Ok(processed) => {
30 if failing {
31 tracing::info!("Analytics projection processing recovered");
32 failing = false;
33 }
34 if processed == 256 {
35 interval.reset_immediately();
36 }
37 },
38 Err(error) => {
39 metrics::counter!("analytics_projection_failures_total").increment(1);
40 if !failing {
41 tracing::error!(
42 error = %error,
43 "Analytics projection paused; pending facts retained for retry"
44 );
45 failing = true;
46 }
47 },
48 }
49 if tokio::time::Instant::now() >= next_metrics {
50 if let Ok(status) = super::status::from_pool(&pool).await {
51 metrics::gauge!("analytics_projection_pending")
52 .set(status.pending_count as f64);
53 let age = status.oldest_pending_at.map_or(0, |oldest| {
54 (chrono::Utc::now() - oldest).num_seconds().max(0)
55 });
56 metrics::gauge!("analytics_projection_oldest_pending_seconds").set(age as f64);
57 }
58 next_metrics = tokio::time::Instant::now() + Duration::from_secs(30);
59 }
60 }
61 }))
62}
63
64pub async fn process_pending(db: &DbPool, limit: usize) -> RuntimeResult<usize> {
65 let pool = db.write_pool_arc()?;
66 drain(&pool, limit).await.map_err(Into::into)
67}
68
69async fn drain(pool: &PgPool, limit: usize) -> Result<usize, AnalyticsError> {
70 let outbox = OutboxConsumer::new(pool.clone());
71 let mut processed = 0;
72 while processed < limit {
73 let Some(mut delivery) = outbox.claim(REPORTING_CONSUMER).await? else {
74 break;
75 };
76 let fact = delivery.fact::<ReportingRow>()?;
77 if fact.consumer != REPORTING_CONSUMER
78 || fact.kind != REPORTING_KIND
79 || fact.version != REPORTING_VERSION
80 {
81 return Err(AnalyticsError::invalid_argument(format!(
82 "unsupported reporting contract in outbox event {}",
83 delivery.id(),
84 )));
85 }
86 projection::lock_projector(delivery.connection()).await?;
87 ReportingProjector::apply_fact(delivery.connection(), &fact.data).await?;
88 delivery.acknowledge().await?;
89 metrics::counter!("analytics_projection_processed_total").increment(1);
90 processed += 1;
91 }
92 Ok(processed)
93}