Skip to main content

systemprompt_runtime/reporting/
worker.rs

1//! Polling durable analytics delivery with atomic projection acknowledgement.
2//! The server's one reporting task: it first makes sure a baseline exists
3//! (building it in the background while the server serves), then drains
4//! captured facts into the projection a batch at a time.
5//!
6//! `BATCH_SIZE` facts are applied and acknowledged per transaction and
7//! `TICK_LIMIT` facts drained per tick before yielding to the metrics gauge.
8//! Waiting for the baseline blocks this task, never the server: a node that
9//! finds another one mid-rebuild waits for it, and a failed attempt is
10//! retried rather than taking the process down. A batch that fails to apply
11//! is re-driven in halves until the failing fact stands alone; that fact
12//! stays pending, is skipped for the rest of the drain, and is reported in
13//! the returned error after everything else applied.
14//!
15//! Copyright (c) systemprompt.io — Business Source License 1.1.
16//! See <https://systemprompt.io> for licensing details.
17
18use std::time::Duration;
19
20use sqlx::PgPool;
21use systemprompt_analytics::AnalyticsError;
22use systemprompt_analytics::projection::{
23    self, REPORTING_CONSUMER, REPORTING_KIND, REPORTING_VERSION, ReportingProjector, ReportingRow,
24};
25use systemprompt_database::DbPool;
26use systemprompt_events::services::durable::{DeliveryBatch, OutboxConsumer};
27use systemprompt_identifiers::EventOutboxId;
28use tokio::task::JoinHandle;
29
30use super::rebuild::RebuildOutcome;
31use crate::RuntimeResult;
32
33const BASELINE_RETRY: Duration = Duration::from_secs(30);
34const BATCH_SIZE: usize = 1000;
35const TICK_LIMIT: usize = 10_000;
36
37pub fn spawn(db: &DbPool) -> RuntimeResult<JoinHandle<()>> {
38    let pool = db.write_pool_arc()?;
39    let db = DbPool::clone(db);
40    Ok(tokio::spawn(async move {
41        ensure_baseline(&db).await;
42        let mut interval = tokio::time::interval(Duration::from_secs(1));
43        interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
44        let mut failing = false;
45        let mut next_metrics = tokio::time::Instant::now();
46        loop {
47            interval.tick().await;
48            match drain(&pool, TICK_LIMIT).await {
49                Ok(processed) => {
50                    if failing {
51                        tracing::info!("Analytics projection processing recovered");
52                        failing = false;
53                    }
54                    if processed == TICK_LIMIT {
55                        interval.reset_immediately();
56                    }
57                },
58                Err(error) => {
59                    metrics::counter!("analytics_projection_failures_total").increment(1);
60                    if !failing {
61                        tracing::error!(
62                            error = %error,
63                            "Analytics projection paused; pending facts retained for retry"
64                        );
65                        failing = true;
66                    }
67                },
68            }
69            if tokio::time::Instant::now() >= next_metrics {
70                if let Ok(status) = super::status::from_pool(&pool).await {
71                    metrics::gauge!("analytics_projection_pending")
72                        .set(status.pending_count as f64);
73                    metrics::gauge!("analytics_projection_applied_last_minute")
74                        .set(status.applied_last_minute as f64);
75                    let age = status.oldest_pending_at.map_or(0, |oldest| {
76                        (chrono::Utc::now() - oldest).num_seconds().max(0)
77                    });
78                    metrics::gauge!("analytics_projection_oldest_pending_seconds").set(age as f64);
79                }
80                next_metrics = tokio::time::Instant::now() + Duration::from_secs(30);
81            }
82        }
83    }))
84}
85
86async fn ensure_baseline(db: &DbPool) {
87    loop {
88        match super::rebuild::initialize(db).await {
89            Ok(RebuildOutcome::Rebuilt | RebuildOutcome::AlreadyInitialized) => return,
90            Ok(RebuildOutcome::InProgressElsewhere) => {
91                tracing::info!("Analytics baseline is being rebuilt by another node; waiting");
92            },
93            Err(error) => {
94                metrics::counter!("analytics_projection_failures_total").increment(1);
95                tracing::error!(error = %error, "Analytics baseline rebuild failed; retrying");
96            },
97        }
98        tokio::time::sleep(BASELINE_RETRY).await;
99    }
100}
101
102pub async fn process_pending(db: &DbPool, limit: usize) -> RuntimeResult<usize> {
103    let pool = db.write_pool_arc()?;
104    drain(&pool, limit).await.map_err(Into::into)
105}
106
107async fn drain(pool: &PgPool, limit: usize) -> Result<usize, AnalyticsError> {
108    let outbox = OutboxConsumer::new(pool.clone());
109    let mut processed = 0;
110    let mut batch_size = BATCH_SIZE;
111    let mut poisoned: Vec<EventOutboxId> = Vec::new();
112    while processed < limit {
113        let want = batch_size.min(limit - processed);
114        let Some(batch) = outbox
115            .claim_batch(
116                REPORTING_CONSUMER,
117                i64::try_from(want).unwrap_or(i64::MAX),
118                &poisoned,
119            )
120            .await?
121        else {
122            break;
123        };
124        let claimed = batch.len();
125        let head = batch.ids().next().cloned();
126        match apply_batch(batch).await {
127            Ok(()) => {
128                processed += claimed;
129                batch_size = BATCH_SIZE;
130                metrics::counter!("analytics_projection_processed_total")
131                    .increment(u64::try_from(claimed).unwrap_or(u64::MAX));
132            },
133            Err(error) if claimed > 1 => {
134                tracing::warn!(
135                    error = %error,
136                    claimed,
137                    "Reporting batch failed; re-driving in halves to isolate the fact"
138                );
139                batch_size = claimed / 2;
140            },
141            Err(error) => {
142                if let Some(id) = head {
143                    tracing::error!(
144                        error = %error,
145                        outbox_id = %id,
146                        "Reporting fact cannot be applied; left pending and skipped"
147                    );
148                    poisoned.push(id);
149                } else {
150                    tracing::error!(
151                        error = %error,
152                        "Reporting batch failed with no fact to isolate"
153                    );
154                }
155                batch_size = BATCH_SIZE;
156            },
157        }
158    }
159    if let Some(first) = poisoned.first() {
160        return Err(AnalyticsError::invalid_argument(format!(
161            "{} reporting fact(s) left pending; first outbox id {first} (applied {processed} others)",
162            poisoned.len()
163        )));
164    }
165    Ok(processed)
166}
167
168async fn apply_batch(mut batch: DeliveryBatch) -> Result<(), AnalyticsError> {
169    let facts = batch.facts::<ReportingRow>()?;
170    for (id, fact) in &facts {
171        if fact.consumer != REPORTING_CONSUMER
172            || fact.kind != REPORTING_KIND
173            || fact.version != REPORTING_VERSION
174        {
175            return Err(AnalyticsError::invalid_argument(format!(
176                "unsupported reporting contract in outbox event {id}",
177            )));
178        }
179    }
180    projection::lock_projector(batch.connection()).await?;
181    for (_, fact) in &facts {
182        ReportingProjector::apply_fact(batch.connection(), &fact.data).await?;
183    }
184    batch.acknowledge_all().await?;
185    Ok(())
186}