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    let mut first_cause: Option<String> = None;
113    while processed < limit {
114        let want = batch_size.min(limit - processed);
115        let Some(batch) = outbox
116            .claim_batch(
117                REPORTING_CONSUMER,
118                i64::try_from(want).unwrap_or(i64::MAX),
119                &poisoned,
120            )
121            .await?
122        else {
123            break;
124        };
125        let claimed = batch.len();
126        let head = batch.ids().next().cloned();
127        match apply_batch(batch).await {
128            Ok(()) => {
129                processed += claimed;
130                batch_size = BATCH_SIZE;
131                metrics::counter!("analytics_projection_processed_total")
132                    .increment(u64::try_from(claimed).unwrap_or(u64::MAX));
133            },
134            Err(error) if claimed > 1 => {
135                tracing::warn!(
136                    error = %error,
137                    claimed,
138                    "Reporting batch failed; re-driving in halves to isolate the fact"
139                );
140                batch_size = claimed / 2;
141            },
142            Err(error) => {
143                if let Some(id) = head {
144                    tracing::error!(
145                        error = %error,
146                        outbox_id = %id,
147                        "Reporting fact cannot be applied; left pending and skipped"
148                    );
149                    poisoned.push(id);
150                    first_cause.get_or_insert_with(|| error.to_string());
151                } else {
152                    tracing::error!(
153                        error = %error,
154                        "Reporting batch failed with no fact to isolate"
155                    );
156                }
157                batch_size = BATCH_SIZE;
158            },
159        }
160    }
161    if let Some(first) = poisoned.first() {
162        return Err(AnalyticsError::invalid_argument(format!(
163            "{} reporting fact(s) left pending; first outbox id {first} (applied {processed} \
164             others): {cause}",
165            poisoned.len(),
166            cause = first_cause.as_deref().unwrap_or("no cause recorded")
167        )));
168    }
169    Ok(processed)
170}
171
172async fn apply_batch(mut batch: DeliveryBatch) -> Result<(), AnalyticsError> {
173    let facts = batch.facts::<ReportingRow>()?;
174    for (id, fact) in &facts {
175        if fact.consumer != REPORTING_CONSUMER
176            || fact.kind != REPORTING_KIND
177            || fact.version != REPORTING_VERSION
178        {
179            return Err(AnalyticsError::invalid_argument(format!(
180                "unsupported reporting contract in outbox event {id}",
181            )));
182        }
183    }
184    projection::lock_projector(batch.connection()).await?;
185    for (_, fact) in &facts {
186        ReportingProjector::apply_fact(batch.connection(), &fact.data).await?;
187    }
188    batch.acknowledge_all().await?;
189    Ok(())
190}