1use std::io::Read;
9use std::path::{Path, PathBuf};
10
11use chrono::{DateTime, Duration, Utc};
12use serde::{Deserialize, Serialize};
13
14use crate::atomic::write_atomic;
15use crate::error::Error;
16use crate::lock::{RunLock, Shared};
17use crate::paths::{nofollow, reject_symlink, RunPaths};
18use crate::projections::read_node;
19use crate::schema::{NodeId, RunId, Status};
20
21pub const TELEMETRY_SCHEMA_VERSION: u32 = 1;
23pub const TELEMETRY_PROTOCOL_VERSION: u32 = 1;
25pub const TELEMETRY_MAX_BYTES: usize = 4 * 1024;
27pub const TELEMETRY_FRESHNESS_SECS: i64 = 90;
29
30pub trait TelemetryClock {
32 fn now(&self) -> DateTime<Utc>;
34}
35
36#[derive(Debug, Clone, Copy, Default)]
38pub struct SystemTelemetryClock;
39
40impl TelemetryClock for SystemTelemetryClock {
41 fn now(&self) -> DateTime<Utc> {
42 Utc::now()
43 }
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
48#[serde(rename_all = "snake_case")]
49pub enum TelemetryState {
50 AgentActive,
52 ToolRunning,
54 Settled,
56 Shutdown,
58}
59
60#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62#[serde(deny_unknown_fields)]
63pub struct TelemetryUpdate {
64 pub schema_version: u32,
66 pub protocol_version: u32,
68 pub run_id: RunId,
70 pub node_id: NodeId,
72 pub attempt: u32,
74 pub state: TelemetryState,
76 #[serde(default, skip_serializing_if = "Option::is_none")]
78 pub active_tool_count: Option<u8>,
79 #[serde(default, skip_serializing_if = "Option::is_none")]
81 pub tool_name: Option<String>,
82}
83
84#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
86pub struct TelemetryAccepted {
87 pub accepted: bool,
89 pub run_id: RunId,
91 pub node_id: NodeId,
93 pub attempt: u32,
95 pub received_at: DateTime<Utc>,
97 pub expires_at: DateTime<Utc>,
99}
100
101#[derive(Debug, thiserror::Error)]
103pub enum TelemetryError {
104 #[error(transparent)]
106 Core(#[from] Error),
107 #[error("invalid telemetry request: {0}")]
109 InvalidRequest(serde_json::Error),
110 #[error("telemetry {what} exceeds {TELEMETRY_MAX_BYTES} bytes (got {bytes})")]
112 TooLarge {
113 what: &'static str,
115 bytes: usize,
117 },
118 #[error("unsupported telemetry schema_version {found}; expected {TELEMETRY_SCHEMA_VERSION}")]
120 UnsupportedSchema {
121 found: u32,
123 },
124 #[error(
126 "unsupported telemetry protocol_version {found}; expected {TELEMETRY_PROTOCOL_VERSION}"
127 )]
128 UnsupportedProtocol {
129 found: u32,
131 },
132 #[error("invalid telemetry tool metadata: {0}")]
134 InvalidMetadata(&'static str),
135 #[error("telemetry run_id {found} does not match current run {expected}")]
137 RunMismatch {
138 expected: RunId,
140 found: RunId,
142 },
143 #[error("canonical run state is not synchronized with the event log")]
146 RunStateNotCurrent,
147 #[error("no node {node_id} in this run")]
149 NodeNotFound {
150 node_id: NodeId,
152 },
153 #[error("node {node_id} is terminal ({status:?})")]
155 TerminalNode {
156 node_id: NodeId,
158 status: Status,
160 },
161 #[error("telemetry attempt {found} does not match current attempt {expected}")]
163 AttemptMismatch {
164 expected: u32,
166 found: u32,
168 },
169 #[error("server clock cannot represent telemetry expiry")]
171 ClockOverflow,
172}
173
174#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(rename_all = "snake_case")]
177pub enum TelemetrySampleStatus {
178 Absent,
180 Current,
182 Stale,
184 ClockUnreliable,
186 Invalid,
188}
189
190#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192pub struct TelemetryView {
193 pub sample: TelemetrySampleStatus,
195 #[serde(skip_serializing_if = "Option::is_none")]
197 pub state: Option<TelemetryState>,
198 #[serde(skip_serializing_if = "Option::is_none")]
200 pub age_ms: Option<i64>,
201 #[serde(skip_serializing_if = "Option::is_none")]
203 pub state_elapsed_ms: Option<i64>,
204 #[serde(skip_serializing_if = "Option::is_none")]
206 pub attempt: Option<u32>,
207 #[serde(skip_serializing_if = "Option::is_none")]
209 pub active_tool_count: Option<u8>,
210 #[serde(skip_serializing_if = "Option::is_none")]
212 pub tool_name: Option<String>,
213}
214
215impl TelemetryView {
216 fn bare(sample: TelemetrySampleStatus) -> Self {
217 Self {
218 sample,
219 state: None,
220 age_ms: None,
221 state_elapsed_ms: None,
222 attempt: None,
223 active_tool_count: None,
224 tool_name: None,
225 }
226 }
227}
228
229#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
231#[serde(deny_unknown_fields)]
232struct StoredTelemetrySample {
233 schema_version: u32,
234 protocol_version: u32,
235 run_id: RunId,
236 node_id: NodeId,
237 attempt: u32,
238 state: TelemetryState,
239 #[serde(default, skip_serializing_if = "Option::is_none")]
240 active_tool_count: Option<u8>,
241 #[serde(default, skip_serializing_if = "Option::is_none")]
242 tool_name: Option<String>,
243 state_since: DateTime<Utc>,
244 received_at: DateTime<Utc>,
245 expires_at: DateTime<Utc>,
246}
247
248pub fn parse_telemetry_update(bytes: &[u8]) -> Result<TelemetryUpdate, TelemetryError> {
251 ensure_size("request", bytes.len())?;
252 let update: TelemetryUpdate =
253 serde_json::from_slice(bytes).map_err(TelemetryError::InvalidRequest)?;
254 validate_update_shape(&update)?;
255 let normalized = serde_json::to_vec(&update).map_err(TelemetryError::InvalidRequest)?;
256 ensure_size("normalized request", normalized.len())?;
257 Ok(update)
258}
259
260pub fn update_telemetry(
263 paths: &RunPaths,
264 update: &TelemetryUpdate,
265) -> Result<TelemetryAccepted, TelemetryError> {
266 update_telemetry_with_clock(paths, update, &SystemTelemetryClock)
267}
268
269pub fn update_telemetry_with_clock(
272 paths: &RunPaths,
273 update: &TelemetryUpdate,
274 clock: &impl TelemetryClock,
275) -> Result<TelemetryAccepted, TelemetryError> {
276 validate_update_shape(update)?;
277 let normalized = serde_json::to_vec(update).map_err(TelemetryError::InvalidRequest)?;
278 ensure_size("normalized request", normalized.len())?;
279 if update.run_id != paths.run_id {
280 return Err(TelemetryError::RunMismatch {
281 expected: paths.run_id.clone(),
282 found: update.run_id.clone(),
283 });
284 }
285
286 let guard = RunLock::acquire_existing(&paths.lock())?;
289 if !canonical_projections_current(paths)? {
293 return Err(TelemetryError::RunStateNotCurrent);
294 }
295 let node = match crate::read_node_opt(paths, &update.node_id)? {
296 Some(node) => node,
297 None => {
298 return Err(TelemetryError::NodeNotFound {
299 node_id: update.node_id.clone(),
300 })
301 }
302 };
303 if node.status.is_terminal() {
304 return Err(TelemetryError::TerminalNode {
305 node_id: node.node_id,
306 status: node.status,
307 });
308 }
309 if update.attempt != node.retry_attempts {
310 return Err(TelemetryError::AttemptMismatch {
311 expected: node.retry_attempts,
312 found: update.attempt,
313 });
314 }
315
316 let received_at = clock.now();
317 let expires_at = received_at
318 .checked_add_signed(Duration::seconds(TELEMETRY_FRESHNESS_SECS))
319 .ok_or(TelemetryError::ClockOverflow)?;
320 let path = checked_telemetry_file(paths, &update.node_id)?;
321 let prior = match read_stored(&path)? {
322 StoredRead::Valid(sample) => Some(sample),
323 StoredRead::Absent | StoredRead::Corrupt => None,
324 };
325 let state_since = prior
326 .filter(|sample| {
327 valid_stored_shape(sample, paths, &update.node_id)
328 && sample.attempt == update.attempt
329 && sample.state == update.state
330 })
331 .map_or(received_at, |sample| sample.state_since);
332 let sample = StoredTelemetrySample {
333 schema_version: TELEMETRY_SCHEMA_VERSION,
334 protocol_version: TELEMETRY_PROTOCOL_VERSION,
335 run_id: update.run_id.clone(),
336 node_id: update.node_id.clone(),
337 attempt: update.attempt,
338 state: update.state,
339 active_tool_count: update.active_tool_count,
340 tool_name: update.tool_name.clone(),
341 state_since,
342 received_at,
343 expires_at,
344 };
345 let stored = serde_json::to_vec(&sample).map_err(TelemetryError::InvalidRequest)?;
346 ensure_size("stored sample", stored.len())?;
347 write_atomic(&path, &stored)?;
348 drop(guard);
349
350 Ok(TelemetryAccepted {
351 accepted: true,
352 run_id: update.run_id.clone(),
353 node_id: update.node_id.clone(),
354 attempt: update.attempt,
355 received_at,
356 expires_at,
357 })
358}
359
360pub fn read_telemetry(paths: &RunPaths, node_id: &NodeId) -> Result<TelemetryView, TelemetryError> {
362 read_telemetry_with_clock(paths, node_id, &SystemTelemetryClock)
363}
364
365pub fn read_telemetry_with_clock(
368 paths: &RunPaths,
369 node_id: &NodeId,
370 clock: &impl TelemetryClock,
371) -> Result<TelemetryView, TelemetryError> {
372 let result = RunLock::<Shared>::with_shared_lock(&paths.lock(), || {
373 if !canonical_projections_current(paths)? {
374 return Ok(None);
375 }
376 read_telemetry_locked(paths, node_id, clock.now()).map(Some)
377 })?;
378 result.ok_or(TelemetryError::RunStateNotCurrent)
379}
380
381pub fn read_all_telemetry(
392 paths: &RunPaths,
393) -> Result<Vec<(NodeId, TelemetryView)>, TelemetryError> {
394 read_all_telemetry_with_clock(paths, &SystemTelemetryClock)
395}
396
397pub fn read_all_telemetry_with_clock(
399 paths: &RunPaths,
400 clock: &impl TelemetryClock,
401) -> Result<Vec<(NodeId, TelemetryView)>, TelemetryError> {
402 let result = RunLock::<Shared>::with_shared_lock(&paths.lock(), || {
403 let node_ids = projected_node_ids(paths)?;
404 if !canonical_projections_current(paths)? {
405 return Ok(None);
406 }
407 let now = clock.now();
408 let mut rows = Vec::with_capacity(node_ids.len());
409 for node_id in node_ids {
410 let view = read_telemetry_locked(paths, &node_id, now)?;
411 rows.push((node_id, view));
412 }
413 Ok(Some(rows))
414 })?;
415 result.ok_or(TelemetryError::RunStateNotCurrent)
416}
417
418fn read_telemetry_locked(
419 paths: &RunPaths,
420 node_id: &NodeId,
421 now: DateTime<Utc>,
422) -> Result<TelemetryView, Error> {
423 let node = read_node(paths, node_id)?;
424 let path = checked_telemetry_file(paths, node_id)?;
425 let sample = match read_stored(&path)? {
426 StoredRead::Absent => return Ok(TelemetryView::bare(TelemetrySampleStatus::Absent)),
427 StoredRead::Corrupt => return Ok(TelemetryView::bare(TelemetrySampleStatus::Invalid)),
428 StoredRead::Valid(sample) => sample,
429 };
430 if !valid_stored_shape(&sample, paths, node_id) {
431 return Ok(TelemetryView::bare(TelemetrySampleStatus::Invalid));
432 }
433 if sample.attempt != node.retry_attempts {
434 return Ok(TelemetryView::bare(TelemetrySampleStatus::Absent));
435 }
436 let clock_bad = now < sample.received_at || now < sample.state_since;
437 let status = if clock_bad {
438 TelemetrySampleStatus::ClockUnreliable
439 } else if now >= sample.expires_at {
440 TelemetrySampleStatus::Stale
441 } else {
442 TelemetrySampleStatus::Current
443 };
444 Ok(TelemetryView {
445 sample: status,
446 state: Some(sample.state),
447 age_ms: (!clock_bad).then(|| (now - sample.received_at).num_milliseconds()),
448 state_elapsed_ms: (!clock_bad).then(|| (now - sample.state_since).num_milliseconds()),
449 attempt: Some(sample.attempt),
450 active_tool_count: sample.active_tool_count,
451 tool_name: sample.tool_name,
452 })
453}
454
455fn canonical_projections_current(paths: &RunPaths) -> Result<bool, Error> {
456 let Some(manifest) = crate::read_manifest_opt(paths)? else {
457 return Ok(false);
458 };
459 let events = paths.checked_events()?;
460 Ok(manifest.applied_seq == crate::recover_last_seq(&events)?)
461}
462
463fn validate_update_shape(update: &TelemetryUpdate) -> Result<(), TelemetryError> {
464 if update.schema_version != TELEMETRY_SCHEMA_VERSION {
465 return Err(TelemetryError::UnsupportedSchema {
466 found: update.schema_version,
467 });
468 }
469 if update.protocol_version != TELEMETRY_PROTOCOL_VERSION {
470 return Err(TelemetryError::UnsupportedProtocol {
471 found: update.protocol_version,
472 });
473 }
474 validate_metadata(
475 update.state,
476 update.active_tool_count,
477 update.tool_name.as_deref(),
478 )
479}
480
481fn validate_metadata(
482 state: TelemetryState,
483 count: Option<u8>,
484 name: Option<&str>,
485) -> Result<(), TelemetryError> {
486 if state != TelemetryState::ToolRunning && (count.is_some() || name.is_some()) {
487 return Err(TelemetryError::InvalidMetadata(
488 "tool metadata is allowed only for tool_running",
489 ));
490 }
491 if let Some(count) = count {
492 if !(1..=32).contains(&count) {
493 return Err(TelemetryError::InvalidMetadata(
494 "active_tool_count must be between 1 and 32",
495 ));
496 }
497 }
498 if let Some(name) = name {
499 if count != Some(1) {
500 return Err(TelemetryError::InvalidMetadata(
501 "tool_name requires active_tool_count=1",
502 ));
503 }
504 if !valid_tool_name(name) {
505 return Err(TelemetryError::InvalidMetadata(
506 "tool_name must match ^[A-Za-z0-9_.:-]{1,64}$",
507 ));
508 }
509 }
510 Ok(())
511}
512
513fn valid_tool_name(name: &str) -> bool {
514 (1..=64).contains(&name.len())
515 && name
516 .bytes()
517 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'_' | b'.' | b':' | b'-'))
518}
519
520fn valid_stored_shape(sample: &StoredTelemetrySample, paths: &RunPaths, node_id: &NodeId) -> bool {
521 sample.schema_version == TELEMETRY_SCHEMA_VERSION
522 && sample.protocol_version == TELEMETRY_PROTOCOL_VERSION
523 && sample.run_id == paths.run_id
524 && sample.node_id == *node_id
525 && sample
526 .received_at
527 .checked_add_signed(Duration::seconds(TELEMETRY_FRESHNESS_SECS))
528 .is_some_and(|expected| sample.expires_at == expected)
529 && validate_metadata(
530 sample.state,
531 sample.active_tool_count,
532 sample.tool_name.as_deref(),
533 )
534 .is_ok()
535}
536
537fn ensure_size(what: &'static str, bytes: usize) -> Result<(), TelemetryError> {
538 if bytes <= TELEMETRY_MAX_BYTES {
539 Ok(())
540 } else {
541 Err(TelemetryError::TooLarge { what, bytes })
542 }
543}
544
545fn projected_node_ids(paths: &RunPaths) -> Result<Vec<NodeId>, Error> {
546 paths.guard_root()?;
547 let dir = paths.nodes_dir();
548 reject_symlink(&dir, || Error::SymlinkSubdir {
549 name: "nodes",
550 path: dir.clone(),
551 })?;
552 let entries = match std::fs::read_dir(&dir) {
553 Ok(entries) => entries,
554 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
555 Err(error) => return Err(Error::io(&dir, error)),
556 };
557 let mut ids = Vec::new();
558 for entry in entries {
559 let entry = entry.map_err(|error| Error::io(&dir, error))?;
560 let path = entry.path();
561 if path.extension().and_then(|value| value.to_str()) != Some("json") {
562 continue;
563 }
564 let Some(stem) = path.file_stem().and_then(|value| value.to_str()) else {
565 continue;
566 };
567 if let Ok(node_id) = NodeId::parse_str(stem) {
568 ids.push(node_id);
569 }
570 }
571 ids.sort_by(|left, right| left.as_str().cmp(right.as_str()));
572 Ok(ids)
573}
574
575fn checked_telemetry_file(paths: &RunPaths, node_id: &NodeId) -> Result<PathBuf, Error> {
576 paths.guard_root()?;
577 let dir = paths.root.join("telemetry");
578 reject_symlink(&dir, || Error::SymlinkSubdir {
579 name: "telemetry",
580 path: dir.clone(),
581 })?;
582 let path = dir.join(format!("{}.json", node_id.as_str()));
583 reject_symlink(&path, || Error::SymlinkStateFile {
584 name: "telemetry",
585 path: path.clone(),
586 })?;
587 Ok(path)
588}
589
590enum StoredRead {
591 Absent,
592 Corrupt,
593 Valid(StoredTelemetrySample),
594}
595
596fn read_stored(path: &Path) -> Result<StoredRead, Error> {
597 let mut options = std::fs::OpenOptions::new();
598 options.read(true);
599 nofollow(&mut options);
600 let mut file = match options.open(path) {
601 Ok(file) => file,
602 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
603 return Ok(StoredRead::Absent)
604 }
605 Err(error) => return Err(Error::io(path, error)),
606 };
607 let mut bytes = Vec::new();
608 file.by_ref()
609 .take((TELEMETRY_MAX_BYTES + 1) as u64)
610 .read_to_end(&mut bytes)
611 .map_err(|error| Error::io(path, error))?;
612 if bytes.len() > TELEMETRY_MAX_BYTES {
613 return Ok(StoredRead::Corrupt);
614 }
615 Ok(match serde_json::from_slice(&bytes) {
616 Ok(sample) => StoredRead::Valid(sample),
617 Err(_) => StoredRead::Corrupt,
618 })
619}