Skip to main content

taskfleet_core/
telemetry.rs

1//! Bounded, advisory worker-telemetry samples.
2//!
3//! Telemetry deliberately lives beside the event-sourced run state rather than
4//! inside it. Updates take the ordinary per-run lock only to validate the node's
5//! current attempt and terminal status atomically with replacing one sample;
6//! they never append an event or rewrite a manifest/node projection.
7
8use 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
21/// Wire and stored-sample schema version supported by this build.
22pub const TELEMETRY_SCHEMA_VERSION: u32 = 1;
23/// Worker telemetry protocol version supported by this build.
24pub const TELEMETRY_PROTOCOL_VERSION: u32 = 1;
25/// Maximum raw request, normalized request, and stored-sample size.
26pub const TELEMETRY_MAX_BYTES: usize = 4 * 1024;
27/// How long a received sample remains current.
28pub const TELEMETRY_FRESHNESS_SECS: i64 = 90;
29
30/// Injectable source of wall-clock time used by telemetry updates and reads.
31pub trait TelemetryClock {
32    /// Return the current server time.
33    fn now(&self) -> DateTime<Utc>;
34}
35
36/// Production telemetry clock.
37#[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/// The four last-told worker activity states.
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
48#[serde(rename_all = "snake_case")]
49pub enum TelemetryState {
50    /// The harness reported an active agent turn.
51    AgentActive,
52    /// One or more tool executions are open.
53    ToolRunning,
54    /// Automatic agent work was reported settled; this is not completion.
55    Settled,
56    /// Session shutdown was reported; this is not completion.
57    Shutdown,
58}
59
60/// Strict v1 telemetry update request.
61#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62#[serde(deny_unknown_fields)]
63pub struct TelemetryUpdate {
64    /// Request schema version; must be [`TELEMETRY_SCHEMA_VERSION`].
65    pub schema_version: u32,
66    /// Protocol version; must be [`TELEMETRY_PROTOCOL_VERSION`].
67    pub protocol_version: u32,
68    /// Exact run identity.
69    pub run_id: RunId,
70    /// Exact node identity.
71    pub node_id: NodeId,
72    /// Absolute current attempt (`Node::retry_attempts`).
73    pub attempt: u32,
74    /// Last-told activity.
75    pub state: TelemetryState,
76    /// Number of active tools, only valid for `tool_running` (1–32).
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub active_tool_count: Option<u8>,
79    /// Sanitized single active tool name, allowed only when count is exactly 1.
80    #[serde(default, skip_serializing_if = "Option::is_none")]
81    pub tool_name: Option<String>,
82}
83
84/// A successfully accepted telemetry update.
85#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
86pub struct TelemetryAccepted {
87    /// The update was stored.
88    pub accepted: bool,
89    /// Run identity.
90    pub run_id: RunId,
91    /// Node identity.
92    pub node_id: NodeId,
93    /// Current attempt accepted.
94    pub attempt: u32,
95    /// Server receive timestamp.
96    pub received_at: DateTime<Utc>,
97    /// Server-computed freshness boundary.
98    pub expires_at: DateTime<Utc>,
99}
100
101/// Validation or state error from a telemetry update.
102#[derive(Debug, thiserror::Error)]
103pub enum TelemetryError {
104    /// Core run-state or filesystem operation failed.
105    #[error(transparent)]
106    Core(#[from] Error),
107    /// Input was not strict, valid JSON for [`TelemetryUpdate`].
108    #[error("invalid telemetry request: {0}")]
109    InvalidRequest(serde_json::Error),
110    /// Raw, normalized, or stored data exceeded the fixed bound.
111    #[error("telemetry {what} exceeds {TELEMETRY_MAX_BYTES} bytes (got {bytes})")]
112    TooLarge {
113        /// Which representation exceeded the bound.
114        what: &'static str,
115        /// Observed byte count.
116        bytes: usize,
117    },
118    /// A request used an unsupported schema version.
119    #[error("unsupported telemetry schema_version {found}; expected {TELEMETRY_SCHEMA_VERSION}")]
120    UnsupportedSchema {
121        /// Rejected version.
122        found: u32,
123    },
124    /// A request used an unsupported protocol version.
125    #[error(
126        "unsupported telemetry protocol_version {found}; expected {TELEMETRY_PROTOCOL_VERSION}"
127    )]
128    UnsupportedProtocol {
129        /// Rejected version.
130        found: u32,
131    },
132    /// Tool metadata did not match the activity state or sanitation rules.
133    #[error("invalid telemetry tool metadata: {0}")]
134    InvalidMetadata(&'static str),
135    /// The request run id does not match the run directory.
136    #[error("telemetry run_id {found} does not match current run {expected}")]
137    RunMismatch {
138        /// Run directory identity.
139        expected: RunId,
140        /// Request identity.
141        found: RunId,
142    },
143    /// Canonical projections and the event log are not synchronized, so exact
144    /// attempt/terminal validation cannot be proven.
145    #[error("canonical run state is not synchronized with the event log")]
146    RunStateNotCurrent,
147    /// The named node does not exist.
148    #[error("no node {node_id} in this run")]
149    NodeNotFound {
150        /// Missing node.
151        node_id: NodeId,
152    },
153    /// The node is terminal and cannot accept telemetry.
154    #[error("node {node_id} is terminal ({status:?})")]
155    TerminalNode {
156        /// Terminal node.
157        node_id: NodeId,
158        /// Current terminal status.
159        status: Status,
160    },
161    /// The supplied attempt is not the exact current attempt.
162    #[error("telemetry attempt {found} does not match current attempt {expected}")]
163    AttemptMismatch {
164        /// Current node attempt.
165        expected: u32,
166        /// Request attempt.
167        found: u32,
168    },
169    /// The server clock could not represent the expiry timestamp.
170    #[error("server clock cannot represent telemetry expiry")]
171    ClockOverflow,
172}
173
174/// Classification of a telemetry sample for the node's current attempt.
175#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(rename_all = "snake_case")]
177pub enum TelemetrySampleStatus {
178    /// No sample exists, or the stored sample belongs to an older attempt.
179    Absent,
180    /// The sample has not reached its expiry boundary.
181    Current,
182    /// The sample has reached its expiry boundary.
183    Stale,
184    /// The read clock is behind the server receive timestamp.
185    ClockUnreliable,
186    /// A sample file exists but is corrupt or violates the strict stored schema.
187    Invalid,
188}
189
190/// Read view of one node's advisory telemetry.
191#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192pub struct TelemetryView {
193    /// Freshness/corruption classification.
194    pub sample: TelemetrySampleStatus,
195    /// Last-told state when a valid current-attempt sample exists.
196    #[serde(skip_serializing_if = "Option::is_none")]
197    pub state: Option<TelemetryState>,
198    /// Age in milliseconds, unavailable when the clock is unreliable.
199    #[serde(skip_serializing_if = "Option::is_none")]
200    pub age_ms: Option<i64>,
201    /// Time since this state was first received, unavailable with a bad clock.
202    #[serde(skip_serializing_if = "Option::is_none")]
203    pub state_elapsed_ms: Option<i64>,
204    /// Sample attempt.
205    #[serde(skip_serializing_if = "Option::is_none")]
206    pub attempt: Option<u32>,
207    /// Bounded active-tool count.
208    #[serde(skip_serializing_if = "Option::is_none")]
209    pub active_tool_count: Option<u8>,
210    /// Sanitized single-tool name.
211    #[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/// Strict stored sample. Private so callers cannot bypass update validation.
230#[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
248/// Parse a strict JSON request, rejecting unknown fields and every representation
249/// larger than 4 KiB.
250pub 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
260/// Validate and atomically replace the node's one advisory sample using the
261/// production clock.
262pub fn update_telemetry(
263    paths: &RunPaths,
264    update: &TelemetryUpdate,
265) -> Result<TelemetryAccepted, TelemetryError> {
266    update_telemetry_with_clock(paths, update, &SystemTelemetryClock)
267}
268
269/// Validate and atomically replace the node's one advisory sample using an
270/// injected clock.
271pub 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    // Unlike the canonical creating lock path, telemetry must never resurrect
287    // a deleted/unknown run merely by trying to validate it.
288    let guard = RunLock::acquire_existing(&paths.lock())?;
289    // Telemetry must neither observe stale canonical state nor heal it by
290    // advancing applied_seq. Fail closed and let an ordinary event writer own
291    // crash-tail recovery.
292    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
360/// Read and classify a sample using the production clock.
361pub fn read_telemetry(paths: &RunPaths, node_id: &NodeId) -> Result<TelemetryView, TelemetryError> {
362    read_telemetry_with_clock(paths, node_id, &SystemTelemetryClock)
363}
364
365/// Read and classify a sample under the run's shared lock using an injected
366/// clock. An old-attempt sample is intentionally rendered as absent.
367pub 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
381/// Read every projected node's telemetry under one shared lock and one
382/// canonical-currency check. This is the bounded read-surface API for callers
383/// that need per-run rows or counts; it avoids recursively locking and rescanning
384/// the event log once per node.
385///
386/// A crash-tailed canonical projection cannot prove current attempts and
387/// returns [`TelemetryError::RunStateNotCurrent`]. Malformed bounded sample
388/// contents still classify as `invalid`; canonical projection, path-integrity,
389/// and operational I/O errors propagate instead of masquerading as sample
390/// corruption.
391pub fn read_all_telemetry(
392    paths: &RunPaths,
393) -> Result<Vec<(NodeId, TelemetryView)>, TelemetryError> {
394    read_all_telemetry_with_clock(paths, &SystemTelemetryClock)
395}
396
397/// [`read_all_telemetry`] with an injected read clock.
398pub 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}