1use crate::{normalize_hook_event, HookAdapterError};
2use gate4agent_types::{
3 AdapterId, ProviderEvent, ProviderEventValidationError, PROVIDER_SUBAGENTS_MAX,
4};
5use serde_json::Value;
6use std::collections::{BTreeMap, BTreeSet, VecDeque};
7use thiserror::Error;
8
9pub const HOOK_EVENT_ID_MAX_BYTES: usize = 256;
10pub const HOOK_SEEN_EVENT_IDS_MAX: usize = 256;
11const CLAUDE_SUBAGENT_ID_MAX_BYTES: usize = 64;
12const HOOK_PROVIDER_CACHE_KEYS_MAX: usize = 32;
13
14#[derive(Clone, Debug, Default, Eq, PartialEq)]
15struct ClaudeTrackedSubagent {
16 agent_type: Option<String>,
17 description: Option<String>,
18 background_tasks_authoritative: bool,
19 listed_as_subagent_task: bool,
20}
21
22#[derive(Clone, Debug, Eq, PartialEq)]
23struct ClaudeBackgroundTask {
24 id: String,
25 agent_type: Option<String>,
26 description: Option<String>,
27 running: bool,
28 teammate: bool,
29}
30
31#[derive(Clone, Debug, Eq, PartialEq)]
32pub struct HookEventEnvelope {
33 pub source_sequence: u64,
34 pub event_id: Option<String>,
35 pub event_name: String,
36 pub payload: Value,
37}
38
39#[derive(Clone, Copy, Debug, Eq, PartialEq)]
40pub enum HookEventDisposition {
41 Applied,
42 Duplicate,
43 StaleSequence,
44 IgnoredUnknown,
45}
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48pub struct HookSubagentSeed {
49 pub provider_agent_id: String,
50 pub agent_type: Option<String>,
51 pub description: Option<String>,
52}
53
54#[derive(Clone, Debug, Eq, PartialEq)]
55pub struct HookReduction {
56 pub source_sequence: u64,
57 pub missed_before: u64,
58 pub disposition: HookEventDisposition,
59 pub events: Vec<ProviderEvent>,
60}
61
62#[derive(Clone, Debug)]
68pub struct HookSessionReducer {
69 adapter_id: AdapterId,
70 last_source_sequence: u64,
71 seen_event_ids: BTreeSet<String>,
72 seen_event_order: VecDeque<String>,
73 tool_correlations: BTreeMap<String, VecDeque<String>>,
74 next_tool_id: u64,
75 claude_subagents: BTreeMap<String, ClaudeTrackedSubagent>,
76 amp_completed_threads: BTreeSet<String>,
77 amp_completed_thread_order: VecDeque<String>,
78 cursor_turn_completed: bool,
79}
80
81impl HookSessionReducer {
82 pub fn new(adapter_id: AdapterId) -> Self {
83 Self {
84 adapter_id,
85 last_source_sequence: 0,
86 seen_event_ids: BTreeSet::new(),
87 seen_event_order: VecDeque::new(),
88 tool_correlations: BTreeMap::new(),
89 next_tool_id: 1,
90 claude_subagents: BTreeMap::new(),
91 amp_completed_threads: BTreeSet::new(),
92 amp_completed_thread_order: VecDeque::new(),
93 cursor_turn_completed: false,
94 }
95 }
96
97 pub fn adapter_id(&self) -> &AdapterId {
98 &self.adapter_id
99 }
100
101 pub fn last_source_sequence(&self) -> u64 {
102 self.last_source_sequence
103 }
104
105 pub fn seed_live_subagents(&mut self, seeds: &[HookSubagentSeed]) -> usize {
111 if self.adapter_id.as_str() != "claude-code" || !self.claude_subagents.is_empty() {
112 return 0;
113 }
114 for seed in seeds.iter().take(PROVIDER_SUBAGENTS_MAX) {
115 self.upsert_claude_subagent(
116 &seed.provider_agent_id,
117 seed.agent_type.clone(),
118 seed.description.clone(),
119 );
120 if let Some(tracked) = self.claude_subagents.get_mut(&seed.provider_agent_id) {
121 tracked.background_tasks_authoritative = true;
122 }
123 }
124 self.claude_subagents.len()
125 }
126
127 pub fn reduce(
128 &mut self,
129 envelope: HookEventEnvelope,
130 ) -> Result<HookReduction, HookSessionReducerError> {
131 if envelope.source_sequence == 0 {
132 return Err(HookSessionReducerError::InvalidSourceSequence);
133 }
134 validate_event_id(envelope.event_id.as_deref())?;
135
136 if envelope.source_sequence <= self.last_source_sequence {
137 return Ok(HookReduction {
138 source_sequence: envelope.source_sequence,
139 missed_before: 0,
140 disposition: HookEventDisposition::StaleSequence,
141 events: Vec::new(),
142 });
143 }
144
145 let missed_before = envelope
146 .source_sequence
147 .saturating_sub(self.last_source_sequence)
148 .saturating_sub(1);
149 if missed_before > 0 {
150 self.tool_correlations.clear();
151 self.cursor_turn_completed = false;
152 }
153 self.last_source_sequence = envelope.source_sequence;
154
155 if let Some(event_id) = envelope.event_id {
156 if self.seen_event_ids.contains(&event_id) {
157 return Ok(HookReduction {
158 source_sequence: envelope.source_sequence,
159 missed_before,
160 disposition: HookEventDisposition::Duplicate,
161 events: Vec::new(),
162 });
163 }
164 self.remember_event_id(event_id);
165 }
166
167 if self.ignore_late_provider_event(&envelope.event_name, &envelope.payload) {
168 return Ok(HookReduction {
169 source_sequence: envelope.source_sequence,
170 missed_before,
171 disposition: HookEventDisposition::IgnoredUnknown,
172 events: Vec::new(),
173 });
174 }
175
176 let mut events =
177 normalize_hook_event(&self.adapter_id, &envelope.event_name, &envelope.payload)?;
178 self.reconcile_cursor_completion(&envelope.event_name, &mut events);
179 self.record_provider_completion(&envelope.event_name, &envelope.payload);
180 if self.adapter_id.as_str() == "claude-code" {
181 self.reconcile_claude_subagents(&envelope.event_name, &envelope.payload, &mut events);
182 }
183 self.correlate_tools(&envelope.event_name, &mut events);
184 for event in &events {
185 event.validate_ingress()?;
186 }
187 let disposition = if events.is_empty() {
188 HookEventDisposition::IgnoredUnknown
189 } else {
190 HookEventDisposition::Applied
191 };
192 Ok(HookReduction {
193 source_sequence: envelope.source_sequence,
194 missed_before,
195 disposition,
196 events,
197 })
198 }
199
200 pub fn clear_protocol_state(&mut self) {
201 self.last_source_sequence = 0;
202 self.seen_event_ids.clear();
203 self.seen_event_order.clear();
204 self.tool_correlations.clear();
205 self.next_tool_id = 1;
206 self.claude_subagents.clear();
207 self.amp_completed_threads.clear();
208 self.amp_completed_thread_order.clear();
209 self.cursor_turn_completed = false;
210 }
211
212 fn reconcile_cursor_completion(&mut self, event_name: &str, events: &mut Vec<ProviderEvent>) {
213 if self.adapter_id.as_str() != "cursor" {
214 return;
215 }
216 match event_name {
217 "sessionStart" | "beforeSubmitPrompt" => self.cursor_turn_completed = false,
218 "afterAgentResponse" if self.cursor_turn_completed => {
219 events.retain(|event| !matches!(event, ProviderEvent::WorkingObserved));
220 }
221 "stop" | "sessionEnd" => self.cursor_turn_completed = true,
222 _ => {}
223 }
224 }
225
226 fn reconcile_claude_subagents(
227 &mut self,
228 event_name: &str,
229 payload: &Value,
230 events: &mut Vec<ProviderEvent>,
231 ) {
232 let before = self.claude_subagents.clone();
233 let record = payload.as_object();
234 let agent_id = record.and_then(|record| read_nonempty_string(record.get("agent_id")));
235 let child_turn_boundary = matches!(event_name, "Stop" | "StopFailure")
236 && agent_id
237 .as_deref()
238 .is_some_and(|agent_id| before.contains_key(agent_id));
239 match event_name {
240 "SubagentStart" => {
241 if let Some(agent_id) = agent_id.as_deref() {
242 self.upsert_claude_subagent(
243 agent_id,
244 record.and_then(|record| read_nonempty_string(record.get("agent_type"))),
245 record.and_then(|record| read_nonempty_string(record.get("description"))),
246 );
247 }
248 }
249 "SubagentStop" => {
250 if let Some(agent_id) = agent_id.as_deref() {
251 self.claude_subagents.remove(agent_id);
252 }
253 }
254 "TeammateIdle" => {
255 if let Some(name) =
256 record.and_then(|record| read_nonempty_string(record.get("teammate_name")))
257 {
258 self.claude_subagents
259 .retain(|id, _| !claude_teammate_id_matches_name(id, &name));
260 }
261 }
262 "Stop" | "StopFailure" if !child_turn_boundary => {
263 if let Some(record) = record {
264 if let Some((tasks, complete)) = read_claude_background_tasks(record) {
265 self.fold_claude_background_tasks(tasks, complete);
266 }
267 }
268 }
269 "PreToolUse" | "PostToolUse" | "PostToolUseFailure" | "PermissionRequest" => {
270 if let Some(agent_id) = agent_id.as_deref() {
271 self.upsert_claude_subagent(
272 agent_id,
273 record.and_then(|record| read_nonempty_string(record.get("agent_type"))),
274 None,
275 );
276 }
277 }
278 _ => {}
279 }
280
281 events.retain(|event| {
282 !matches!(
283 event,
284 ProviderEvent::SubagentStarted { .. } | ProviderEvent::SubagentStopped { .. }
285 )
286 });
287 if child_turn_boundary {
288 events.clear();
289 }
290 let mut lifecycle = Vec::new();
291 for id in before.keys() {
292 if !self.claude_subagents.contains_key(id) {
293 lifecycle.push(ProviderEvent::SubagentStopped {
294 agent_id: id.clone(),
295 });
296 }
297 }
298 for (id, tracked) in &self.claude_subagents {
299 let changed = before.get(id).is_none_or(|previous| {
300 previous.agent_type != tracked.agent_type
301 || previous.description != tracked.description
302 });
303 if changed {
304 lifecycle.push(ProviderEvent::SubagentStarted {
305 agent_id: id.clone(),
306 agent_type: tracked.agent_type.clone(),
307 description: tracked.description.clone(),
308 });
309 }
310 }
311 lifecycle.append(events);
312 *events = lifecycle;
313 }
314
315 fn upsert_claude_subagent(
316 &mut self,
317 id: &str,
318 agent_type: Option<String>,
319 description: Option<String>,
320 ) {
321 if id.is_empty() || id.len() > CLAUDE_SUBAGENT_ID_MAX_BYTES {
322 return;
323 }
324 if let Some(existing) = self.claude_subagents.get_mut(id) {
325 existing.agent_type = agent_type.or(existing.agent_type.take());
326 existing.description = description.or(existing.description.take());
327 existing.background_tasks_authoritative = false;
328 return;
329 }
330 if self.claude_subagents.len() >= PROVIDER_SUBAGENTS_MAX {
331 return;
332 }
333 self.claude_subagents.insert(
334 id.to_owned(),
335 ClaudeTrackedSubagent {
336 agent_type,
337 description,
338 ..ClaudeTrackedSubagent::default()
339 },
340 );
341 }
342
343 fn fold_claude_background_tasks(
344 &mut self,
345 tasks: Vec<ClaudeBackgroundTask>,
346 inventory_complete: bool,
347 ) {
348 if tasks.is_empty() {
349 if inventory_complete {
350 self.claude_subagents.clear();
351 }
352 return;
353 }
354 let has_teammate_task = tasks.iter().any(|task| task.teammate);
355 let mut listed_ids = BTreeSet::new();
356 let mut pending_running_tasks = Vec::new();
357 for task in tasks.iter().filter(|task| !task.teammate) {
358 listed_ids.insert(task.id.clone());
359 if !task.running {
360 self.claude_subagents.remove(&task.id);
361 continue;
362 }
363 self.upsert_claude_subagent(
364 &task.id,
365 task.agent_type.clone(),
366 task.description.clone(),
367 );
368 if let Some(tracked) = self.claude_subagents.get_mut(&task.id) {
369 tracked.background_tasks_authoritative = true;
370 tracked.listed_as_subagent_task = true;
371 } else {
372 pending_running_tasks.push(task.clone());
373 }
374 }
375 if inventory_complete {
376 self.claude_subagents.retain(|id, tracked| {
377 listed_ids.contains(id)
378 || (has_teammate_task
379 && !tracked.background_tasks_authoritative
380 && !tracked.listed_as_subagent_task
381 && is_claude_teammate_lifecycle_id(id))
382 });
383 }
384 for task in pending_running_tasks {
385 if self.claude_subagents.len() >= PROVIDER_SUBAGENTS_MAX {
386 break;
387 }
388 self.upsert_claude_subagent(&task.id, task.agent_type, task.description);
389 if let Some(tracked) = self.claude_subagents.get_mut(&task.id) {
390 tracked.background_tasks_authoritative = true;
391 tracked.listed_as_subagent_task = true;
392 }
393 }
394 }
395
396 fn remember_event_id(&mut self, event_id: String) {
397 self.seen_event_ids.insert(event_id.clone());
398 self.seen_event_order.push_back(event_id);
399 while self.seen_event_order.len() > HOOK_SEEN_EVENT_IDS_MAX {
400 if let Some(expired) = self.seen_event_order.pop_front() {
401 self.seen_event_ids.remove(&expired);
402 }
403 }
404 }
405
406 fn ignore_late_provider_event(&mut self, event_name: &str, payload: &Value) -> bool {
407 match self.adapter_id.as_str() {
408 "amp" => {
409 let thread = amp_thread_key(payload);
410 if matches!(event_name, "session.start" | "agent.start") {
411 self.amp_completed_threads.remove(&thread);
412 self.amp_completed_thread_order
413 .retain(|existing| existing != &thread);
414 return false;
415 }
416 matches!(event_name, "tool.call" | "tool.result")
417 && self.amp_completed_threads.contains(&thread)
418 }
419 _ => false,
420 }
421 }
422
423 fn record_provider_completion(&mut self, event_name: &str, payload: &Value) {
424 match self.adapter_id.as_str() {
425 "amp" if event_name == "agent.end" => {
426 let thread = amp_thread_key(payload);
427 if self.amp_completed_threads.insert(thread.clone()) {
428 self.amp_completed_thread_order.push_back(thread);
429 }
430 while self.amp_completed_thread_order.len() > HOOK_PROVIDER_CACHE_KEYS_MAX {
431 if let Some(expired) = self.amp_completed_thread_order.pop_front() {
432 self.amp_completed_threads.remove(&expired);
433 }
434 }
435 }
436 _ => {}
437 }
438 }
439
440 fn correlate_tools(&mut self, event_name: &str, events: &mut [ProviderEvent]) {
441 for event in events {
442 match event {
443 ProviderEvent::SessionStarted { .. } | ProviderEvent::TurnStarted { .. } => {
444 self.tool_correlations.clear();
445 }
446 ProviderEvent::ToolStarted {
447 id, name, agent_id, ..
448 } => {
449 let raw_id = id.clone();
450 let correlation_key = tool_correlation_key(agent_id.as_deref(), &raw_id);
451 let coalesced_id = (matches!(self.adapter_id.as_str(), "pi" | "omp")
452 && event_name == "tool_execution_start")
453 .then(|| {
454 self.tool_correlations
455 .get(&correlation_key)
456 .and_then(|correlations| correlations.front())
457 .cloned()
458 })
459 .flatten();
460 if let Some(coalesced_id) = coalesced_id {
461 *id = coalesced_id;
462 continue;
463 }
464 let correlated_id = if raw_id.is_empty() || raw_id == *name {
465 let value = format!("hook-tool-{}", self.next_tool_id);
466 self.next_tool_id = self.next_tool_id.saturating_add(1);
467 value
468 } else {
469 raw_id.clone()
470 };
471 self.tool_correlations
472 .entry(correlation_key)
473 .or_default()
474 .push_back(correlated_id.clone());
475 *id = correlated_id;
476 }
477 ProviderEvent::ToolCompleted { id, agent_id, .. } => {
478 let raw_id = id.clone();
479 let correlation_key = tool_correlation_key(agent_id.as_deref(), &raw_id);
480 if let Some(correlations) = self.tool_correlations.get_mut(&correlation_key) {
481 if let Some(correlated_id) = correlations.pop_front() {
482 *id = correlated_id;
483 }
484 if correlations.is_empty() {
485 self.tool_correlations.remove(&correlation_key);
486 }
487 }
488 }
489 ProviderEvent::TurnCompleted { .. }
490 | ProviderEvent::TurnInterrupted
491 | ProviderEvent::SessionEnded { .. } => {
492 self.tool_correlations.clear();
493 }
494 ProviderEvent::Text { .. }
495 | ProviderEvent::ContextWindowUsage { .. }
496 | ProviderEvent::SessionIdentityObserved { .. }
497 | ProviderEvent::WorkingObserved
498 | ProviderEvent::Thinking { .. }
499 | ProviderEvent::Error { .. }
500 | ProviderEvent::Ready
501 | ProviderEvent::InteractionRequested { .. }
502 | ProviderEvent::InteractionResolved { .. }
503 | ProviderEvent::SubagentStarted { .. }
504 | ProviderEvent::SubagentStopped { .. }
505 | ProviderEvent::RateLimited { .. }
506 | ProviderEvent::HostRequestObserved { .. }
507 | ProviderEvent::UnrecognizedNotification { .. }
508 | ProviderEvent::UserMessage { .. }
511 | ProviderEvent::Plan { .. }
512 | ProviderEvent::AvailableCommandsUpdated { .. }
513 | ProviderEvent::ModeChanged { .. }
514 | ProviderEvent::SessionInfoUpdated { .. }
515 | ProviderEvent::UsageUpdated { .. }
516 | ProviderEvent::ConfigOptionsUpdated { .. } => {}
517 }
518 }
519 }
520}
521
522fn read_nonempty_string(value: Option<&Value>) -> Option<String> {
523 let value = value?.as_str()?.trim();
524 (!value.is_empty()).then(|| value.to_owned())
525}
526
527fn payload_string(payload: &Value, keys: &[&str], max_bytes: usize) -> Option<String> {
528 let record = payload.as_object()?;
529 let value = keys
530 .iter()
531 .find_map(|key| read_nonempty_string(record.get(*key)))?;
532 (value.len() <= max_bytes && !value.chars().any(char::is_control)).then_some(value)
533}
534
535fn amp_thread_key(payload: &Value) -> String {
536 payload_string(
537 payload,
538 &["threadId", "threadID", "thread_id"],
539 HOOK_EVENT_ID_MAX_BYTES,
540 )
541 .unwrap_or_else(|| "amp-default-thread".to_owned())
542}
543
544fn tool_correlation_key(agent_id: Option<&str>, raw_id: &str) -> String {
545 format!("{}\0{raw_id}", agent_id.unwrap_or("lead"))
546}
547
548fn read_claude_background_tasks(
549 record: &serde_json::Map<String, Value>,
550) -> Option<(Vec<ClaudeBackgroundTask>, bool)> {
551 let raw = record.get("background_tasks")?.as_array()?;
552 let mut tasks = Vec::new();
553 let mut complete = true;
554 for value in raw {
555 let Some(task) = value.as_object() else {
556 continue;
557 };
558 let Some(task_type) = task.get("type").and_then(Value::as_str) else {
559 continue;
560 };
561 if task_type != "subagent" && task_type != "teammate" {
562 continue;
563 }
564 let Some(id) = read_nonempty_string(task.get("id")) else {
565 continue;
566 };
567 if id.len() > CLAUDE_SUBAGENT_ID_MAX_BYTES {
568 continue;
569 }
570 if tasks.len() >= PROVIDER_SUBAGENTS_MAX {
571 complete = false;
572 break;
573 }
574 tasks.push(ClaudeBackgroundTask {
575 id,
576 agent_type: read_nonempty_string(task.get("agent_type")),
577 description: read_nonempty_string(task.get("description")),
578 running: task.get("status").and_then(Value::as_str) == Some("running"),
579 teammate: task_type == "teammate",
580 });
581 }
582 Some((tasks, complete))
583}
584
585fn is_claude_teammate_lifecycle_id(id: &str) -> bool {
586 let Some(separator) = id.rfind('-') else {
587 return false;
588 };
589 separator > 1
590 && id.starts_with('a')
591 && id[separator + 1..]
592 .chars()
593 .all(|character| character.is_ascii_hexdigit())
594}
595
596fn claude_teammate_id_matches_name(id: &str, name: &str) -> bool {
597 let prefix = format!("a{name}-");
598 id.strip_prefix(&prefix)
599 .is_some_and(|suffix| !suffix.is_empty() && !suffix.contains('-'))
600}
601
602fn validate_event_id(value: Option<&str>) -> Result<(), HookSessionReducerError> {
603 if value.is_some_and(|value| {
604 value.trim().is_empty()
605 || value.len() > HOOK_EVENT_ID_MAX_BYTES
606 || value.chars().any(char::is_control)
607 }) {
608 return Err(HookSessionReducerError::InvalidEventId);
609 }
610 Ok(())
611}
612
613#[derive(Clone, Debug, Error, Eq, PartialEq)]
614pub enum HookSessionReducerError {
615 #[error("hook source sequence must be greater than zero")]
616 InvalidSourceSequence,
617 #[error("hook event ID is empty, unsafe, or too large")]
618 InvalidEventId,
619 #[error(transparent)]
620 Normalize(#[from] HookAdapterError),
621 #[error(transparent)]
622 InvalidCanonicalEvent(#[from] ProviderEventValidationError),
623}
624
625#[cfg(test)]
626mod tests {
627 use super::*;
628 use serde_json::json;
629
630 fn reducer() -> HookSessionReducer {
631 HookSessionReducer::new(AdapterId::new("grok").unwrap())
632 }
633
634 fn envelope(
635 source_sequence: u64,
636 event_id: &str,
637 event_name: &str,
638 payload: Value,
639 ) -> HookEventEnvelope {
640 HookEventEnvelope {
641 source_sequence,
642 event_id: Some(event_id.to_owned()),
643 event_name: event_name.to_owned(),
644 payload,
645 }
646 }
647
648 #[test]
649 fn correlates_provider_tools_without_explicit_ids() {
650 let mut reducer = reducer();
651 let started = reducer
652 .reduce(envelope(
653 1,
654 "e1",
655 "PreToolUse",
656 json!({"toolName": "shell", "toolInput": {"command": "pwd"}}),
657 ))
658 .unwrap();
659 let [ProviderEvent::ToolStarted { id, .. }] = started.events.as_slice() else {
660 panic!("expected tool start");
661 };
662 assert_eq!(id, "hook-tool-1");
663
664 let completed = reducer
665 .reduce(envelope(
666 2,
667 "e2",
668 "PostToolUse",
669 json!({"toolName": "shell", "toolResponse": "ok"}),
670 ))
671 .unwrap();
672 let [ProviderEvent::ToolCompleted { id, .. }] = completed.events.as_slice() else {
673 panic!("expected tool completion");
674 };
675 assert_eq!(id, "hook-tool-1");
676 }
677
678 #[test]
679 fn pi_coalesces_call_and_execution_start_before_exact_completion() {
680 for adapter_id in ["pi", "omp"] {
681 let mut reducer = HookSessionReducer::new(AdapterId::new(adapter_id).unwrap());
682 let called = reducer
683 .reduce(envelope(
684 1,
685 "e1",
686 "tool_call",
687 json!({"tool_name": "bash", "tool_input": {"command": "cargo check"}}),
688 ))
689 .unwrap();
690 let [ProviderEvent::ToolStarted { id: called_id, .. }] = called.events.as_slice()
691 else {
692 panic!("expected tool call");
693 };
694
695 let executing = reducer
696 .reduce(envelope(
697 2,
698 "e2",
699 "tool_execution_start",
700 json!({"tool_name": "bash", "tool_input": {"command": "cargo check"}}),
701 ))
702 .unwrap();
703 let [ProviderEvent::ToolStarted {
704 id: executing_id, ..
705 }] = executing.events.as_slice()
706 else {
707 panic!("expected execution start");
708 };
709 assert_eq!(executing_id, called_id);
710
711 let completed = reducer
712 .reduce(envelope(
713 3,
714 "e3",
715 "tool_execution_end",
716 json!({"tool_name": "bash"}),
717 ))
718 .unwrap();
719 let [ProviderEvent::ToolCompleted {
720 id: completed_id, ..
721 }] = completed.events.as_slice()
722 else {
723 panic!("expected execution completion");
724 };
725 assert_eq!(completed_id, called_id);
726 }
727 }
728
729 #[test]
730 fn amp_drops_late_tool_events_per_thread_until_restart() {
731 let mut reducer = HookSessionReducer::new(AdapterId::new("amp").unwrap());
732 reducer
733 .reduce(envelope(
734 1,
735 "e1",
736 "agent.end",
737 json!({"threadId": "thread-1", "status": "completed"}),
738 ))
739 .unwrap();
740 let late = reducer
741 .reduce(envelope(
742 2,
743 "e2",
744 "tool.result",
745 json!({"threadId": "thread-1", "toolUseId": "t1", "tool": "bash"}),
746 ))
747 .unwrap();
748 assert_eq!(late.disposition, HookEventDisposition::IgnoredUnknown);
749
750 let other_thread = reducer
751 .reduce(envelope(
752 3,
753 "e3",
754 "tool.call",
755 json!({"threadId": "thread-2", "toolUseId": "t2", "tool": "read"}),
756 ))
757 .unwrap();
758 assert_eq!(other_thread.disposition, HookEventDisposition::Applied);
759
760 reducer
761 .reduce(envelope(
762 4,
763 "e4",
764 "agent.start",
765 json!({"threadId": "thread-1", "message": "retry"}),
766 ))
767 .unwrap();
768 let resumed = reducer
769 .reduce(envelope(
770 5,
771 "e5",
772 "tool.call",
773 json!({"threadId": "thread-1", "toolUseId": "t3", "tool": "read"}),
774 ))
775 .unwrap();
776 assert_eq!(resumed.disposition, HookEventDisposition::Applied);
777 }
778
779 #[test]
780 fn cursor_late_response_enriches_completion_without_resurrecting_work() {
781 let mut reducer = HookSessionReducer::new(AdapterId::new("cursor").unwrap());
782 reducer
783 .reduce(envelope(
784 1,
785 "cursor-1",
786 "beforeSubmitPrompt",
787 json!({"prompt": "add tests"}),
788 ))
789 .unwrap();
790 reducer
791 .reduce(envelope(
792 2,
793 "cursor-2",
794 "stop",
795 json!({"status": "completed"}),
796 ))
797 .unwrap();
798
799 let late = reducer
800 .reduce(envelope(
801 3,
802 "cursor-3",
803 "afterAgentResponse",
804 json!({"text": "All set"}),
805 ))
806 .unwrap();
807 assert!(matches!(
808 late.events.as_slice(),
809 [ProviderEvent::Text { text, .. }] if text == "All set"
810 ));
811
812 reducer
813 .reduce(envelope(
814 4,
815 "cursor-4",
816 "beforeSubmitPrompt",
817 json!({"prompt": "next"}),
818 ))
819 .unwrap();
820 let active = reducer
821 .reduce(envelope(
822 5,
823 "cursor-5",
824 "afterAgentResponse",
825 json!({"text": "Draft"}),
826 ))
827 .unwrap();
828 assert!(matches!(
829 active.events.as_slice(),
830 [ProviderEvent::WorkingObserved, ProviderEvent::Text { text, .. }]
831 if text == "Draft"
832 ));
833 }
834
835 #[test]
836 fn clearing_protocol_state_clears_provider_late_event_caches() {
837 let mut reducer = HookSessionReducer::new(AdapterId::new("cursor").unwrap());
838 reducer
839 .reduce(envelope(
840 1,
841 "cursor-1",
842 "stop",
843 json!({"status": "completed"}),
844 ))
845 .unwrap();
846 reducer.clear_protocol_state();
847
848 let response = reducer
849 .reduce(envelope(
850 1,
851 "cursor-2",
852 "afterAgentResponse",
853 json!({"text": "fresh"}),
854 ))
855 .unwrap();
856 assert!(matches!(
857 response.events.as_slice(),
858 [ProviderEvent::WorkingObserved, ProviderEvent::Text { text, .. }]
859 if text == "fresh"
860 ));
861 }
862
863 #[test]
864 fn correlates_idless_tools_independently_per_claude_child() {
865 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
866 for (sequence, agent_id) in ["a1", "a2"].into_iter().enumerate() {
867 reducer
868 .reduce(envelope(
869 u64::try_from(sequence + 1).unwrap(),
870 &format!("start-{agent_id}"),
871 "PreToolUse",
872 json!({"agent_id": agent_id, "tool_name": "shell"}),
873 ))
874 .unwrap();
875 }
876 let child_two = reducer
877 .reduce(envelope(
878 3,
879 "done-a2",
880 "PostToolUse",
881 json!({"agent_id": "a2", "tool_name": "shell", "tool_response": "ok"}),
882 ))
883 .unwrap();
884 assert!(child_two.events.iter().any(|event| matches!(
885 event,
886 ProviderEvent::ToolCompleted { id, agent_id: Some(agent_id), .. }
887 if id == "hook-tool-2" && agent_id == "a2"
888 )));
889
890 let child_one = reducer
891 .reduce(envelope(
892 4,
893 "done-a1",
894 "PostToolUse",
895 json!({"agent_id": "a1", "tool_name": "shell", "tool_response": "ok"}),
896 ))
897 .unwrap();
898 assert!(child_one.events.iter().any(|event| matches!(
899 event,
900 ProviderEvent::ToolCompleted { id, agent_id: Some(agent_id), .. }
901 if id == "hook-tool-1" && agent_id == "a1"
902 )));
903 }
904
905 #[test]
906 fn suppresses_replayed_ids_and_stale_sequences() {
907 let mut reducer = reducer();
908 reducer
909 .reduce(envelope(
910 1,
911 "e1",
912 "UserPromptSubmit",
913 json!({"prompt": "one"}),
914 ))
915 .unwrap();
916 let duplicate = reducer
917 .reduce(envelope(
918 2,
919 "e1",
920 "UserPromptSubmit",
921 json!({"prompt": "one"}),
922 ))
923 .unwrap();
924 assert_eq!(duplicate.disposition, HookEventDisposition::Duplicate);
925 assert!(duplicate.events.is_empty());
926
927 let stale = reducer
928 .reduce(envelope(1, "e2", "Stop", json!({})))
929 .unwrap();
930 assert_eq!(stale.disposition, HookEventDisposition::StaleSequence);
931 assert!(stale.events.is_empty());
932 }
933
934 #[test]
935 fn reports_gaps_and_drops_unsafe_tool_correlation() {
936 let mut reducer = reducer();
937 reducer
938 .reduce(envelope(
939 1,
940 "e1",
941 "PreToolUse",
942 json!({"toolName": "shell"}),
943 ))
944 .unwrap();
945 let after_gap = reducer
946 .reduce(envelope(
947 4,
948 "e4",
949 "PostToolUse",
950 json!({"toolName": "shell", "toolResponse": "unknown start"}),
951 ))
952 .unwrap();
953 assert_eq!(after_gap.missed_before, 2);
954 let [ProviderEvent::ToolCompleted { id, .. }] = after_gap.events.as_slice() else {
955 panic!("expected partial completion");
956 };
957 assert_eq!(id, "shell");
958 }
959
960 #[test]
961 fn emits_canonical_turn_start_and_completion() {
962 let mut reducer = reducer();
963 let started = reducer
964 .reduce(envelope(
965 1,
966 "e1",
967 "userPromptSubmit",
968 json!({"prompt": "fix tests"}),
969 ))
970 .unwrap();
971 assert!(matches!(
972 started.events.as_slice(),
973 [ProviderEvent::TurnStarted { prompt }] if prompt.as_deref() == Some("fix tests")
974 ));
975 let completed = reducer
976 .reduce(envelope(2, "e2", "stop", json!({})))
977 .unwrap();
978 assert!(matches!(
979 completed.events.last(),
980 Some(ProviderEvent::TurnCompleted { .. })
981 ));
982 }
983
984 #[test]
985 fn claude_stop_reconciles_one_shot_inventory_before_turn_completion() {
986 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
987 let started = reducer
988 .reduce(envelope(
989 1,
990 "c1",
991 "SubagentStart",
992 json!({"agent_id": "a1", "agent_type": "reviewer"}),
993 ))
994 .unwrap();
995 assert!(matches!(
996 started.events.as_slice(),
997 [ProviderEvent::SubagentStarted { agent_id, .. }] if agent_id == "a1"
998 ));
999
1000 let completed = reducer
1001 .reduce(envelope(2, "c2", "Stop", json!({"background_tasks": []})))
1002 .unwrap();
1003 assert!(matches!(
1004 completed.events.as_slice(),
1005 [ProviderEvent::SubagentStopped { agent_id }, ProviderEvent::TurnCompleted { .. }]
1006 if agent_id == "a1"
1007 ));
1008 }
1009
1010 #[test]
1011 fn claude_inventory_recovers_running_child_after_listener_restart() {
1012 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
1013 let recovered = reducer
1014 .reduce(envelope(
1015 1,
1016 "c1",
1017 "Stop",
1018 json!({
1019 "background_tasks": [{
1020 "id": "a77",
1021 "type": "subagent",
1022 "status": "running",
1023 "agent_type": "probe",
1024 "description": "verify restart"
1025 }]
1026 }),
1027 ))
1028 .unwrap();
1029 assert!(matches!(
1030 recovered.events.as_slice(),
1031 [
1032 ProviderEvent::SubagentStarted {
1033 agent_id,
1034 agent_type: Some(agent_type),
1035 ..
1036 },
1037 ProviderEvent::TurnCompleted { .. }
1038 ] if agent_id == "a77" && agent_type == "probe"
1039 ));
1040 }
1041
1042 #[test]
1043 fn claude_persisted_seed_is_reaped_by_complete_inventory() {
1044 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
1045 assert_eq!(
1046 reducer.seed_live_subagents(&[HookSubagentSeed {
1047 provider_agent_id: "areviewer-6d3cb5b5".to_owned(),
1048 agent_type: Some("reviewer".to_owned()),
1049 description: None,
1050 }]),
1051 1
1052 );
1053 let reconciled = reducer
1054 .reduce(envelope(
1055 1,
1056 "c1",
1057 "Stop",
1058 json!({
1059 "background_tasks": [{
1060 "id": "team-reviewer",
1061 "type": "teammate",
1062 "status": "running"
1063 }]
1064 }),
1065 ))
1066 .unwrap();
1067 assert!(matches!(
1068 reconciled.events.as_slice(),
1069 [ProviderEvent::SubagentStopped { agent_id }, ProviderEvent::TurnCompleted { .. }]
1070 if agent_id == "areviewer-6d3cb5b5"
1071 ));
1072 }
1073
1074 #[test]
1075 fn claude_child_stop_is_not_misclassified_as_lead_completion() {
1076 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
1077 reducer
1078 .reduce(envelope(
1079 1,
1080 "c1",
1081 "SubagentStart",
1082 json!({"agent_id": "a1"}),
1083 ))
1084 .unwrap();
1085 let child_stop = reducer
1086 .reduce(envelope(2, "c2", "Stop", json!({"agent_id": "a1"})))
1087 .unwrap();
1088 assert_eq!(child_stop.disposition, HookEventDisposition::IgnoredUnknown);
1089 assert!(child_stop.events.is_empty());
1090
1091 let stopped = reducer
1092 .reduce(envelope(3, "c3", "SubagentStop", json!({"agent_id": "a1"})))
1093 .unwrap();
1094 assert!(matches!(
1095 stopped.events.as_slice(),
1096 [ProviderEvent::SubagentStopped { agent_id }] if agent_id == "a1"
1097 ));
1098 }
1099
1100 #[test]
1101 fn claude_teammate_idle_removes_only_exact_named_lifecycle_rows() {
1102 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
1103 for (sequence, id) in ["arev-6d3cb5b5", "arev-two-6d3cb5b5"]
1104 .into_iter()
1105 .enumerate()
1106 {
1107 reducer
1108 .reduce(envelope(
1109 u64::try_from(sequence + 1).unwrap(),
1110 &format!("c{}", sequence + 1),
1111 "SubagentStart",
1112 json!({"agent_id": id}),
1113 ))
1114 .unwrap();
1115 }
1116 let idled = reducer
1117 .reduce(envelope(
1118 3,
1119 "c3",
1120 "TeammateIdle",
1121 json!({"teammate_name": "rev"}),
1122 ))
1123 .unwrap();
1124 assert!(matches!(
1125 idled.events.as_slice(),
1126 [ProviderEvent::SubagentStopped { agent_id }] if agent_id == "arev-6d3cb5b5"
1127 ));
1128 }
1129
1130 #[test]
1131 fn claude_complete_inventory_retries_replacement_after_stale_cleanup() {
1132 let mut reducer = HookSessionReducer::new(AdapterId::new("claude-code").unwrap());
1133 for index in 0..PROVIDER_SUBAGENTS_MAX {
1134 reducer
1135 .reduce(envelope(
1136 u64::try_from(index + 1).unwrap(),
1137 &format!("start-{index}"),
1138 "SubagentStart",
1139 json!({"agent_id": format!("stale{index}")}),
1140 ))
1141 .unwrap();
1142 }
1143 let reconciled = reducer
1144 .reduce(envelope(
1145 u64::try_from(PROVIDER_SUBAGENTS_MAX + 1).unwrap(),
1146 "inventory",
1147 "Stop",
1148 json!({
1149 "background_tasks": [{
1150 "id": "replacement",
1151 "type": "subagent",
1152 "status": "running"
1153 }]
1154 }),
1155 ))
1156 .unwrap();
1157 assert!(reconciled.events.iter().any(|event| matches!(
1158 event,
1159 ProviderEvent::SubagentStarted { agent_id, .. } if agent_id == "replacement"
1160 )));
1161 assert_eq!(
1162 reconciled
1163 .events
1164 .iter()
1165 .filter(|event| matches!(event, ProviderEvent::SubagentStopped { .. }))
1166 .count(),
1167 PROVIDER_SUBAGENTS_MAX
1168 );
1169 }
1170}