systemprompt_runtime/reporting/
worker.rs1use 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}