Skip to main content

codex_helper_core/fleet/
observer.rs

1use std::collections::HashMap;
2
3use crate::dashboard_core::build_dashboard_snapshot;
4use crate::dashboard_core::snapshot::DashboardSnapshot;
5use crate::state::{ActiveRequest, FinishedRequest, ProxyState, SessionIdentityCard};
6
7use super::model::{
8    FleetConfidence, FleetEvidence, FleetEvidenceSource, FleetGraphStatus, FleetNodeHealth,
9    FleetNodeKind, FleetNodeSnapshot, FleetProcessSummary, FleetSnapshot, FleetTopology,
10    FleetUsageSummary, FleetWorkUnit, FleetWorkUnitKind, FleetWorkUnitState, now_ms,
11};
12use super::process_scan::scan_codex_processes;
13
14pub async fn build_local_fleet_snapshot(
15    state: &ProxyState,
16    service_name: &str,
17    node_id: impl Into<String>,
18    label: impl Into<String>,
19) -> FleetSnapshot {
20    let refreshed_at_ms = now_ms();
21    let dashboard = build_dashboard_snapshot(
22        state,
23        service_name,
24        crate::state::recent_finished_max(),
25        7,
26        None,
27        None,
28    )
29    .await;
30    let process_scan = scan_codex_processes();
31    let node = build_local_fleet_node_from_parts(LocalFleetNodeParts {
32        node_id: node_id.into(),
33        label: label.into(),
34        refreshed_at_ms,
35        processes: process_scan.summary(),
36        session_cards: &dashboard.session_cards,
37        active: &dashboard.active,
38        recent: &dashboard.recent,
39    });
40
41    FleetSnapshot {
42        api_version: 1,
43        service_name: service_name.to_string(),
44        refreshed_at_ms,
45        nodes: vec![node],
46    }
47}
48
49pub fn build_local_fleet_snapshot_from_parts(
50    service_name: &str,
51    node_id: impl Into<String>,
52    label: impl Into<String>,
53    refreshed_at_ms: u64,
54    session_cards: &[SessionIdentityCard],
55    active: &[ActiveRequest],
56    recent: &[FinishedRequest],
57) -> FleetSnapshot {
58    let node = build_local_fleet_node_from_parts(LocalFleetNodeParts {
59        node_id: node_id.into(),
60        label: label.into(),
61        refreshed_at_ms,
62        processes: FleetProcessSummary::default(),
63        session_cards,
64        active,
65        recent,
66    });
67
68    FleetSnapshot {
69        api_version: 1,
70        service_name: service_name.to_string(),
71        refreshed_at_ms,
72        nodes: vec![node],
73    }
74}
75
76pub fn build_local_fleet_snapshot_from_dashboard(
77    service_name: &str,
78    node_id: impl Into<String>,
79    label: impl Into<String>,
80    dashboard: &DashboardSnapshot,
81) -> FleetSnapshot {
82    let node = build_local_fleet_node_from_parts(LocalFleetNodeParts {
83        node_id: node_id.into(),
84        label: label.into(),
85        refreshed_at_ms: dashboard.refreshed_at_ms,
86        processes: scan_codex_processes().summary(),
87        session_cards: &dashboard.session_cards,
88        active: &dashboard.active,
89        recent: &dashboard.recent,
90    });
91
92    FleetSnapshot {
93        api_version: 1,
94        service_name: service_name.to_string(),
95        refreshed_at_ms: dashboard.refreshed_at_ms,
96        nodes: vec![node],
97    }
98}
99
100struct LocalFleetNodeParts<'a> {
101    node_id: String,
102    label: String,
103    refreshed_at_ms: u64,
104    processes: FleetProcessSummary,
105    session_cards: &'a [SessionIdentityCard],
106    active: &'a [ActiveRequest],
107    recent: &'a [FinishedRequest],
108}
109
110fn build_local_fleet_node_from_parts(parts: LocalFleetNodeParts<'_>) -> FleetNodeSnapshot {
111    let LocalFleetNodeParts {
112        node_id,
113        label,
114        refreshed_at_ms,
115        processes,
116        session_cards,
117        active,
118        recent,
119    } = parts;
120
121    let active_by_session = active
122        .iter()
123        .filter_map(|request| request.session_id.as_deref().map(|sid| (sid, request)))
124        .fold(
125            HashMap::<&str, Vec<&ActiveRequest>>::new(),
126            |mut acc, (sid, request)| {
127                acc.entry(sid).or_default().push(request);
128                acc
129            },
130        );
131    let recent_by_session = recent
132        .iter()
133        .filter_map(|request| request.session_id.as_deref().map(|sid| (sid, request)))
134        .fold(
135            HashMap::<&str, Vec<&FinishedRequest>>::new(),
136            |mut acc, (sid, request)| {
137                acc.entry(sid).or_default().push(request);
138                acc
139            },
140        );
141
142    let mut work_units = session_cards
143        .iter()
144        .enumerate()
145        .map(|(idx, card)| {
146            work_unit_from_session_card(
147                &node_id,
148                idx,
149                card,
150                card.session_id
151                    .as_deref()
152                    .and_then(|sid| active_by_session.get(sid).map(Vec::as_slice)),
153                card.session_id
154                    .as_deref()
155                    .and_then(|sid| recent_by_session.get(sid).map(Vec::as_slice)),
156            )
157        })
158        .collect::<Vec<_>>();
159
160    work_units.sort_by_key(|unit| std::cmp::Reverse(unit.last_activity_ms.unwrap_or(0)));
161
162    FleetNodeSnapshot {
163        node_id,
164        label,
165        kind: FleetNodeKind::Local,
166        health: FleetNodeHealth::Fresh,
167        refreshed_at_ms,
168        stale_since_ms: None,
169        snapshot_age_ms: Some(0),
170        active_endpoint: None,
171        last_error: None,
172        processes,
173        topology: FleetTopology {
174            status: FleetGraphStatus::Unavailable,
175            edges: Vec::new(),
176            note: Some(
177                "subagent graph source unavailable; node-local session rows preserved".into(),
178            ),
179        },
180        work_units,
181    }
182}
183
184fn work_unit_from_session_card(
185    node_id: &str,
186    idx: usize,
187    card: &SessionIdentityCard,
188    active: Option<&[&ActiveRequest]>,
189    recent: Option<&[&FinishedRequest]>,
190) -> FleetWorkUnit {
191    let id = card
192        .session_id
193        .as_ref()
194        .map(|sid| format!("session:{sid}"))
195        .unwrap_or_else(|| format!("unknown-session:{idx}"));
196    let active_started_at_ms = active
197        .and_then(|requests| requests.iter().map(|r| r.started_at_ms).min())
198        .or(card.active_started_at_ms_min);
199    let last_recent = recent.and_then(|requests| {
200        requests
201            .iter()
202            .max_by_key(|request| request.ended_at_ms)
203            .copied()
204    });
205    let last_activity_ms = active_started_at_ms
206        .max(card.last_ended_at_ms)
207        .or(card.last_ended_at_ms)
208        .or(active_started_at_ms);
209    let state = infer_work_unit_state(card, last_recent);
210    let evidence = infer_work_unit_evidence(card, active, recent);
211    let last_error = last_recent
212        .filter(|request| request.status_code >= 400)
213        .map(|request| format!("last status {}", request.status_code));
214
215    FleetWorkUnit {
216        node_id: node_id.to_string(),
217        id,
218        parent_id: None,
219        kind: FleetWorkUnitKind::Root,
220        state,
221        evidence,
222        session_id: card.session_id.clone(),
223        local_thread_id: card.session_id.clone(),
224        task_name: None,
225        cwd: card.cwd.clone(),
226        model: card
227            .effective_model
228            .as_ref()
229            .map(|value| value.value.clone())
230            .or_else(|| card.last_model.clone()),
231        station_name: card
232            .effective_station
233            .as_ref()
234            .map(|value| value.value.clone())
235            .or_else(|| card.last_station_name.clone()),
236        provider_id: card.last_provider_id.clone(),
237        last_status: card.last_status,
238        active_started_at_ms,
239        last_activity_ms,
240        last_error,
241        usage: FleetUsageSummary {
242            last_usage: card.last_usage.clone(),
243            total_usage: card.total_usage.clone(),
244            turns_total: card.turns_total,
245            turns_with_usage: card.turns_with_usage,
246            last_output_tokens_per_second: card.last_output_tokens_per_second,
247            avg_output_tokens_per_second: card.avg_output_tokens_per_second,
248        },
249    }
250}
251
252fn infer_work_unit_state(
253    card: &SessionIdentityCard,
254    last_recent: Option<&FinishedRequest>,
255) -> FleetWorkUnitState {
256    if card.active_count > 0 {
257        return FleetWorkUnitState::Running;
258    }
259    if let Some(status) = card.last_status
260        && status >= 400
261    {
262        return FleetWorkUnitState::Errored;
263    }
264    if let Some(request) = last_recent
265        && request.status_code >= 400
266    {
267        return FleetWorkUnitState::Errored;
268    }
269    if card.last_ended_at_ms.is_some() || card.turns_total.unwrap_or(0) > 0 {
270        return FleetWorkUnitState::Completed;
271    }
272    FleetWorkUnitState::Unknown
273}
274
275fn infer_work_unit_evidence(
276    card: &SessionIdentityCard,
277    active: Option<&[&ActiveRequest]>,
278    recent: Option<&[&FinishedRequest]>,
279) -> FleetEvidence {
280    if card.active_count > 0 || card.active_started_at_ms_min.is_some() {
281        return FleetEvidence::high(FleetEvidenceSource::RuntimeStatus);
282    }
283    if active.is_some_and(|requests| !requests.is_empty()) {
284        return FleetEvidence::high(FleetEvidenceSource::RuntimeStatus);
285    }
286    if recent.is_some_and(|requests| !requests.is_empty()) || card.last_ended_at_ms.is_some() {
287        return FleetEvidence::with_detail(
288            FleetEvidenceSource::RuntimeStatus,
289            FleetConfidence::Medium,
290            "derived from request/session runtime history",
291        );
292    }
293    if card.host_local_transcript_path.is_some() {
294        return FleetEvidence::with_detail(
295            FleetEvidenceSource::SessionLog,
296            FleetConfidence::Medium,
297            "derived from local Codex transcript path",
298        );
299    }
300    FleetEvidence::default()
301}
302
303#[cfg(test)]
304mod tests {
305    use crate::state::{FinishedRequest, RequestObservability, SessionIdentityCard};
306    use crate::usage::UsageMetrics;
307
308    use super::*;
309
310    fn card(session_id: &str) -> SessionIdentityCard {
311        SessionIdentityCard {
312            session_id: Some(session_id.to_string()),
313            last_model: Some("gpt-5".to_string()),
314            total_usage: Some(UsageMetrics {
315                total_tokens: 12,
316                output_tokens: 5,
317                ..UsageMetrics::default()
318            }),
319            avg_output_tokens_per_second: Some(42.0),
320            turns_total: Some(2),
321            ..SessionIdentityCard::default()
322        }
323    }
324
325    fn finished(id: u64, session_id: &str, status_code: u16) -> FinishedRequest {
326        FinishedRequest {
327            id,
328            trace_id: None,
329            session_id: Some(session_id.to_string()),
330            session_identity_source: None,
331            client_name: None,
332            client_addr: None,
333            cwd: None,
334            model: None,
335            reasoning_effort: None,
336            service_tier: None,
337            station_name: None,
338            provider_id: None,
339            upstream_base_url: None,
340            route_decision: None,
341            usage: None,
342            cost: crate::pricing::CostBreakdown::default(),
343            retry: None,
344            provider_signals: Vec::new(),
345            policy_actions: Vec::new(),
346            observability: RequestObservability::default(),
347            service: "codex".to_string(),
348            method: "POST".to_string(),
349            path: "/v1/responses".to_string(),
350            status_code,
351            duration_ms: 10,
352            ttfb_ms: None,
353            streaming: false,
354            ended_at_ms: id,
355        }
356    }
357
358    #[test]
359    fn local_snapshot_preserves_source_confidence_and_usage() {
360        let mut c = card("sid-1");
361        c.active_count = 1;
362        c.active_started_at_ms_min = Some(100);
363
364        let snapshot =
365            build_local_fleet_snapshot_from_parts("codex", "local", "Local", 200, &[c], &[], &[]);
366
367        let node = snapshot.nodes.first().expect("node");
368        let unit = node.work_units.first().expect("unit");
369        assert_eq!(unit.node_id, "local");
370        assert_eq!(unit.state, FleetWorkUnitState::Running);
371        assert_eq!(unit.evidence.source, FleetEvidenceSource::RuntimeStatus);
372        assert_eq!(unit.evidence.confidence, FleetConfidence::High);
373        assert_eq!(unit.usage.avg_output_tokens_per_second, Some(42.0));
374        assert_eq!(
375            unit.usage
376                .total_usage
377                .as_ref()
378                .map(|usage| usage.total_tokens),
379            Some(12)
380        );
381        assert_eq!(node.topology.status, FleetGraphStatus::Unavailable);
382    }
383
384    #[test]
385    fn local_snapshot_marks_recent_error_without_active_request() {
386        let mut c = card("sid-err");
387        c.active_count = 0;
388        c.last_status = Some(500);
389        c.last_ended_at_ms = Some(300);
390        let recent = vec![finished(300, "sid-err", 500)];
391
392        let snapshot = build_local_fleet_snapshot_from_parts(
393            "codex",
394            "local",
395            "Local",
396            400,
397            &[c],
398            &[],
399            &recent,
400        );
401
402        let unit = &snapshot.nodes[0].work_units[0];
403        assert_eq!(unit.state, FleetWorkUnitState::Errored);
404        assert_eq!(unit.last_error.as_deref(), Some("last status 500"));
405        assert_eq!(unit.evidence.confidence, FleetConfidence::Medium);
406    }
407}