Skip to main content

omni_dev/sessions/
stream.rs

1//! The stream-json tracker (Feed 4): the pure state machine behind
2//! `omni-dev claude-wrap`, turning Claude's `--output-format stream-json` stdio
3//! into [`ObserveRequest`]s for the [`SessionsRegistry`].
4//!
5//! Unlike the other three feeds this one is **authoritative**. Hooks and the
6//! transcript watcher observe a session from the outside and infer what it is
7//! doing; the wrapper sits *in* the stream the Claude VS Code extension itself
8//! reads, so it sees the exact protocol events — including `can_use_tool`, the
9//! permission prompt that never reaches a transcript and is therefore invisible
10//! to Feeds 1–3 unless the user has the `Notification` hook installed. See
11//! ADR-0057.
12//!
13//! Both directions matter. A permission request travels **CLI → editor** on the
14//! child's stdout as a `control_request`; its resolution travels **editor → CLI**
15//! on the child's stdin as a `control_response`. Correlating the two by
16//! `request_id` is what makes `waiting_for_permission` exact rather than a guess,
17//! so [`StreamTracker::observe_line`] takes the [`Direction`] a line was seen in.
18//!
19//! Nothing here does I/O and nothing here retains conversation content: only the
20//! message `type`/`subtype`, the identity fields (`session_id`, `cwd`, `model`)
21//! and outstanding permission ids are ever read out of a line. Every parse is
22//! best-effort — an unparseable or unrecognized line is ignored, never fatal,
23//! because the wrapper must fail open (see [`crate::cli::claude_wrap`]).
24//!
25//! [`SessionsRegistry`]: super::SessionsRegistry
26
27use std::collections::HashSet;
28use std::path::PathBuf;
29
30use serde::Deserialize;
31
32use super::{ObserveRequest, SessionEvent, SessionState};
33
34/// Ceiling on simultaneously-outstanding permission requests tracked at once.
35///
36/// Claude asks about one tool at a time in practice, so this only bounds memory
37/// against a malformed or adversarial stream; ids past the cap are dropped,
38/// which can only ever make the tracker return to `working` early.
39const MAX_PENDING_PERMISSIONS: usize = 64;
40
41/// Which side of the wrapped process's stdio a line was observed on.
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum Direction {
44    /// A line the wrapped `claude` wrote to its **stdout** (CLI → editor).
45    FromClaude,
46    /// A line the editor wrote to the wrapped `claude`'s **stdin** (editor → CLI).
47    ToClaude,
48}
49
50/// The subset of a stream-json line this tracker reads.
51///
52/// Deliberately minimal and fully optional: the stream schema is Claude Code's
53/// internal protocol, so anything unrecognized must deserialize successfully and
54/// be ignored rather than break the feed. Conversation content (`message`,
55/// tool inputs, results) is never named here, so it is never even materialized.
56#[derive(Debug, Default, Deserialize)]
57struct StreamLine {
58    /// The message kind: `system`, `assistant`, `user`, `result`, `stream_event`,
59    /// `control_request`, `control_response`, …
60    #[serde(rename = "type", default)]
61    kind: Option<String>,
62    /// The `system` message's discriminator, notably `init`.
63    #[serde(default)]
64    subtype: Option<String>,
65    /// The Claude session id, carried on most messages.
66    #[serde(default)]
67    session_id: Option<String>,
68    /// The session's working directory, carried on `system`/`init`.
69    #[serde(default)]
70    cwd: Option<PathBuf>,
71    /// The model id, carried on `system`/`init`.
72    #[serde(default)]
73    model: Option<String>,
74    /// A control message's correlation id, when carried at the top level.
75    #[serde(default)]
76    request_id: Option<String>,
77    /// A `control_request`'s body.
78    #[serde(default)]
79    request: Option<ControlBody>,
80    /// A `control_response`'s body.
81    #[serde(default)]
82    response: Option<ControlBody>,
83}
84
85/// The shared shape of a `control_request` / `control_response` body.
86#[derive(Debug, Default, Deserialize)]
87struct ControlBody {
88    /// The control kind, e.g. `can_use_tool`, `initialize`, `hook_callback`.
89    #[serde(default)]
90    subtype: Option<String>,
91    /// The correlation id, when carried inside the body rather than at the top
92    /// level (a `control_response` echoes it here).
93    #[serde(default)]
94    request_id: Option<String>,
95}
96
97/// The authoritative session-state machine over one wrapped `claude` process.
98///
99/// Feed it every line of both stdio directions; it returns an [`ObserveRequest`]
100/// exactly when the session's effective state *changes*, so a long turn costs one
101/// report at its start and one at its end rather than one per streamed token.
102#[derive(Debug)]
103pub struct StreamTracker {
104    /// The session id, learned from the first line that carries one.
105    session_id: Option<String>,
106    /// The session's working directory, learned from `system`/`init`.
107    cwd: Option<PathBuf>,
108    /// The model id, learned from `system`/`init`.
109    model: Option<String>,
110    /// The state implied by the most recent content message, before the
111    /// permission overlay is applied.
112    base: SessionState,
113    /// The `request_id`s of permission prompts asked but not yet answered.
114    pending: HashSet<String>,
115    /// The state most recently returned to the caller, for change detection.
116    reported: Option<SessionState>,
117}
118
119impl StreamTracker {
120    /// Creates a tracker for one wrapped process, before any line is seen.
121    #[must_use]
122    pub fn new() -> Self {
123        Self {
124            session_id: None,
125            cwd: None,
126            model: None,
127            // A `claude` that has started but has not been prompted is sitting at
128            // the prompt: idle, not "starting". `Starting` stays the hook feed's
129            // (very brief) `SessionStart` state, so both feeds agree that a tab
130            // the user opened and has not typed into is not doing work.
131            base: SessionState::Idle,
132            pending: HashSet::new(),
133            reported: None,
134        }
135    }
136
137    /// The session id, once a line has carried one.
138    #[must_use]
139    pub fn session_id(&self) -> Option<&str> {
140        self.session_id.as_deref()
141    }
142
143    /// Feeds one stdio line and returns a sighting when the effective state
144    /// changed as a result.
145    ///
146    /// Returns `None` for every line that is unparseable, unrecognized, seen
147    /// before the session id is known, or that leaves the state unchanged.
148    pub fn observe_line(&mut self, direction: Direction, line: &str) -> Option<ObserveRequest> {
149        let line = line.trim();
150        if line.is_empty() {
151            return None;
152        }
153        let parsed: StreamLine = serde_json::from_str(line).ok()?;
154        self.absorb_identity(&parsed);
155        self.apply(direction, &parsed);
156        self.emit_if_changed()
157    }
158
159    /// Re-reports the current state, so a session that has been silent for a
160    /// while does not age out of the registry on its TTL.
161    ///
162    /// The wrapper lives exactly as long as the `claude` process does, so this
163    /// is real liveness rather than the activity-based approximation the hook and
164    /// transcript feeds are limited to.
165    #[must_use]
166    pub fn keepalive(&self) -> Option<ObserveRequest> {
167        self.request(self.state())
168    }
169
170    /// Records the identity fields carried on a line, never overwriting a value
171    /// already learned with a later absent one.
172    fn absorb_identity(&mut self, parsed: &StreamLine) {
173        if self.session_id.is_none() {
174            if let Some(id) = parsed.session_id.as_deref() {
175                if !id.trim().is_empty() {
176                    self.session_id = Some(id.to_string());
177                }
178            }
179        }
180        if self.cwd.is_none() {
181            self.cwd.clone_from(&parsed.cwd);
182        }
183        if self.model.is_none() {
184            self.model.clone_from(&parsed.model);
185        }
186    }
187
188    /// Applies a line's state effect: content messages move [`Self::base`],
189    /// control messages open and close permission prompts.
190    fn apply(&mut self, direction: Direction, parsed: &StreamLine) {
191        match parsed.kind.as_deref() {
192            // The session announced itself but has not been prompted yet.
193            Some("system") if parsed.subtype.as_deref() == Some("init") => {
194                self.base = SessionState::Idle;
195            }
196            // A replayed user prompt, a streamed assistant reply, or a tool
197            // result: the turn is running.
198            Some("assistant" | "user" | "stream_event") => self.base = SessionState::Working,
199            // The turn finished. Also the drift backstop: if the stream ever
200            // stops answering a permission request in a shape this tracker
201            // recognizes, a completed turn unwedges it rather than pinning the
202            // session on `waiting_for_permission` forever.
203            Some("result") => {
204                self.base = SessionState::Idle;
205                self.pending.clear();
206            }
207            Some("control_request") if direction == Direction::FromClaude => {
208                self.open_permission(parsed);
209            }
210            Some("control_response") if direction == Direction::ToClaude => {
211                self.close_permission(parsed);
212            }
213            _ => {}
214        }
215    }
216
217    /// Records a `can_use_tool` request as outstanding; every other control
218    /// subtype (`initialize`, `hook_callback`, `mcp_message`, …) carries no
219    /// state signal and is ignored.
220    fn open_permission(&mut self, parsed: &StreamLine) {
221        let body = parsed.request.as_ref();
222        if body.and_then(|b| b.subtype.as_deref()) != Some("can_use_tool") {
223            return;
224        }
225        let Some(id) = correlation_id(parsed, body) else {
226            return;
227        };
228        if self.pending.len() < MAX_PENDING_PERMISSIONS {
229            self.pending.insert(id);
230        }
231    }
232
233    /// Clears the permission request a `control_response` answers.
234    fn close_permission(&mut self, parsed: &StreamLine) {
235        let body = parsed.response.as_ref();
236        if let Some(id) = correlation_id(parsed, body) {
237            self.pending.remove(&id);
238        }
239    }
240
241    /// The effective state: an outstanding permission prompt outranks whatever
242    /// the content messages last implied, so a reply that keeps streaming while
243    /// the user is being asked to approve a tool cannot mask the prompt.
244    fn state(&self) -> SessionState {
245        if self.pending.is_empty() {
246            self.base
247        } else {
248            SessionState::WaitingForPermission
249        }
250    }
251
252    /// Returns a sighting when the effective state differs from the last one
253    /// reported, and records it as reported.
254    fn emit_if_changed(&mut self) -> Option<ObserveRequest> {
255        let state = self.state();
256        if self.reported == Some(state) {
257            return None;
258        }
259        let request = self.request(state)?;
260        self.reported = Some(state);
261        Some(request)
262    }
263
264    /// Builds the sighting for `state`, or `None` while the session id is still
265    /// unknown (nothing can be keyed without it).
266    fn request(&self, state: SessionState) -> Option<ObserveRequest> {
267        Some(ObserveRequest {
268            session_id: self.session_id.clone()?,
269            cwd: self.cwd.clone(),
270            transcript_path: None,
271            event: SessionEvent::StreamState(state),
272            repo: None,
273            model: self.model.clone(),
274        })
275    }
276}
277
278impl Default for StreamTracker {
279    fn default() -> Self {
280        Self::new()
281    }
282}
283
284/// A control message's correlation id, from the body when present (where a
285/// `control_response` echoes it) and otherwise from the top level.
286fn correlation_id(parsed: &StreamLine, body: Option<&ControlBody>) -> Option<String> {
287    body.and_then(|b| b.request_id.clone())
288        .or_else(|| parsed.request_id.clone())
289}
290
291#[cfg(test)]
292#[allow(clippy::unwrap_used, clippy::expect_used)]
293mod tests {
294    use super::*;
295
296    const INIT: &str = r#"{"type":"system","subtype":"init","session_id":"sess-1","cwd":"/w/repo","model":"claude-opus-5","tools":["Read"]}"#;
297
298    fn tracker_after_init() -> StreamTracker {
299        let mut tracker = StreamTracker::new();
300        let first = tracker
301            .observe_line(Direction::FromClaude, INIT)
302            .expect("init announces the session");
303        assert_eq!(first.session_id, "sess-1");
304        tracker
305    }
306
307    fn state_of(request: &ObserveRequest) -> SessionState {
308        match request.event {
309            SessionEvent::StreamState(state) => state,
310            other => panic!("expected a stream state, got {other:?}"),
311        }
312    }
313
314    #[test]
315    fn init_announces_the_session_as_idle_with_its_identity() {
316        let mut tracker = StreamTracker::new();
317        let request = tracker.observe_line(Direction::FromClaude, INIT).unwrap();
318        assert_eq!(request.session_id, "sess-1");
319        assert_eq!(
320            request.cwd.as_deref(),
321            Some(std::path::Path::new("/w/repo"))
322        );
323        assert_eq!(request.model.as_deref(), Some("claude-opus-5"));
324        // A started-but-unprompted session sits at the prompt.
325        assert_eq!(state_of(&request), SessionState::Idle);
326        assert_eq!(tracker.session_id(), Some("sess-1"));
327    }
328
329    #[test]
330    fn nothing_is_reported_before_a_session_id_is_known() {
331        let mut tracker = StreamTracker::new();
332        // A content message with no session id moves the state but cannot be keyed.
333        assert!(tracker
334            .observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
335            .is_none());
336        assert!(tracker.keepalive().is_none());
337        // …and the state it moved to is reported as soon as an id arrives.
338        let request = tracker
339            .observe_line(
340                Direction::FromClaude,
341                r#"{"type":"assistant","session_id":"sess-1"}"#,
342            )
343            .unwrap();
344        assert_eq!(state_of(&request), SessionState::Working);
345    }
346
347    #[test]
348    fn a_turn_reports_working_then_idle_once_each() {
349        let mut tracker = tracker_after_init();
350        let working = tracker
351            .observe_line(
352                Direction::FromClaude,
353                r#"{"type":"user","session_id":"sess-1"}"#,
354            )
355            .unwrap();
356        assert_eq!(state_of(&working), SessionState::Working);
357        // Streaming does not re-report: the state has not changed.
358        assert!(tracker
359            .observe_line(Direction::FromClaude, r#"{"type":"stream_event"}"#)
360            .is_none());
361        assert!(tracker
362            .observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
363            .is_none());
364        let idle = tracker
365            .observe_line(
366                Direction::FromClaude,
367                r#"{"type":"result","subtype":"success"}"#,
368            )
369            .unwrap();
370        assert_eq!(state_of(&idle), SessionState::Idle);
371    }
372
373    #[test]
374    fn a_permission_prompt_reports_waiting_until_it_is_answered() {
375        let mut tracker = tracker_after_init();
376        tracker.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#);
377        let waiting = tracker
378            .observe_line(
379                Direction::FromClaude,
380                r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool","tool_name":"Bash"}}"#,
381            )
382            .unwrap();
383        assert_eq!(state_of(&waiting), SessionState::WaitingForPermission);
384        // Content still streaming while the user is asked must not mask the prompt.
385        assert!(tracker
386            .observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
387            .is_none());
388        // The answer arrives on the *other* direction, echoing the id in its body.
389        let resumed = tracker
390            .observe_line(
391                Direction::ToClaude,
392                r#"{"type":"control_response","response":{"subtype":"success","request_id":"req-1"}}"#,
393            )
394            .unwrap();
395        assert_eq!(state_of(&resumed), SessionState::Working);
396    }
397
398    #[test]
399    fn a_permission_response_is_only_honored_from_the_editor() {
400        let mut tracker = tracker_after_init();
401        tracker.observe_line(
402            Direction::FromClaude,
403            r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool"}}"#,
404        );
405        // The same line seen on the wrong direction is not the answer.
406        assert!(tracker
407            .observe_line(
408                Direction::FromClaude,
409                r#"{"type":"control_response","response":{"request_id":"req-1"}}"#,
410            )
411            .is_none());
412        assert_eq!(
413            state_of(&tracker.keepalive().unwrap()),
414            SessionState::WaitingForPermission
415        );
416    }
417
418    #[test]
419    fn other_control_subtypes_carry_no_state_signal() {
420        let mut tracker = tracker_after_init();
421        for subtype in ["initialize", "hook_callback", "mcp_message", "interrupt"] {
422            let line = format!(
423                r#"{{"type":"control_request","request_id":"c","request":{{"subtype":"{subtype}"}}}}"#
424            );
425            assert!(tracker.observe_line(Direction::FromClaude, &line).is_none());
426        }
427        assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
428    }
429
430    #[test]
431    fn a_finished_turn_unwedges_a_stranded_permission() {
432        let mut tracker = tracker_after_init();
433        tracker.observe_line(
434            Direction::FromClaude,
435            r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool"}}"#,
436        );
437        // No matching response ever arrives (protocol drift); `result` still ends
438        // the turn rather than pinning the session on `waiting_for_permission`.
439        let idle = tracker
440            .observe_line(Direction::FromClaude, r#"{"type":"result"}"#)
441            .unwrap();
442        assert_eq!(state_of(&idle), SessionState::Idle);
443    }
444
445    #[test]
446    fn outstanding_permissions_are_capped() {
447        let mut tracker = tracker_after_init();
448        for i in 0..(MAX_PENDING_PERMISSIONS + 10) {
449            let line = format!(
450                r#"{{"type":"control_request","request_id":"req-{i}","request":{{"subtype":"can_use_tool"}}}}"#
451            );
452            tracker.observe_line(Direction::FromClaude, &line);
453        }
454        assert_eq!(tracker.pending.len(), MAX_PENDING_PERMISSIONS);
455    }
456
457    #[test]
458    fn a_blank_session_id_is_not_taken_as_identity() {
459        let mut tracker = StreamTracker::new();
460        assert!(tracker
461            .observe_line(
462                Direction::FromClaude,
463                r#"{"type":"assistant","session_id":"   "}"#,
464            )
465            .is_none());
466        assert_eq!(tracker.session_id(), None);
467        // …and a real id later still lands.
468        tracker.observe_line(Direction::FromClaude, INIT);
469        assert_eq!(tracker.session_id(), Some("sess-1"));
470    }
471
472    #[test]
473    fn control_messages_with_no_correlation_id_are_ignored() {
474        let mut tracker = tracker_after_init();
475        // A permission request that cannot be correlated is not tracked, rather
476        // than pinning the session on a prompt nothing can ever answer.
477        assert!(tracker
478            .observe_line(
479                Direction::FromClaude,
480                r#"{"type":"control_request","request":{"subtype":"can_use_tool"}}"#,
481            )
482            .is_none());
483        assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
484        // Likewise an answer that names no request clears nothing.
485        tracker.observe_line(
486            Direction::FromClaude,
487            r#"{"type":"control_request","request_id":"r1","request":{"subtype":"can_use_tool"}}"#,
488        );
489        assert!(tracker
490            .observe_line(Direction::ToClaude, r#"{"type":"control_response"}"#)
491            .is_none());
492        assert_eq!(
493            state_of(&tracker.keepalive().unwrap()),
494            SessionState::WaitingForPermission
495        );
496    }
497
498    #[test]
499    fn default_matches_a_fresh_tracker() {
500        let tracker = StreamTracker::default();
501        assert_eq!(tracker.session_id(), None);
502        assert!(tracker.keepalive().is_none());
503    }
504
505    #[test]
506    fn unparseable_and_unknown_lines_are_ignored() {
507        let mut tracker = tracker_after_init();
508        for line in [
509            "",
510            "   ",
511            "not json at all",
512            "{",
513            "[]",
514            r#"{"type":"nonsense"}"#,
515            r#"{"no_type":true}"#,
516        ] {
517            assert!(tracker.observe_line(Direction::FromClaude, line).is_none());
518        }
519        assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
520    }
521
522    #[test]
523    fn keepalive_re_reports_the_current_state_without_a_change() {
524        let mut tracker = tracker_after_init();
525        tracker.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#);
526        let first = tracker.keepalive().unwrap();
527        let second = tracker.keepalive().unwrap();
528        assert_eq!(state_of(&first), SessionState::Working);
529        assert_eq!(state_of(&second), SessionState::Working);
530        assert_eq!(first.session_id, "sess-1");
531    }
532}