use std::collections::HashSet;
use std::path::PathBuf;
use serde::Deserialize;
use super::{ObserveRequest, SessionEvent, SessionState};
const MAX_PENDING_PERMISSIONS: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Direction {
FromClaude,
ToClaude,
}
#[derive(Debug, Default, Deserialize)]
struct StreamLine {
#[serde(rename = "type", default)]
kind: Option<String>,
#[serde(default)]
subtype: Option<String>,
#[serde(default)]
session_id: Option<String>,
#[serde(default)]
cwd: Option<PathBuf>,
#[serde(default)]
model: Option<String>,
#[serde(default)]
request_id: Option<String>,
#[serde(default)]
request: Option<ControlBody>,
#[serde(default)]
response: Option<ControlBody>,
}
#[derive(Debug, Default, Deserialize)]
struct ControlBody {
#[serde(default)]
subtype: Option<String>,
#[serde(default)]
request_id: Option<String>,
}
#[derive(Debug)]
pub struct StreamTracker {
session_id: Option<String>,
cwd: Option<PathBuf>,
model: Option<String>,
base: SessionState,
pending: HashSet<String>,
reported: Option<SessionState>,
}
impl StreamTracker {
#[must_use]
pub fn new() -> Self {
Self {
session_id: None,
cwd: None,
model: None,
base: SessionState::Idle,
pending: HashSet::new(),
reported: None,
}
}
#[must_use]
pub fn session_id(&self) -> Option<&str> {
self.session_id.as_deref()
}
pub fn observe_line(&mut self, direction: Direction, line: &str) -> Option<ObserveRequest> {
let line = line.trim();
if line.is_empty() {
return None;
}
let parsed: StreamLine = serde_json::from_str(line).ok()?;
self.absorb_identity(&parsed);
self.apply(direction, &parsed);
self.emit_if_changed()
}
#[must_use]
pub fn keepalive(&self) -> Option<ObserveRequest> {
self.request(self.state())
}
fn absorb_identity(&mut self, parsed: &StreamLine) {
if self.session_id.is_none() {
if let Some(id) = parsed.session_id.as_deref() {
if !id.trim().is_empty() {
self.session_id = Some(id.to_string());
}
}
}
if self.cwd.is_none() {
self.cwd.clone_from(&parsed.cwd);
}
if self.model.is_none() {
self.model.clone_from(&parsed.model);
}
}
fn apply(&mut self, direction: Direction, parsed: &StreamLine) {
match parsed.kind.as_deref() {
Some("system") if parsed.subtype.as_deref() == Some("init") => {
self.base = SessionState::Idle;
}
Some("assistant" | "user" | "stream_event") => self.base = SessionState::Working,
Some("result") => {
self.base = SessionState::Idle;
self.pending.clear();
}
Some("control_request") if direction == Direction::FromClaude => {
self.open_permission(parsed);
}
Some("control_response") if direction == Direction::ToClaude => {
self.close_permission(parsed);
}
_ => {}
}
}
fn open_permission(&mut self, parsed: &StreamLine) {
let body = parsed.request.as_ref();
if body.and_then(|b| b.subtype.as_deref()) != Some("can_use_tool") {
return;
}
let Some(id) = correlation_id(parsed, body) else {
return;
};
if self.pending.len() < MAX_PENDING_PERMISSIONS {
self.pending.insert(id);
}
}
fn close_permission(&mut self, parsed: &StreamLine) {
let body = parsed.response.as_ref();
if let Some(id) = correlation_id(parsed, body) {
self.pending.remove(&id);
}
}
fn state(&self) -> SessionState {
if self.pending.is_empty() {
self.base
} else {
SessionState::WaitingForPermission
}
}
fn emit_if_changed(&mut self) -> Option<ObserveRequest> {
let state = self.state();
if self.reported == Some(state) {
return None;
}
let request = self.request(state)?;
self.reported = Some(state);
Some(request)
}
fn request(&self, state: SessionState) -> Option<ObserveRequest> {
Some(ObserveRequest {
session_id: self.session_id.clone()?,
cwd: self.cwd.clone(),
transcript_path: None,
event: SessionEvent::StreamState(state),
repo: None,
model: self.model.clone(),
})
}
}
impl Default for StreamTracker {
fn default() -> Self {
Self::new()
}
}
fn correlation_id(parsed: &StreamLine, body: Option<&ControlBody>) -> Option<String> {
body.and_then(|b| b.request_id.clone())
.or_else(|| parsed.request_id.clone())
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
const INIT: &str = r#"{"type":"system","subtype":"init","session_id":"sess-1","cwd":"/w/repo","model":"claude-opus-5","tools":["Read"]}"#;
fn tracker_after_init() -> StreamTracker {
let mut tracker = StreamTracker::new();
let first = tracker
.observe_line(Direction::FromClaude, INIT)
.expect("init announces the session");
assert_eq!(first.session_id, "sess-1");
tracker
}
fn state_of(request: &ObserveRequest) -> SessionState {
match request.event {
SessionEvent::StreamState(state) => state,
other => panic!("expected a stream state, got {other:?}"),
}
}
#[test]
fn init_announces_the_session_as_idle_with_its_identity() {
let mut tracker = StreamTracker::new();
let request = tracker.observe_line(Direction::FromClaude, INIT).unwrap();
assert_eq!(request.session_id, "sess-1");
assert_eq!(
request.cwd.as_deref(),
Some(std::path::Path::new("/w/repo"))
);
assert_eq!(request.model.as_deref(), Some("claude-opus-5"));
assert_eq!(state_of(&request), SessionState::Idle);
assert_eq!(tracker.session_id(), Some("sess-1"));
}
#[test]
fn nothing_is_reported_before_a_session_id_is_known() {
let mut tracker = StreamTracker::new();
assert!(tracker
.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
.is_none());
assert!(tracker.keepalive().is_none());
let request = tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"assistant","session_id":"sess-1"}"#,
)
.unwrap();
assert_eq!(state_of(&request), SessionState::Working);
}
#[test]
fn a_turn_reports_working_then_idle_once_each() {
let mut tracker = tracker_after_init();
let working = tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"user","session_id":"sess-1"}"#,
)
.unwrap();
assert_eq!(state_of(&working), SessionState::Working);
assert!(tracker
.observe_line(Direction::FromClaude, r#"{"type":"stream_event"}"#)
.is_none());
assert!(tracker
.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
.is_none());
let idle = tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"result","subtype":"success"}"#,
)
.unwrap();
assert_eq!(state_of(&idle), SessionState::Idle);
}
#[test]
fn a_permission_prompt_reports_waiting_until_it_is_answered() {
let mut tracker = tracker_after_init();
tracker.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#);
let waiting = tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool","tool_name":"Bash"}}"#,
)
.unwrap();
assert_eq!(state_of(&waiting), SessionState::WaitingForPermission);
assert!(tracker
.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#)
.is_none());
let resumed = tracker
.observe_line(
Direction::ToClaude,
r#"{"type":"control_response","response":{"subtype":"success","request_id":"req-1"}}"#,
)
.unwrap();
assert_eq!(state_of(&resumed), SessionState::Working);
}
#[test]
fn a_permission_response_is_only_honored_from_the_editor() {
let mut tracker = tracker_after_init();
tracker.observe_line(
Direction::FromClaude,
r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool"}}"#,
);
assert!(tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"control_response","response":{"request_id":"req-1"}}"#,
)
.is_none());
assert_eq!(
state_of(&tracker.keepalive().unwrap()),
SessionState::WaitingForPermission
);
}
#[test]
fn other_control_subtypes_carry_no_state_signal() {
let mut tracker = tracker_after_init();
for subtype in ["initialize", "hook_callback", "mcp_message", "interrupt"] {
let line = format!(
r#"{{"type":"control_request","request_id":"c","request":{{"subtype":"{subtype}"}}}}"#
);
assert!(tracker.observe_line(Direction::FromClaude, &line).is_none());
}
assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
}
#[test]
fn a_finished_turn_unwedges_a_stranded_permission() {
let mut tracker = tracker_after_init();
tracker.observe_line(
Direction::FromClaude,
r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool"}}"#,
);
let idle = tracker
.observe_line(Direction::FromClaude, r#"{"type":"result"}"#)
.unwrap();
assert_eq!(state_of(&idle), SessionState::Idle);
}
#[test]
fn outstanding_permissions_are_capped() {
let mut tracker = tracker_after_init();
for i in 0..(MAX_PENDING_PERMISSIONS + 10) {
let line = format!(
r#"{{"type":"control_request","request_id":"req-{i}","request":{{"subtype":"can_use_tool"}}}}"#
);
tracker.observe_line(Direction::FromClaude, &line);
}
assert_eq!(tracker.pending.len(), MAX_PENDING_PERMISSIONS);
}
#[test]
fn a_blank_session_id_is_not_taken_as_identity() {
let mut tracker = StreamTracker::new();
assert!(tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"assistant","session_id":" "}"#,
)
.is_none());
assert_eq!(tracker.session_id(), None);
tracker.observe_line(Direction::FromClaude, INIT);
assert_eq!(tracker.session_id(), Some("sess-1"));
}
#[test]
fn control_messages_with_no_correlation_id_are_ignored() {
let mut tracker = tracker_after_init();
assert!(tracker
.observe_line(
Direction::FromClaude,
r#"{"type":"control_request","request":{"subtype":"can_use_tool"}}"#,
)
.is_none());
assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
tracker.observe_line(
Direction::FromClaude,
r#"{"type":"control_request","request_id":"r1","request":{"subtype":"can_use_tool"}}"#,
);
assert!(tracker
.observe_line(Direction::ToClaude, r#"{"type":"control_response"}"#)
.is_none());
assert_eq!(
state_of(&tracker.keepalive().unwrap()),
SessionState::WaitingForPermission
);
}
#[test]
fn default_matches_a_fresh_tracker() {
let tracker = StreamTracker::default();
assert_eq!(tracker.session_id(), None);
assert!(tracker.keepalive().is_none());
}
#[test]
fn unparseable_and_unknown_lines_are_ignored() {
let mut tracker = tracker_after_init();
for line in [
"",
" ",
"not json at all",
"{",
"[]",
r#"{"type":"nonsense"}"#,
r#"{"no_type":true}"#,
] {
assert!(tracker.observe_line(Direction::FromClaude, line).is_none());
}
assert_eq!(state_of(&tracker.keepalive().unwrap()), SessionState::Idle);
}
#[test]
fn keepalive_re_reports_the_current_state_without_a_change() {
let mut tracker = tracker_after_init();
tracker.observe_line(Direction::FromClaude, r#"{"type":"assistant"}"#);
let first = tracker.keepalive().unwrap();
let second = tracker.keepalive().unwrap();
assert_eq!(state_of(&first), SessionState::Working);
assert_eq!(state_of(&second), SessionState::Working);
assert_eq!(first.session_id, "sess-1");
}
}