1use std::collections::HashSet;
28use std::path::PathBuf;
29
30use serde::Deserialize;
31
32use super::{ObserveRequest, SessionEvent, SessionState};
33
34const MAX_PENDING_PERMISSIONS: usize = 64;
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum Direction {
44 FromClaude,
46 ToClaude,
48}
49
50#[derive(Debug, Default, Deserialize)]
57struct StreamLine {
58 #[serde(rename = "type", default)]
61 kind: Option<String>,
62 #[serde(default)]
64 subtype: Option<String>,
65 #[serde(default)]
67 session_id: Option<String>,
68 #[serde(default)]
70 cwd: Option<PathBuf>,
71 #[serde(default)]
73 model: Option<String>,
74 #[serde(default)]
76 request_id: Option<String>,
77 #[serde(default)]
79 request: Option<ControlBody>,
80 #[serde(default)]
82 response: Option<ControlBody>,
83}
84
85#[derive(Debug, Default, Deserialize)]
87struct ControlBody {
88 #[serde(default)]
90 subtype: Option<String>,
91 #[serde(default)]
94 request_id: Option<String>,
95}
96
97#[derive(Debug)]
103pub struct StreamTracker {
104 session_id: Option<String>,
106 cwd: Option<PathBuf>,
108 model: Option<String>,
110 base: SessionState,
113 pending: HashSet<String>,
115 reported: Option<SessionState>,
117}
118
119impl StreamTracker {
120 #[must_use]
122 pub fn new() -> Self {
123 Self {
124 session_id: None,
125 cwd: None,
126 model: None,
127 base: SessionState::Idle,
132 pending: HashSet::new(),
133 reported: None,
134 }
135 }
136
137 #[must_use]
139 pub fn session_id(&self) -> Option<&str> {
140 self.session_id.as_deref()
141 }
142
143 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 #[must_use]
166 pub fn keepalive(&self) -> Option<ObserveRequest> {
167 self.request(self.state())
168 }
169
170 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 fn apply(&mut self, direction: Direction, parsed: &StreamLine) {
191 match parsed.kind.as_deref() {
192 Some("system") if parsed.subtype.as_deref() == Some("init") => {
194 self.base = SessionState::Idle;
195 }
196 Some("assistant" | "user" | "stream_event") => self.base = SessionState::Working,
199 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 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 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 fn state(&self) -> SessionState {
245 if self.pending.is_empty() {
246 self.base
247 } else {
248 SessionState::WaitingForPermission
249 }
250 }
251
252 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 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
284fn 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 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 assert!(tracker
334 .observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
335 .is_none());
336 assert!(tracker.keepalive().is_none());
337 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 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 assert!(tracker
386 .observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
387 .is_none());
388 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 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 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 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 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 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}