返回 CodeWhale
alerts.rs
根目录 / crates / tui / src / fleet / alerts.rs
1 //! Opt-in fleet alert routing and adapter payloads.
2
3 #![allow(dead_code)]
4
5 use std::collections::BTreeMap;
6 use std::time::Duration;
7
8 use anyhow::{Context, Result, anyhow};
9 use codewhale_protocol::fleet::{
10 FleetAlertEventClass, FleetReceipt, FleetRunId, FleetTaskFailureKind, FleetWorkerEvent,
11 FleetWorkerEventPayload,
12 };
13 use serde::{Deserialize, Serialize};
14 use serde_json::{Value, json};
15
16 const DEFAULT_ALERT_TIMEOUT_SECONDS: u64 = 10;
17
18 #[derive(Debug, Clone, Serialize, Deserialize, Default)]
19 pub struct FleetAlertConfig {
20 #[serde(default)]
21 pub enabled: bool,
22 #[serde(default)]
23 pub dry_run: bool,
24 #[serde(default)]
25 pub routes: Vec<FleetAlertRoute>,
26 #[serde(default)]
27 pub adapters: BTreeMap<String, FleetAlertAdapterConfig>,
28 }
29
30 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
31 pub struct FleetAlertRoute {
32 #[serde(default)]
33 #[serde(skip_serializing_if = "Vec::is_empty")]
34 pub events: Vec<FleetAlertEventClass>,
35 pub adapter: String,
36 }
37
38 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
39 #[serde(tag = "kind", rename_all = "snake_case")]
40 pub enum FleetAlertAdapterConfig {
41 Slack {
42 webhook_env: String,
43 #[serde(skip_serializing_if = "Option::is_none")]
44 channel: Option<String>,
45 },
46 Webhook {
47 url_env: String,
48 #[serde(skip_serializing_if = "Option::is_none")]
49 secret_env: Option<String>,
50 },
51 PagerDuty {
52 routing_key_env: String,
53 #[serde(default = "default_pagerduty_severity")]
54 severity: String,
55 },
56 }
57
58 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
59 pub struct FleetAlertEvent {
60 pub class: FleetAlertEventClass,
61 pub run_id: FleetRunId,
62 #[serde(skip_serializing_if = "Option::is_none")]
63 pub worker_id: Option<String>,
64 #[serde(skip_serializing_if = "Option::is_none")]
65 pub task_id: Option<String>,
66 pub status: String,
67 pub reason: String,
68 }
69
70 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
71 pub struct FleetAlertDelivery {
72 pub adapter: String,
73 pub event_class: FleetAlertEventClass,
74 pub dry_run: bool,
75 pub sent: bool,
76 pub redacted_payload: Value,
77 }
78
79 pub trait FleetAlertSecretResolver {
80 fn resolve(&self, name: &str) -> Option<String>;
81 }
82
83 #[derive(Debug, Clone, Copy, Default)]
84 pub struct FleetEnvSecretResolver;
85
86 impl FleetAlertSecretResolver for FleetEnvSecretResolver {
87 fn resolve(&self, name: &str) -> Option<String> {
88 std::env::var(name).ok().filter(|value| !value.is_empty())
89 }
90 }
91
92 pub struct FleetAlertDispatcher<R = FleetEnvSecretResolver> {
93 config: FleetAlertConfig,
94 resolver: R,
95 }
96
97 impl FleetAlertConfig {
98 pub fn dry_run_for_adapter(adapter: FleetAlertAdapterConfig) -> Self {
99 let mut adapters = BTreeMap::new();
100 adapters.insert("dry-run".to_string(), adapter);
101 Self {
102 enabled: true,
103 dry_run: true,
104 routes: vec![FleetAlertRoute {
105 events: Vec::new(),
106 adapter: "dry-run".to_string(),
107 }],
108 adapters,
109 }
110 }
111 }
112
113 impl<R> FleetAlertDispatcher<R>
114 where
115 R: FleetAlertSecretResolver,
116 {
117 pub fn new(config: FleetAlertConfig, resolver: R) -> Self {
118 Self { config, resolver }
119 }
120
121 pub fn dispatch(&self, event: &FleetAlertEvent) -> Result<Vec<FleetAlertDelivery>> {
122 if !self.config.enabled {
123 return Ok(Vec::new());
124 }
125 let mut deliveries = Vec::new();
126 for route in self
127 .config
128 .routes
129 .iter()
130 .filter(|route| route_matches(route, event.class))
131 {
132 let adapter = self.config.adapters.get(&route.adapter).ok_or_else(|| {
133 anyhow!("fleet alert adapter {} is not configured", route.adapter)
134 })?;
135 let prepared = prepare_alert(&route.adapter, adapter, event, self.config.dry_run)?;
136 let sent = if self.config.dry_run {
137 false
138 } else {
139 send_alert(adapter, &prepared.body, &self.resolver)?
140 };
141 deliveries.push(FleetAlertDelivery {
142 adapter: route.adapter.clone(),
143 event_class: event.class,
144 dry_run: self.config.dry_run,
145 sent,
146 redacted_payload: prepared.redacted_payload,
147 });
148 }
149 Ok(deliveries)
150 }
151 }
152
153 impl FleetAlertEvent {
154 pub fn stale_from_worker_event(event: &FleetWorkerEvent) -> Option<Self> {
155 let FleetWorkerEventPayload::Stale { last_heartbeat_at } = &event.payload else {
156 return None;
157 };
158 Some(Self {
159 class: FleetAlertEventClass::Stale,
160 run_id: event.run_id.clone(),
161 worker_id: Some(event.worker_id.clone()),
162 task_id: Some(event.task_id.clone()),
163 status: "stale".to_string(),
164 reason: last_heartbeat_at
165 .as_ref()
166 .map(|ts| format!("worker heartbeat stale since {ts}"))
167 .unwrap_or_else(|| "worker heartbeat is stale".to_string()),
168 })
169 }
170
171 pub fn restart_exhausted(
172 run_id: FleetRunId,
173 worker_id: impl Into<String>,
174 task_id: impl Into<String>,
175 reason: impl Into<String>,
176 ) -> Self {
177 Self {
178 class: FleetAlertEventClass::RestartExhausted,
179 run_id,
180 worker_id: Some(worker_id.into()),
181 task_id: Some(task_id.into()),
182 status: "failed".to_string(),
183 reason: reason.into(),
184 }
185 }
186
187 pub fn needs_human(
188 run_id: FleetRunId,
189 worker_id: Option<String>,
190 task_id: Option<String>,
191 reason: impl Into<String>,
192 ) -> Self {
193 Self {
194 class: FleetAlertEventClass::NeedsHuman,
195 run_id,
196 worker_id,
197 task_id,
198 status: "needs_human".to_string(),
199 reason: reason.into(),
200 }
201 }
202
203 pub fn budget_exceeded(
204 run_id: FleetRunId,
205 worker_id: Option<String>,
206 task_id: Option<String>,
207 reason: impl Into<String>,
208 ) -> Self {
209 Self {
210 class: FleetAlertEventClass::BudgetExceeded,
211 run_id,
212 worker_id,
213 task_id,
214 status: "budget_exceeded".to_string(),
215 reason: reason.into(),
216 }
217 }
218
219 pub fn verifier_failed(receipt: &FleetReceipt) -> Option<Self> {
220 if receipt.failure_kind != Some(FleetTaskFailureKind::Verifier) {
221 return None;
222 }
223 Some(Self {
224 class: FleetAlertEventClass::VerifierFailed,
225 run_id: receipt.run_id.clone(),
226 worker_id: Some(receipt.worker_id.clone()),
227 task_id: Some(receipt.task_id.clone()),
228 status: "verifier_failed".to_string(),
229 reason: receipt
230 .score
231 .as_ref()
232 .and_then(|score| score.notes.clone())
233 .unwrap_or_else(|| "verifier failed".to_string()),
234 })
235 }
236
237 pub fn run_completed(run_id: FleetRunId, reason: impl Into<String>) -> Self {
238 Self {
239 class: FleetAlertEventClass::RunCompleted,
240 run_id,
241 worker_id: None,
242 task_id: None,
243 status: "completed".to_string(),
244 reason: reason.into(),
245 }
246 }
247
248 pub fn inspection_commands(&self) -> Vec<String> {
249 let mut commands = vec!["codewhale fleet status".to_string()];
250 if let Some(worker_id) = &self.worker_id {
251 commands.push(format!("codewhale fleet inspect {worker_id}"));
252 }
253 commands
254 }
255 }
256
257 struct PreparedAlert {
258 body: Value,
259 redacted_payload: Value,
260 }
261
262 fn prepare_alert(
263 adapter_name: &str,
264 adapter: &FleetAlertAdapterConfig,
265 event: &FleetAlertEvent,
266 dry_run: bool,
267 ) -> Result<PreparedAlert> {
268 let safe_event = safe_event_payload(event);
269 let prepared = match adapter {
270 FleetAlertAdapterConfig::Slack {
271 webhook_env,
272 channel,
273 } => {
274 let body = slack_body(event, channel.as_deref());
275 let redacted_payload = json!({
276 "adapter": adapter_name,
277 "kind": "slack",
278 "dry_run": dry_run,
279 "target": redacted_env(webhook_env),
280 "event": safe_event,
281 "body": body,
282 });
283 PreparedAlert {
284 body,
285 redacted_payload,
286 }
287 }
288 FleetAlertAdapterConfig::Webhook {
289 url_env,
290 secret_env,
291 } => {
292 let body = json!({
293 "source": "codewhale",
294 "event": safe_event,
295 });
296 let redacted_payload = json!({
297 "adapter": adapter_name,
298 "kind": "webhook",
299 "dry_run": dry_run,
300 "target": redacted_env(url_env),
301 "headers": redacted_secret_header(secret_env.as_deref()),
302 "body": body,
303 });
304 PreparedAlert {
305 body,
306 redacted_payload,
307 }
308 }
309 FleetAlertAdapterConfig::PagerDuty {
310 routing_key_env,
311 severity,
312 } => {
313 let body = pagerduty_body(event, severity, redacted_env(routing_key_env));
314 let redacted_payload = json!({
315 "adapter": adapter_name,
316 "kind": "pagerduty",
317 "dry_run": dry_run,
318 "target": "https://events.pagerduty.com/v2/enqueue",
319 "body": body,
320 });
321 PreparedAlert {
322 body,
323 redacted_payload,
324 }
325 }
326 };
327 Ok(prepared)
328 }
329
330 fn send_alert<R>(
331 adapter: &FleetAlertAdapterConfig,
332 redacted_body: &Value,
333 resolver: &R,
334 ) -> Result<bool>
335 where
336 R: FleetAlertSecretResolver,
337 {
338 let client = crate::tls::reqwest_blocking_client_builder()
339 .timeout(Duration::from_secs(DEFAULT_ALERT_TIMEOUT_SECONDS))
340 .build()
341 .context("building fleet alert HTTP client")?;
342 match adapter {
343 FleetAlertAdapterConfig::Slack { webhook_env, .. } => {
344 let url = required_https_url(resolver, webhook_env)?;
345 client
346 .post(url)
347 .json(redacted_body)
348 .send()
349 .context("sending fleet Slack alert")?
350 .error_for_status()
351 .context("Slack alert rejected")?;
352 }
353 FleetAlertAdapterConfig::Webhook {
354 url_env,
355 secret_env,
356 } => {
357 let url = required_https_url(resolver, url_env)?;
358 let mut request = client.post(url).json(redacted_body);
359 if let Some(secret_env) = secret_env {
360 request = request.header(
361 "X-CodeWhale-Webhook-Secret",
362 required_secret(resolver, secret_env)?,
363 );
364 }
365 request
366 .send()
367 .context("sending fleet webhook alert")?
368 .error_for_status()
369 .context("webhook alert rejected")?;
370 }
371 FleetAlertAdapterConfig::PagerDuty {
372 routing_key_env,
373 severity,
374 } => {
375 let routing_key = required_secret(resolver, routing_key_env)?;
376 let mut body = redacted_body.clone();
377 if let Some(map) = body.as_object_mut() {
378 map.insert("routing_key".to_string(), Value::String(routing_key));
379 }
380 if let Some(payload) = body.get_mut("payload").and_then(Value::as_object_mut) {
381 payload.insert("severity".to_string(), Value::String(severity.clone()));
382 }
383 client
384 .post("https://events.pagerduty.com/v2/enqueue")
385 .json(&body)
386 .send()
387 .context("sending fleet PagerDuty alert")?
388 .error_for_status()
389 .context("PagerDuty alert rejected")?;
390 }
391 }
392 Ok(true)
393 }
394
395 fn route_matches(route: &FleetAlertRoute, class: FleetAlertEventClass) -> bool {
396 route.events.is_empty() || route.events.contains(&class)
397 }
398
399 fn safe_event_payload(event: &FleetAlertEvent) -> Value {
400 json!({
401 "class": event.class,
402 "run_id": event.run_id.0.clone(),
403 "worker_id": event.worker_id.clone(),
404 "task_id": event.task_id.clone(),
405 "status": event.status.clone(),
406 "reason": short_reason(&event.reason),
407 "commands": event.inspection_commands(),
408 })
409 }
410
411 fn slack_body(event: &FleetAlertEvent, channel: Option<&str>) -> Value {
412 let text = format!(
413 "Codewhale fleet {}: run={} task={} reason={}",
414 alert_class_label(event.class),
415 event.run_id.0,
416 event.task_id.as_deref().unwrap_or("-"),
417 short_reason(&event.reason)
418 );
419 let mut body = json!({
420 "text": text,
421 "blocks": [
422 {
423 "type": "section",
424 "text": {
425 "type": "mrkdwn",
426 "text": text
427 }
428 },
429 {
430 "type": "context",
431 "elements": [
432 {
433 "type": "mrkdwn",
434 "text": event.inspection_commands().join(" | ")
435 }
436 ]
437 }
438 ]
439 });
440 if let Some(channel) = channel
441 && let Some(map) = body.as_object_mut()
442 {
443 map.insert("channel".to_string(), Value::String(channel.to_string()));
444 }
445 body
446 }
447
448 fn pagerduty_body(event: &FleetAlertEvent, severity: &str, routing_key: String) -> Value {
449 json!({
450 "routing_key": routing_key,
451 "event_action": "trigger",
452 "payload": {
453 "summary": format!("Codewhale fleet {}: {}", alert_class_label(event.class), short_reason(&event.reason)),
454 "severity": severity,
455 "source": "codewhale",
456 "custom_details": safe_event_payload(event),
457 }
458 })
459 }
460
461 fn redacted_env(name: &str) -> String {
462 format!("<redacted:env:{name}>")
463 }
464
465 fn alert_class_label(class: FleetAlertEventClass) -> &'static str {
466 match class {
467 FleetAlertEventClass::Stale => "stale",
468 FleetAlertEventClass::RestartExhausted => "restart_exhausted",
469 FleetAlertEventClass::NeedsHuman => "needs_human",
470 FleetAlertEventClass::BudgetExceeded => "budget_exceeded",
471 FleetAlertEventClass::VerifierFailed => "verifier_failed",
472 FleetAlertEventClass::RunCompleted => "run_completed",
473 }
474 }
475
476 fn redacted_secret_header(secret_env: Option<&str>) -> Value {
477 match secret_env {
478 Some(name) => json!({ "X-CodeWhale-Webhook-Secret": redacted_env(name) }),
479 None => json!({}),
480 }
481 }
482
483 fn required_secret<R>(resolver: &R, name: &str) -> Result<String>
484 where
485 R: FleetAlertSecretResolver,
486 {
487 resolver
488 .resolve(name)
489 .ok_or_else(|| anyhow!("fleet alert secret {name} is not configured"))
490 }
491
492 fn required_https_url<R>(resolver: &R, name: &str) -> Result<String>
493 where
494 R: FleetAlertSecretResolver,
495 {
496 let url = resolver
497 .resolve(name)
498 .ok_or_else(|| anyhow!("fleet alert URL {name} is not configured"))?;
499 validate_https_alert_url(name, &url)?;
500 Ok(url)
501 }
502
503 fn validate_https_alert_url(name: &str, url: &str) -> Result<()> {
504 let parsed = reqwest::Url::parse(url)
505 .with_context(|| format!("fleet alert URL from {name} is not a valid URL"))?;
506 if parsed.scheme() != "https" {
507 return Err(anyhow!("fleet alert URL from {name} must use https"));
508 }
509 Ok(())
510 }
511
512 fn short_reason(reason: &str) -> String {
513 let trimmed = reason.trim();
514 if trimmed.len() <= 240 {
515 return trimmed.to_string();
516 }
517 let prefix: String = trimmed.chars().take(237).collect();
518 format!("{prefix}...")
519 }
520
521 fn default_pagerduty_severity() -> String {
522 "error".to_string()
523 }
524
525 #[cfg(test)]
526 mod tests {
527 use super::*;
528 use codewhale_protocol::fleet::{FleetScore, FleetTaskResult};
529
530 #[derive(Default)]
531 struct MapResolver {
532 values: BTreeMap<String, String>,
533 }
534
535 impl FleetAlertSecretResolver for MapResolver {
536 fn resolve(&self, name: &str) -> Option<String> {
537 self.values.get(name).cloned()
538 }
539 }
540
541 fn event(class: FleetAlertEventClass) -> FleetAlertEvent {
542 FleetAlertEvent {
543 class,
544 run_id: FleetRunId::from("run-1"),
545 worker_id: Some("worker-1".to_string()),
546 task_id: Some("task-a".to_string()),
547 status: "stale".to_string(),
548 reason: "worker heartbeat stale".to_string(),
549 }
550 }
551
552 #[test]
553 fn fleet_alert_disabled_by_default() {
554 let dispatcher =
555 FleetAlertDispatcher::new(FleetAlertConfig::default(), MapResolver::default());
556
557 let deliveries = dispatcher
558 .dispatch(&event(FleetAlertEventClass::Stale))
559 .unwrap();
560
561 assert!(deliveries.is_empty());
562 }
563
564 #[test]
565 fn fleet_alert_policy_routes_event_classes_to_adapters() {
566 let mut adapters = BTreeMap::new();
567 adapters.insert(
568 "ops-slack".to_string(),
569 FleetAlertAdapterConfig::Slack {
570 webhook_env: "FLEET_SLACK_WEBHOOK".to_string(),
571 channel: Some("#fleet".to_string()),
572 },
573 );
574 adapters.insert(
575 "release-webhook".to_string(),
576 FleetAlertAdapterConfig::Webhook {
577 url_env: "FLEET_WEBHOOK_URL".to_string(),
578 secret_env: None,
579 },
580 );
581 let dispatcher = FleetAlertDispatcher::new(
582 FleetAlertConfig {
583 enabled: true,
584 dry_run: true,
585 routes: vec![
586 FleetAlertRoute {
587 events: vec![FleetAlertEventClass::Stale],
588 adapter: "ops-slack".to_string(),
589 },
590 FleetAlertRoute {
591 events: vec![FleetAlertEventClass::RunCompleted],
592 adapter: "release-webhook".to_string(),
593 },
594 ],
595 adapters,
596 },
597 MapResolver::default(),
598 );
599
600 let deliveries = dispatcher
601 .dispatch(&event(FleetAlertEventClass::Stale))
602 .unwrap();
603
604 assert_eq!(deliveries.len(), 1);
605 assert_eq!(deliveries[0].adapter, "ops-slack");
606 assert_eq!(deliveries[0].event_class, FleetAlertEventClass::Stale);
607 assert!(!deliveries[0].sent);
608 assert_eq!(deliveries[0].redacted_payload["kind"], "slack");
609 }
610
611 #[test]
612 fn fleet_alert_dry_run_redacts_secrets() {
613 let mut adapters = BTreeMap::new();
614 adapters.insert(
615 "pager".to_string(),
616 FleetAlertAdapterConfig::PagerDuty {
617 routing_key_env: "FLEET_PD_ROUTING_KEY".to_string(),
618 severity: "critical".to_string(),
619 },
620 );
621 let mut resolver = MapResolver::default();
622 resolver.values.insert(
623 "FLEET_PD_ROUTING_KEY".to_string(),
624 "real-routing-key-secret".to_string(),
625 );
626 let dispatcher = FleetAlertDispatcher::new(
627 FleetAlertConfig {
628 enabled: true,
629 dry_run: true,
630 routes: vec![FleetAlertRoute {
631 events: vec![FleetAlertEventClass::RestartExhausted],
632 adapter: "pager".to_string(),
633 }],
634 adapters,
635 },
636 resolver,
637 );
638
639 let deliveries = dispatcher
640 .dispatch(&event(FleetAlertEventClass::RestartExhausted))
641 .unwrap();
642 let payload = serde_json::to_string(&deliveries[0].redacted_payload).unwrap();
643
644 assert!(payload.contains("<redacted:env:FLEET_PD_ROUTING_KEY>"));
645 assert!(!payload.contains("real-routing-key-secret"));
646 assert!(payload.contains("codewhale fleet inspect worker-1"));
647 }
648
649 #[test]
650 fn fleet_alert_url_validation_requires_https() {
651 validate_https_alert_url("FLEET_WEBHOOK_URL", "https://hooks.example.invalid/fleet")
652 .expect("https alert URL should be accepted");
653
654 let err =
655 validate_https_alert_url("FLEET_WEBHOOK_URL", "http://hooks.example.invalid/fleet")
656 .expect_err("cleartext alert URL should fail");
657 assert!(format!("{err:#}").contains("must use https"));
658 }
659
660 #[test]
661 fn required_https_url_uses_secret_resolver() {
662 let mut resolver = MapResolver::default();
663 resolver.values.insert(
664 "FLEET_WEBHOOK_URL".to_string(),
665 "https://hooks.example.invalid/fleet".to_string(),
666 );
667
668 let url = required_https_url(&resolver, "FLEET_WEBHOOK_URL").expect("resolve URL");
669 assert_eq!(url, "https://hooks.example.invalid/fleet");
670 }
671
672 #[test]
673 fn fleet_alert_event_is_derived_from_ledgered_stale_worker_event() {
674 let worker_event = FleetWorkerEvent {
675 seq: 4,
676 run_id: FleetRunId::from("run-1"),
677 worker_id: "worker-1".to_string(),
678 task_id: "task-a".to_string(),
679 timestamp: "2026-06-13T02:00:00Z".to_string(),
680 payload: FleetWorkerEventPayload::Stale {
681 last_heartbeat_at: Some("2026-06-13T01:57:00Z".to_string()),
682 },
683 extra: BTreeMap::new(),
684 };
685
686 let alert = FleetAlertEvent::stale_from_worker_event(&worker_event).unwrap();
687
688 assert_eq!(alert.class, FleetAlertEventClass::Stale);
689 assert_eq!(alert.worker_id.as_deref(), Some("worker-1"));
690 assert!(alert.reason.contains("2026-06-13T01:57:00Z"));
691 assert_eq!(
692 alert.inspection_commands(),
693 vec![
694 "codewhale fleet status".to_string(),
695 "codewhale fleet inspect worker-1".to_string()
696 ]
697 );
698 }
699
700 #[test]
701 fn fleet_alert_verifier_failed_event_is_derived_from_receipt() {
702 let receipt = FleetReceipt {
703 run_id: FleetRunId::from("run-1"),
704 task_id: "task-a".to_string(),
705 worker_id: "worker-1".to_string(),
706 attempt: None,
707 terminal_seq: None,
708 completed_at: "2026-06-13T02:00:00Z".to_string(),
709 result: FleetTaskResult::Fail,
710 failure_kind: Some(FleetTaskFailureKind::Verifier),
711 artifacts: vec![],
712 score: Some(FleetScore {
713 value: 0.0,
714 max: Some(1.0),
715 notes: Some("regex scorer could not be compiled".to_string()),
716 }),
717 resolved_route: None,
718 effective_permissions: None,
719 };
720
721 let alert = FleetAlertEvent::verifier_failed(&receipt).unwrap();
722
723 assert_eq!(alert.class, FleetAlertEventClass::VerifierFailed);
724 assert_eq!(alert.status, "verifier_failed");
725 assert!(alert.reason.contains("regex scorer"));
726 }
727 }
728
728 lines RUST