Skip to main content

systemprompt_runtime/reporting/
worker.rs

1//! Polling durable analytics delivery with atomic projection acknowledgement.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}