1use gate4agent_types::{
4 normalize_semantic_prompt, prepare_agent_command, prepare_input, prepare_shell_command,
5 validate_candidate_id, validate_capability_models, validate_history_error,
6 validate_resume_error, validate_session_config_value_json, validate_session_control_id,
7 ActiveProviderTool, AgentInstanceId,
8 CapabilityProbeRequest,
9 CapabilitySnapshot, CommandEnvelope, CommandId, ControlCommand, ControlEffect, ControlError,
10 ControlEvent, ControlEventKind, ControlHealth, ControlObservation, ControlSnapshot,
11 EffectEnvelope, ForegroundAuthority, ForegroundRequirement, ForegroundSnapshot,
12 HistoryOperation, HistoryQuery, HistorySnapshot, InputAction, ObservationEnvelope,
13 ObservationIgnoredReason, OperationId, PendingCapabilityProbe, PendingHistoryOperation,
14 PendingResumeOperation, PreparedInputKind, ProviderActivity, ProviderEvent,
15 ProviderInteraction, ProviderInteractionId, ProviderInteractionKind,
16 ProviderInteractionOutcome, ProviderInteractionResponse, ProviderInteractionResponseKind,
17 ProviderInteractionStatus, ProviderInteractionTarget, ProviderSessionIdentity,
18 ProviderRuntimeCapability, ProviderRuntimePolicy, ProviderSessionKey, ProviderSnapshot,
19 ProviderSource, ProviderSourceCursor, ProviderSubagent, PtyScreenState,
20 ResumeAuthorityTarget, ResumeLaunchRequest, ResumePhase, ResumeSessionSummary, ResumeSnapshot,
21 ResumeTarget, SessionGeneration, SessionSnapshot, SessionStatus, StartRequest, TerminalControl,
22 TerminalSize, TokenUsage, TransportKind, CONTROL_INSTANCE_IDENTITIES_CAPACITY,
23 CONTROL_INSTANCE_IDENTITIES_MAX, CONTROL_SESSIONS_MAX,
24 PROVIDER_INGRESS_EVENTS_MAX, PROVIDER_INTERACTIONS_MAX, PROVIDER_INTERACTION_FAILURE_MAX_BYTES,
25 PROVIDER_SUBAGENTS_MAX, WORKING_DIRECTORY_MAX_BYTES,
26};
27use std::collections::{BTreeMap, BTreeSet};
28
29const CONTROL_REVISION_HEADROOM: u64 = PROVIDER_INGRESS_EVENTS_MAX as u64 + 1;
30const CONTROL_EVENT_HEADROOM: u64 = PROVIDER_INGRESS_EVENTS_MAX as u64
31 * (PROVIDER_INTERACTIONS_MAX as u64 + 1)
32 + PROVIDER_INTERACTIONS_MAX as u64
33 + 2;
34
35#[derive(Clone, Debug, Eq, PartialEq)]
36struct SessionState {
37 snapshot: SessionSnapshot,
38 runtime_policy: ProviderRuntimePolicy,
39 pending_terminal_size: Option<TerminalSize>,
40 pending_interrupt: bool,
41 pending_resume_identity: Option<ProviderSessionIdentity>,
42}
43
44#[derive(Clone, Debug, Eq, PartialEq)]
47pub struct Gate4AgentEngine {
48 sessions: BTreeMap<AgentInstanceId, SessionState>,
49 generation_watermarks: BTreeMap<AgentInstanceId, SessionGeneration>,
50 next_operation_id: Option<u64>,
51 next_event_sequence: Option<u64>,
52 revision: u64,
53 counter_error: Option<ControlError>,
54 effects: Vec<EffectEnvelope>,
55 events: Vec<ControlEvent>,
56}
57
58impl Gate4AgentEngine {
59 pub fn new() -> Self {
60 Self {
61 sessions: BTreeMap::new(),
62 generation_watermarks: BTreeMap::new(),
63 next_operation_id: Some(1),
64 next_event_sequence: Some(1),
65 revision: 0,
66 counter_error: None,
67 effects: Vec::new(),
68 events: Vec::new(),
69 }
70 }
71
72 pub fn apply_command(&mut self, envelope: CommandEnvelope) -> Result<(), ControlError> {
73 if self.has_counter_headroom() {
74 let result = self.apply_command_in_place(envelope);
75 debug_assert!(
76 self.counter_error.is_none(),
77 "bounded command exceeded reserved control counter headroom"
78 );
79 return result;
80 }
81 let mut candidate = self.clone();
82 candidate.apply_command_in_place(envelope)?;
83 if let Some(error) = candidate.counter_error.take() {
84 self.retire_exhausted_counter(&error);
85 return Err(error);
86 }
87 *self = candidate;
88 Ok(())
89 }
90
91 fn apply_command_in_place(&mut self, envelope: CommandEnvelope) -> Result<(), ControlError> {
92 let command_id = envelope.id;
93 match envelope.command {
94 ControlCommand::Register {
95 instance_id,
96 agent_id,
97 transport,
98 } => {
99 if self.sessions.contains_key(&instance_id) {
100 return Err(ControlError::DuplicateInstance { instance_id });
101 }
102 if self.sessions.len() >= CONTROL_SESSIONS_MAX {
103 return Err(ControlError::SessionCapacityExceeded {
104 instance_id,
105 max: CONTROL_SESSIONS_MAX,
106 });
107 }
108 if !self.generation_watermarks.contains_key(&instance_id)
109 && self.generation_watermarks.len() >= CONTROL_INSTANCE_IDENTITIES_MAX
110 {
111 return Err(ControlError::InstanceIdentityCapacityExceeded {
112 instance_id,
113 max: CONTROL_INSTANCE_IDENTITIES_MAX,
114 });
115 }
116 let generation = match self.generation_watermarks.get(&instance_id).copied() {
117 Some(watermark) => checked_next_generation(watermark).ok_or(
118 ControlError::GenerationExhausted {
119 instance_id,
120 generation: watermark,
121 },
122 )?,
123 None => SessionGeneration::default(),
124 };
125 self.sessions.insert(
126 instance_id,
127 SessionState {
128 snapshot: SessionSnapshot {
129 instance_id,
130 agent_id,
131 transport,
132 generation,
133 status: SessionStatus::Registered,
134 pending_operation: None,
135 pending_input: None,
136 process_id: None,
137 terminal_size: None,
138 terminal_frame: None,
139 terminal_stale: None,
140 session_options: None,
141 capabilities: CapabilitySnapshot::default(),
142 history: HistorySnapshot::default(),
143 resume: ResumeSnapshot::default(),
144 foreground: ForegroundSnapshot::default(),
145 provider: ProviderSnapshot::default(),
146 screen_state: PtyScreenState::Unknown,
148 },
149 runtime_policy: ProviderRuntimePolicy::raw_pty(),
150 pending_terminal_size: None,
151 pending_interrupt: false,
152 pending_resume_identity: None,
153 },
154 );
155 self.generation_watermarks.insert(instance_id, generation);
156 self.bump_revision();
157 self.emit_event(
158 Some(command_id),
159 instance_id,
160 generation,
161 ControlEventKind::Registered,
162 );
163 Ok(())
164 }
165 ControlCommand::Start {
166 instance_id,
167 runtime_policy,
168 request,
169 } => self.start(command_id, instance_id, runtime_policy, request),
170 ControlCommand::Stop { instance_id, force } => {
171 self.stop(command_id, instance_id, force)
172 }
173 ControlCommand::SendInput {
174 instance_id,
175 action,
176 } => self.send_input(command_id, instance_id, action),
177 ControlCommand::Resize { instance_id, size } => {
178 self.resize(command_id, instance_id, size)
179 }
180 ControlCommand::RefreshForeground { instance_id } => {
181 self.refresh_foreground(command_id, instance_id)
182 }
183 ControlCommand::ProbeCapabilities {
184 instance_id,
185 request,
186 } => self.probe_capabilities(command_id, instance_id, request),
187 ControlCommand::DiscoverHistory { instance_id, query } => {
188 self.discover_history(command_id, instance_id, query)
189 }
190 ControlCommand::LoadHistory {
191 instance_id,
192 candidate_id,
193 } => self.load_history(command_id, instance_id, candidate_id),
194 ControlCommand::Resume {
195 instance_id,
196 target,
197 runtime_policy,
198 request,
199 } => self.resume(command_id, instance_id, target, runtime_policy, request),
200 ControlCommand::ResolveInteraction {
201 instance_id,
202 generation,
203 interaction_id,
204 response,
205 } => self.resolve_interaction(
206 command_id,
207 instance_id,
208 generation,
209 interaction_id,
210 response,
211 ),
212 ControlCommand::SetSessionMode {
213 instance_id,
214 mode_id,
215 } => self.set_session_mode(command_id, instance_id, mode_id),
216 ControlCommand::SetSessionConfigOption {
217 instance_id,
218 option_id,
219 value_json,
220 } => self.set_session_config_option(command_id, instance_id, option_id, value_json),
221 ControlCommand::SetSessionModel {
222 instance_id,
223 model_id,
224 } => self.set_session_model(command_id, instance_id, model_id),
225 ControlCommand::IngestProvider {
226 instance_id,
227 generation,
228 source,
229 source_sequence,
230 events,
231 } => self.ingest_provider(
232 command_id,
233 instance_id,
234 generation,
235 source,
236 source_sequence,
237 events,
238 ),
239 ControlCommand::Remove { instance_id } => self.remove(command_id, instance_id),
240 }
241 }
242
243 pub fn apply_observation(&mut self, envelope: ObservationEnvelope) {
244 let _ = self.try_apply_observation(envelope);
245 }
246
247 pub fn try_apply_observation(
248 &mut self,
249 envelope: ObservationEnvelope,
250 ) -> Result<(), ControlError> {
251 if self.has_counter_headroom() && self.observation_has_provider_headroom(&envelope) {
252 self.apply_observation_in_place(envelope);
253 debug_assert!(
254 self.counter_error.is_none(),
255 "bounded observation exceeded reserved control counter headroom"
256 );
257 return Ok(());
258 }
259 let mut candidate = self.clone();
260 candidate.apply_observation_in_place(envelope);
261 if let Some(error) = candidate.counter_error.take() {
262 self.retire_exhausted_counter(&error);
263 return Err(error);
264 }
265 *self = candidate;
266 Ok(())
267 }
268
269 fn apply_observation_in_place(&mut self, envelope: ObservationEnvelope) {
270 let instance_id = envelope.instance_id;
271 let generation = envelope.generation;
272 let Some(current_state) = self.sessions.get(&instance_id) else {
273 self.emit_ignored(
274 instance_id,
275 generation,
276 ObservationIgnoredReason::UnknownInstance,
277 );
278 return;
279 };
280 let runtime_policy = current_state.runtime_policy;
281 let current = ¤t_state.snapshot;
282
283 let capability_generation_matches = is_capability_observation(&envelope.observation)
284 && current
285 .capabilities
286 .pending
287 .as_ref()
288 .is_some_and(|pending| pending.generation == generation);
289 if current.generation != generation && !capability_generation_matches {
290 self.emit_ignored(
291 instance_id,
292 generation,
293 ObservationIgnoredReason::StaleGeneration,
294 );
295 return;
296 }
297
298 if let Some(capability) = denied_observation_capability(runtime_policy, &envelope.observation)
299 {
300 self.emit_ignored(
301 instance_id,
302 generation,
303 ObservationIgnoredReason::ProviderRuntimePolicyDenied { capability },
304 );
305 return;
306 }
307
308 if envelope.observation.requires_operation_id() {
309 let Some(operation_id) = envelope.operation_id else {
310 self.emit_ignored(
311 instance_id,
312 generation,
313 ObservationIgnoredReason::MissingOperation,
314 );
315 return;
316 };
317 let expected_operation = if is_capability_observation(&envelope.observation) {
318 current
319 .capabilities
320 .pending
321 .as_ref()
322 .map(|pending| pending.operation_id)
323 } else if is_history_observation(&envelope.observation) {
324 current
325 .history
326 .pending
327 .as_ref()
328 .map(|pending| pending.operation_id)
329 } else {
330 current.pending_operation
331 };
332 if expected_operation != Some(operation_id) {
333 self.emit_ignored(
334 instance_id,
335 generation,
336 ObservationIgnoredReason::OperationMismatch,
337 );
338 return;
339 }
340 }
341
342 if let ControlObservation::ForegroundObserved { process } = &envelope.observation {
343 if !process.is_valid_for(¤t.agent_id)
344 || current
345 .process_id
346 .is_some_and(|root_process_id| root_process_id != process.root_process_id)
347 {
348 self.emit_ignored(
349 instance_id,
350 generation,
351 ObservationIgnoredReason::InvalidForegroundObservation,
352 );
353 return;
354 }
355 }
356
357 let valid = matches!(
358 (¤t.status, &envelope.observation),
359 (SessionStatus::Starting, ControlObservation::Spawned { .. })
360 | (
361 SessionStatus::Starting,
362 ControlObservation::SpawnFailed { .. }
363 )
364 | (
365 SessionStatus::Stopping,
366 ControlObservation::StopCompleted { .. }
367 )
368 | (
369 SessionStatus::Stopping,
370 ControlObservation::StopFailed { .. }
371 )
372 | (SessionStatus::Running, ControlObservation::InputCompleted)
373 | (
374 SessionStatus::Running,
375 ControlObservation::InputFailed { .. }
376 )
377 | (
378 SessionStatus::Running,
379 ControlObservation::ResizeCompleted { .. }
380 )
381 | (
382 SessionStatus::Running,
383 ControlObservation::ResizeFailed { .. }
384 )
385 | (
386 SessionStatus::Running,
387 ControlObservation::ForegroundObserved { .. }
388 | ControlObservation::ForegroundFailed { .. },
389 )
390 | (
391 SessionStatus::Running,
392 ControlObservation::InteractionResolutionCompleted { .. }
393 | ControlObservation::InteractionResolutionFailed { .. },
394 )
395 | (
396 SessionStatus::Running,
397 ControlObservation::SessionModeSet { .. }
398 | ControlObservation::SessionModeSetFailed { .. },
399 )
400 | (
401 SessionStatus::Running,
402 ControlObservation::SessionConfigOptionSet { .. }
403 | ControlObservation::SessionConfigOptionSetFailed { .. },
404 )
405 | (
406 SessionStatus::Running,
407 ControlObservation::SessionModelSet { .. }
408 | ControlObservation::SessionModelSetFailed { .. },
409 )
410 | (
411 SessionStatus::Running | SessionStatus::Stopping,
412 ControlObservation::TerminalFrame { .. }
413 | ControlObservation::TerminalStale { .. }
414 | ControlObservation::ScreenState { .. },
415 )
416 | (
417 SessionStatus::Starting | SessionStatus::Running | SessionStatus::Stopping,
418 ControlObservation::ProviderEvent { .. }
419 | ControlObservation::ProviderGap { .. },
420 )
421 | (
422 SessionStatus::Starting | SessionStatus::Running | SessionStatus::Stopping,
423 ControlObservation::ProcessExited { .. },
424 )
425 | (
426 _,
427 ControlObservation::CapabilitiesProbed { .. }
428 | ControlObservation::CapabilityProbeFailed { .. },
429 )
430 | (
431 _,
432 ControlObservation::HistoryDiscovered { .. }
433 | ControlObservation::HistoryLoaded { .. }
434 | ControlObservation::HistoryFailed { .. },
435 )
436 | (
437 _,
438 ControlObservation::ResumeAuthorized { .. }
439 | ControlObservation::ResumeDenied { .. }
440 | ControlObservation::ResumeFailed { .. },
441 )
442 );
443 if !valid {
444 self.emit_ignored(
445 instance_id,
446 generation,
447 ObservationIgnoredReason::InvalidState,
448 );
449 return;
450 }
451
452 if is_capability_observation(&envelope.observation) {
453 let pending = current
454 .capabilities
455 .pending
456 .as_ref()
457 .expect("capability operation correlation was validated")
458 .clone();
459 let event = match &envelope.observation {
460 ControlObservation::CapabilitiesProbed {
461 session_option_models,
462 } if validate_capability_models(session_option_models).is_ok() => {
463 ControlEventKind::CapabilitiesProbed {
464 count: session_option_models.len(),
465 }
466 }
467 ControlObservation::CapabilityProbeFailed { failure } => {
468 ControlEventKind::CapabilityProbeFailed { failure: *failure }
469 }
470 _ => {
471 self.emit_ignored(
472 instance_id,
473 generation,
474 ObservationIgnoredReason::InvalidCapabilityObservation,
475 );
476 return;
477 }
478 };
479 let event_generation = self
480 .sessions
481 .get(&instance_id)
482 .expect("validated session")
483 .snapshot
484 .generation;
485 let state = self
486 .sessions
487 .get_mut(&instance_id)
488 .expect("validated session");
489 debug_assert_eq!(state.snapshot.capabilities.pending.as_ref(), Some(&pending));
490 state.snapshot.capabilities.pending = None;
491 state.snapshot.capabilities.settled = true;
492 match envelope.observation {
493 ControlObservation::CapabilitiesProbed {
494 session_option_models,
495 } => {
496 state.snapshot.capabilities.session_option_models = session_option_models;
497 state.snapshot.capabilities.last_failure = None;
498 }
499 ControlObservation::CapabilityProbeFailed { failure } => {
500 state.snapshot.capabilities.session_option_models.clear();
501 state.snapshot.capabilities.last_failure = Some(failure);
502 }
503 _ => unreachable!("capability observation was matched above"),
504 }
505 self.bump_revision();
506 self.emit_event(None, instance_id, event_generation, event);
507 return;
508 }
509
510 if is_history_observation(&envelope.observation) {
511 let pending = current
512 .history
513 .pending
514 .as_ref()
515 .expect("history operation correlation was validated")
516 .clone();
517 let event = match (&envelope.observation, &pending.operation) {
518 (
519 ControlObservation::HistoryDiscovered { candidates },
520 HistoryOperation::Discover { query },
521 ) if history_candidates_are_valid(candidates, query.limit) => {
522 ControlEventKind::HistoryDiscovered {
523 count: candidates.len(),
524 }
525 }
526 (ControlObservation::HistoryLoaded { session }, HistoryOperation::Load { .. })
527 if session.validate().is_ok() =>
528 {
529 ControlEventKind::HistoryLoaded {
530 session_id: session.session_id.clone(),
531 }
532 }
533 (ControlObservation::HistoryFailed { message }, _)
534 if validate_history_error(message).is_ok() =>
535 {
536 ControlEventKind::HistoryFailed {
537 message: message.clone(),
538 }
539 }
540 _ => {
541 self.emit_ignored(
542 instance_id,
543 generation,
544 ObservationIgnoredReason::InvalidHistoryObservation,
545 );
546 return;
547 }
548 };
549 let state = self
550 .sessions
551 .get_mut(&instance_id)
552 .expect("validated session");
553 state.snapshot.history.pending = None;
554 match envelope.observation {
555 ControlObservation::HistoryDiscovered { candidates } => {
556 state.snapshot.history.candidates = candidates;
557 state.snapshot.history.loaded_candidate_id = None;
558 state.snapshot.history.loaded = None;
559 state.snapshot.history.last_error = None;
560 }
561 ControlObservation::HistoryLoaded { session } => {
562 state.snapshot.history.loaded_candidate_id = match pending.operation {
563 HistoryOperation::Load { candidate_id } => Some(candidate_id),
564 HistoryOperation::Discover { .. } => {
565 unreachable!("loaded history matched a load operation")
566 }
567 };
568 state.snapshot.history.loaded = Some(session);
569 state.snapshot.history.last_error = None;
570 }
571 ControlObservation::HistoryFailed { message } => {
572 state.snapshot.history.last_error = Some(message);
573 }
574 _ => unreachable!("history observation was matched above"),
575 }
576 self.bump_revision();
577 self.emit_event(None, instance_id, generation, event);
578 return;
579 }
580
581 if is_resume_authority_observation(&envelope.observation) {
582 let Some(pending) = current.resume.pending.as_ref().cloned() else {
583 self.emit_ignored(
584 instance_id,
585 generation,
586 ObservationIgnoredReason::InvalidResumeObservation,
587 );
588 return;
589 };
590 if pending.phase != ResumePhase::Authorizing {
591 self.emit_ignored(
592 instance_id,
593 generation,
594 ObservationIgnoredReason::InvalidResumeObservation,
595 );
596 return;
597 }
598 let operation_id = pending.operation_id;
599 match envelope.observation {
600 ControlObservation::ResumeAuthorized { provider_session }
601 if provider_session.validate().is_ok()
602 && resume_identity_matches_target(
603 current,
604 &pending.target,
605 &provider_session,
606 ) =>
607 {
608 let generation_watermark =
609 self.generation_watermark(instance_id, current.generation);
610 let Some(next_generation) = checked_next_generation(generation_watermark)
611 else {
612 self.emit_ignored(
613 instance_id,
614 generation,
615 ObservationIgnoredReason::GenerationExhausted,
616 );
617 return;
618 };
619 let summary = ResumeSessionSummary::from(&provider_session);
620 let transport = current.transport;
621 let agent_id = {
622 let state = self
623 .sessions
624 .get_mut(&instance_id)
625 .expect("validated session");
626 state.pending_terminal_size = Some(pending.request.terminal_size);
627 state.pending_interrupt = false;
628 state.pending_resume_identity = Some(provider_session.clone());
629 let session = &mut state.snapshot;
630 session.generation = next_generation;
631 session.status = SessionStatus::Starting;
632 session.pending_operation = Some(operation_id);
633 session.pending_input = None;
634 session.process_id = None;
635 session.terminal_size = None;
636 session.terminal_frame = None;
637 session.terminal_stale = None;
638 session.session_options = None;
639 session.history = HistorySnapshot::default();
640 session.foreground = ForegroundSnapshot::default();
641 session.screen_state = PtyScreenState::default();
645 session.provider = ProviderSnapshot::default();
646 session.resume.pending = Some(PendingResumeOperation {
647 phase: ResumePhase::Spawning,
648 ..pending.clone()
649 });
650 session.resume.last_error = None;
651 session.agent_id.clone()
652 };
653 self.generation_watermarks
654 .insert(instance_id, next_generation);
655 self.effects.push(EffectEnvelope {
656 operation_id,
657 instance_id,
658 generation: next_generation,
659 effect: ControlEffect::SpawnResume {
660 agent_id,
661 transport,
662 provider_session,
663 runtime_policy,
664 request: pending.request,
665 },
666 });
667 self.bump_revision();
668 self.emit_event(
669 None,
670 instance_id,
671 next_generation,
672 ControlEventKind::ResumeAuthorized { session: summary },
673 );
674 }
675 ControlObservation::ResumeDenied { reason }
676 if validate_resume_error(&reason).is_ok() =>
677 {
678 let state = self
679 .sessions
680 .get_mut(&instance_id)
681 .expect("validated session");
682 state.pending_resume_identity = None;
683 state.snapshot.pending_operation = None;
684 state.snapshot.resume.pending = None;
685 state.snapshot.resume.last_error = Some(reason.clone());
686 self.bump_revision();
687 self.emit_event(
688 None,
689 instance_id,
690 generation,
691 ControlEventKind::ResumeDenied { reason },
692 );
693 }
694 ControlObservation::ResumeFailed { message }
695 if validate_resume_error(&message).is_ok() =>
696 {
697 let state = self
698 .sessions
699 .get_mut(&instance_id)
700 .expect("validated session");
701 state.pending_resume_identity = None;
702 state.snapshot.pending_operation = None;
703 state.snapshot.resume.pending = None;
704 state.snapshot.resume.last_error = Some(message.clone());
705 self.bump_revision();
706 self.emit_event(
707 None,
708 instance_id,
709 generation,
710 ControlEventKind::ResumeFailed { message },
711 );
712 }
713 _ => {
714 self.emit_ignored(
715 instance_id,
716 generation,
717 ObservationIgnoredReason::InvalidResumeObservation,
718 );
719 }
720 }
721 return;
722 }
723
724 if is_interaction_resolution_observation(&envelope.observation) {
725 let operation_id = envelope
726 .operation_id
727 .expect("interaction resolution requires an operation id");
728 let (interaction_id, failure) = match &envelope.observation {
729 ControlObservation::InteractionResolutionCompleted { interaction_id } => {
730 (*interaction_id, None)
731 }
732 ControlObservation::InteractionResolutionFailed {
733 interaction_id,
734 message,
735 } if interaction_failure_is_valid(message) => {
736 (*interaction_id, Some(message.clone()))
737 }
738 ControlObservation::InteractionResolutionFailed { .. } => {
739 self.emit_ignored(
740 instance_id,
741 generation,
742 ObservationIgnoredReason::InvalidInteractionObservation,
743 );
744 return;
745 }
746 _ => unreachable!("interaction observation was classified above"),
747 };
748 let Some(interaction) = current
749 .provider
750 .interactions
751 .iter()
752 .find(|interaction| interaction.id == interaction_id)
753 else {
754 self.emit_ignored(
755 instance_id,
756 generation,
757 ObservationIgnoredReason::InvalidInteractionObservation,
758 );
759 return;
760 };
761 let ProviderInteractionStatus::Resolving {
762 operation_id: interaction_operation_id,
763 response_kind,
764 } = interaction.status
765 else {
766 self.emit_ignored(
767 instance_id,
768 generation,
769 ObservationIgnoredReason::InvalidInteractionObservation,
770 );
771 return;
772 };
773 if interaction_operation_id != operation_id {
774 self.emit_ignored(
775 instance_id,
776 generation,
777 ObservationIgnoredReason::InvalidInteractionObservation,
778 );
779 return;
780 }
781 let resume_activity = interaction.resume_lead_activity;
782 let state = self
783 .sessions
784 .get_mut(&instance_id)
785 .expect("validated session");
786 state.snapshot.pending_operation = None;
787 let interaction = state
788 .snapshot
789 .provider
790 .interactions
791 .iter_mut()
792 .find(|interaction| interaction.id == interaction_id)
793 .expect("validated interaction");
794 if let Some(message) = failure {
795 interaction.status = ProviderInteractionStatus::Pending;
796 state.snapshot.provider.lead_activity = ProviderActivity::WaitingForInput;
797 refresh_provider_activity(&mut state.snapshot.provider);
798 self.bump_revision();
799 self.emit_event(
800 None,
801 instance_id,
802 generation,
803 ControlEventKind::InteractionResolutionFailed {
804 interaction_id,
805 message,
806 },
807 );
808 } else {
809 let outcome = interaction_response_outcome(response_kind);
810 interaction.status = ProviderInteractionStatus::Resolved { outcome };
811 state.snapshot.provider.lead_activity = if state
812 .snapshot
813 .provider
814 .interactions
815 .iter()
816 .any(interaction_is_unresolved)
817 {
818 ProviderActivity::WaitingForInput
819 } else {
820 resume_activity.unwrap_or(ProviderActivity::Working)
821 };
822 refresh_provider_activity(&mut state.snapshot.provider);
823 self.bump_revision();
824 self.emit_event(
825 None,
826 instance_id,
827 generation,
828 ControlEventKind::InteractionResolved {
829 interaction_id,
830 outcome,
831 },
832 );
833 }
834 return;
835 }
836
837 if let ControlObservation::TerminalFrame { frame } = &envelope.observation {
838 if current
839 .terminal_frame
840 .as_ref()
841 .is_some_and(|existing| existing.sequence >= frame.sequence)
842 {
843 self.emit_ignored(
844 instance_id,
845 generation,
846 ObservationIgnoredReason::StaleTerminalFrame,
847 );
848 return;
849 }
850 let state = self
851 .sessions
852 .get_mut(&instance_id)
853 .expect("validated session");
854 state.snapshot.terminal_size = Some(frame.size);
855 state.snapshot.terminal_frame = Some(frame.clone());
856 state.snapshot.terminal_stale = None;
857 self.bump_revision();
858 return;
859 }
860 if let ControlObservation::TerminalStale { message } = &envelope.observation {
861 let state = self
862 .sessions
863 .get_mut(&instance_id)
864 .expect("validated session");
865 state.snapshot.terminal_stale = Some(message.clone());
866 invalidate_foreground(&mut state.snapshot, message.clone());
867 self.bump_revision();
868 self.emit_event(
869 None,
870 instance_id,
871 generation,
872 ControlEventKind::TerminalStale {
873 message: message.clone(),
874 },
875 );
876 return;
877 }
878 if let ControlObservation::ScreenState { state: screen_state } = &envelope.observation {
879 let state = self
885 .sessions
886 .get_mut(&instance_id)
887 .expect("validated session");
888 state.snapshot.screen_state = screen_state.clone();
889 self.bump_revision();
890 return;
891 }
892 if let ControlObservation::ProviderEvent {
893 source,
894 sequence,
895 event,
896 } = &envelope.observation
897 {
898 if current.provider.sequence == u64::MAX {
899 self.counter_error = Some(ControlError::ProviderSequenceExhausted {
900 instance_id,
901 generation,
902 });
903 return;
904 }
905 if provider_source_sequence(¤t.provider, source) >= *sequence {
906 self.emit_ignored(
907 instance_id,
908 generation,
909 ObservationIgnoredReason::StaleProviderEvent,
910 );
911 return;
912 }
913 if let Some(capability) = provider_event_denied_capability(runtime_policy, event) {
914 let state = self
930 .sessions
931 .get_mut(&instance_id)
932 .expect("validated session");
933 let canonical_sequence =
934 reduce_provider_gap(&mut state.snapshot.provider, source, *sequence, 1);
935 self.bump_revision();
936 self.emit_event(
937 None,
938 instance_id,
939 generation,
940 ControlEventKind::ProviderEvent {
941 sequence: canonical_sequence,
942 source: source.clone(),
943 source_sequence: *sequence,
944 event: ProviderEvent::Error {
945 message: provider_refusal_message(capability, 1),
946 },
947 },
948 );
949 return;
950 }
951 let state = self
952 .sessions
953 .get_mut(&instance_id)
954 .expect("validated session");
955 let pending_resolution = pending_interaction_resolution(&state.snapshot);
956 let reduction = reduce_provider_event(
957 &mut state.snapshot.provider,
958 source,
959 *sequence,
960 event.clone(),
961 );
962 let superseded_resolution =
963 pending_resolution.filter(|(operation_id, interaction_id)| {
964 !interaction_resolution_is_pending(
965 &state.snapshot,
966 *operation_id,
967 *interaction_id,
968 )
969 });
970 if superseded_resolution.is_some() {
971 state.snapshot.pending_operation = None;
972 }
973 if let Some((operation_id, _)) = superseded_resolution {
974 self.effects.retain(|effect| {
975 effect.instance_id != instance_id
976 || effect.generation != generation
977 || effect.operation_id != operation_id
978 });
979 }
980 self.bump_revision();
981 self.emit_event(
982 None,
983 instance_id,
984 generation,
985 ControlEventKind::ProviderEvent {
986 sequence: reduction.sequence,
987 source: source.clone(),
988 source_sequence: *sequence,
989 event: event.clone(),
990 },
991 );
992 self.emit_interaction_transitions(
993 None,
994 instance_id,
995 generation,
996 reduction.interaction_transitions,
997 );
998 return;
999 }
1000 if let ControlObservation::ProviderGap {
1001 source,
1002 source_sequence,
1003 missed,
1004 } = &envelope.observation
1005 {
1006 if current.provider.sequence == u64::MAX {
1007 self.counter_error = Some(ControlError::ProviderSequenceExhausted {
1008 instance_id,
1009 generation,
1010 });
1011 return;
1012 }
1013 let current_source_sequence = provider_source_sequence(¤t.provider, source);
1014 let expected_source_sequence = current_source_sequence.checked_add(*missed);
1015 if *source_sequence == 0
1016 || *missed == 0
1017 || expected_source_sequence != Some(*source_sequence)
1018 {
1019 if expected_source_sequence.is_none() {
1020 self.counter_error = Some(ControlError::ProviderSourceSequenceExhausted {
1021 instance_id,
1022 generation,
1023 provider_source: source.clone(),
1024 });
1025 } else {
1026 self.emit_ignored(
1027 instance_id,
1028 generation,
1029 ObservationIgnoredReason::StaleProviderEvent,
1030 );
1031 }
1032 return;
1033 }
1034 let missed = *missed;
1035 let state = self
1036 .sessions
1037 .get_mut(&instance_id)
1038 .expect("validated session");
1039 let canonical_sequence = reduce_provider_gap(
1040 &mut state.snapshot.provider,
1041 source,
1042 *source_sequence,
1043 missed,
1044 );
1045 self.bump_revision();
1046 self.emit_event(
1047 None,
1048 instance_id,
1049 generation,
1050 ControlEventKind::ProviderGap {
1051 sequence: canonical_sequence,
1052 source: source.clone(),
1053 source_sequence: *source_sequence,
1054 missed,
1055 },
1056 );
1057 return;
1058 }
1059
1060 let pending_input = current.pending_input;
1061 let pending_interrupt = self
1062 .sessions
1063 .get(&instance_id)
1064 .expect("validated session")
1065 .pending_interrupt;
1066 let pending_resume = current.resume.pending.clone();
1067 let pending_resume_identity = self
1068 .sessions
1069 .get(&instance_id)
1070 .expect("validated session")
1071 .pending_resume_identity
1072 .clone();
1073 let mut interaction_transitions = Vec::new();
1074 let event = match envelope.observation {
1075 ControlObservation::Spawned { process_id } => {
1076 let state = self
1077 .sessions
1078 .get_mut(&instance_id)
1079 .expect("validated session");
1080 state.pending_interrupt = false;
1081 state.pending_resume_identity = None;
1082 let terminal_size = state.pending_terminal_size.take();
1083 let session = &mut state.snapshot;
1084 session.status = SessionStatus::Running;
1085 session.pending_operation = None;
1086 session.pending_input = None;
1087 session.process_id = process_id;
1088 session.terminal_size = terminal_size;
1089 if pending_resume
1090 .as_ref()
1091 .is_some_and(|pending| pending.phase == ResumePhase::Spawning)
1092 {
1093 let identity = pending_resume_identity
1094 .as_ref()
1095 .expect("resume spawn must retain its authorized identity");
1096 let summary = ResumeSessionSummary::from(identity);
1097 session.provider.session = Some(identity.clone());
1098 session.resume.pending = None;
1099 session.resume.last_session = Some(summary.clone());
1100 session.resume.last_error = None;
1101 ControlEventKind::Resumed {
1102 session: summary,
1103 process_id,
1104 }
1105 } else {
1106 ControlEventKind::Running { process_id }
1107 }
1108 }
1109 ControlObservation::SpawnFailed { message } => {
1110 let state = self
1111 .sessions
1112 .get_mut(&instance_id)
1113 .expect("validated session");
1114 state.pending_terminal_size = None;
1115 state.pending_interrupt = false;
1116 state.pending_resume_identity = None;
1117 let session = &mut state.snapshot;
1118 session.status = SessionStatus::Failed {
1119 message: message.clone(),
1120 };
1121 session.pending_operation = None;
1122 session.pending_input = None;
1123 session.process_id = None;
1124 session.foreground = ForegroundSnapshot::default();
1125 interaction_transitions = resolve_all_pending_interactions(
1126 &mut session.provider,
1127 ProviderInteractionOutcome::TurnEnded,
1128 );
1129 if pending_resume
1130 .as_ref()
1131 .is_some_and(|pending| pending.phase == ResumePhase::Spawning)
1132 {
1133 session.provider.session = pending_resume_identity.clone();
1134 session.resume.pending = None;
1135 session.resume.last_error = Some(message.clone());
1136 ControlEventKind::ResumeFailed { message }
1137 } else {
1138 ControlEventKind::Failed { message }
1139 }
1140 }
1141 ControlObservation::ProcessExited {
1142 exit_code,
1143 final_terminal,
1144 } => {
1145 self.effects.retain(|effect| {
1146 effect.instance_id != instance_id || effect.generation != generation
1147 });
1148 let state = self
1149 .sessions
1150 .get_mut(&instance_id)
1151 .expect("validated session");
1152 state.pending_terminal_size = None;
1153 state.pending_interrupt = false;
1154 state.pending_resume_identity = None;
1155 let session = &mut state.snapshot;
1156 session.status = SessionStatus::Exited { exit_code };
1157 session.pending_operation = None;
1158 session.pending_input = None;
1159 session.process_id = None;
1160 session.foreground = ForegroundSnapshot::default();
1161 session.resume.pending = None;
1162 interaction_transitions = resolve_all_pending_interactions(
1163 &mut session.provider,
1164 ProviderInteractionOutcome::TurnEnded,
1165 );
1166 if let Some(frame) = final_terminal.filter(|frame| {
1167 session
1168 .terminal_frame
1169 .as_ref()
1170 .is_none_or(|existing| existing.sequence < frame.sequence)
1171 }) {
1172 session.terminal_size = Some(frame.size);
1173 session.terminal_frame = Some(frame);
1174 session.terminal_stale = None;
1175 }
1176 ControlEventKind::Exited {
1177 exit_code,
1178 forced: false,
1179 }
1180 }
1181 ControlObservation::StopCompleted {
1182 forced,
1183 exit_code,
1184 final_terminal,
1185 } => {
1186 let state = self
1187 .sessions
1188 .get_mut(&instance_id)
1189 .expect("validated session");
1190 state.pending_terminal_size = None;
1191 state.pending_interrupt = false;
1192 state.pending_resume_identity = None;
1193 let session = &mut state.snapshot;
1194 session.status = SessionStatus::Exited { exit_code };
1195 session.pending_operation = None;
1196 session.pending_input = None;
1197 session.process_id = None;
1198 session.foreground = ForegroundSnapshot::default();
1199 session.resume.pending = None;
1200 interaction_transitions = resolve_all_pending_interactions(
1201 &mut session.provider,
1202 ProviderInteractionOutcome::TurnEnded,
1203 );
1204 if let Some(frame) = final_terminal.filter(|frame| {
1205 session
1206 .terminal_frame
1207 .as_ref()
1208 .is_none_or(|existing| existing.sequence < frame.sequence)
1209 }) {
1210 session.terminal_size = Some(frame.size);
1211 session.terminal_frame = Some(frame);
1212 session.terminal_stale = None;
1213 }
1214 ControlEventKind::Exited { exit_code, forced }
1215 }
1216 ControlObservation::StopFailed { message } => {
1217 let state = self
1218 .sessions
1219 .get_mut(&instance_id)
1220 .expect("validated session");
1221 state.pending_terminal_size = None;
1222 state.pending_interrupt = false;
1223 state.pending_resume_identity = None;
1224 let session = &mut state.snapshot;
1225 session.status = SessionStatus::Failed {
1226 message: message.clone(),
1227 };
1228 session.pending_operation = None;
1229 session.pending_input = None;
1230 session.process_id = None;
1231 session.foreground = ForegroundSnapshot::default();
1232 session.resume.pending = None;
1233 interaction_transitions = resolve_all_pending_interactions(
1234 &mut session.provider,
1235 ProviderInteractionOutcome::TurnEnded,
1236 );
1237 ControlEventKind::Failed { message }
1238 }
1239 ControlObservation::InputCompleted => {
1240 let input_kind = pending_input
1241 .expect("validated input completion must have a pending input kind");
1242 let state = self
1243 .sessions
1244 .get_mut(&instance_id)
1245 .expect("validated session");
1246 state.pending_interrupt = false;
1247 let session = &mut state.snapshot;
1248 session.pending_operation = None;
1249 session.pending_input = None;
1250 if pending_interrupt {
1251 interaction_transitions = resolve_all_pending_interactions(
1252 &mut session.provider,
1253 ProviderInteractionOutcome::Interrupted,
1254 );
1255 session.provider.lead_activity = ProviderActivity::Idle;
1256 refresh_provider_activity(&mut session.provider);
1257 session.provider.current_prompt = None;
1258 session.provider.active_tools.clear();
1259 }
1260 ControlEventKind::InputCompleted { input_kind }
1261 }
1262 ControlObservation::InputFailed { message } => {
1263 let input_kind =
1264 pending_input.expect("validated input failure must have a pending input kind");
1265 let state = self
1266 .sessions
1267 .get_mut(&instance_id)
1268 .expect("validated session");
1269 state.pending_interrupt = false;
1270 let session = &mut state.snapshot;
1271 session.pending_operation = None;
1272 session.pending_input = None;
1273 ControlEventKind::InputFailed {
1274 input_kind,
1275 message,
1276 }
1277 }
1278 ControlObservation::ResizeCompleted { size } => {
1279 let state = self
1280 .sessions
1281 .get_mut(&instance_id)
1282 .expect("validated session");
1283 state.pending_terminal_size = None;
1284 let session = &mut state.snapshot;
1285 session.pending_operation = None;
1286 session.terminal_size = Some(size);
1287 ControlEventKind::Resized { size }
1288 }
1289 ControlObservation::ResizeFailed { message } => {
1290 let state = self
1291 .sessions
1292 .get_mut(&instance_id)
1293 .expect("validated session");
1294 state.pending_terminal_size = None;
1295 let session = &mut state.snapshot;
1296 session.pending_operation = None;
1297 ControlEventKind::ResizeFailed { message }
1298 }
1299 ControlObservation::ForegroundObserved { process } => {
1300 let session = self.session_mut(instance_id);
1301 session.pending_operation = None;
1302 session.foreground = ForegroundSnapshot {
1303 authority: ForegroundAuthority::Confirmed,
1304 process: Some(process.clone()),
1305 stale_reason: None,
1306 };
1307 ControlEventKind::ForegroundObserved { process }
1308 }
1309 ControlObservation::ForegroundFailed { message } => {
1310 let session = self.session_mut(instance_id);
1311 session.pending_operation = None;
1312 invalidate_foreground(session, message.clone());
1313 ControlEventKind::ForegroundFailed { message }
1314 }
1315 ControlObservation::SessionModeSet { mode_id } => {
1316 let session = self.session_mut(instance_id);
1317 session.pending_operation = None;
1318 ControlEventKind::SessionModeSet { mode_id }
1319 }
1320 ControlObservation::SessionModeSetFailed { message } => {
1321 let session = self.session_mut(instance_id);
1322 session.pending_operation = None;
1323 ControlEventKind::SessionModeSetFailed { message }
1324 }
1325 ControlObservation::SessionConfigOptionSet { option_id } => {
1326 let session = self.session_mut(instance_id);
1327 session.pending_operation = None;
1328 ControlEventKind::SessionConfigOptionSet { option_id }
1329 }
1330 ControlObservation::SessionConfigOptionSetFailed { message } => {
1331 let session = self.session_mut(instance_id);
1332 session.pending_operation = None;
1333 ControlEventKind::SessionConfigOptionSetFailed { message }
1334 }
1335 ControlObservation::SessionModelSet { model_id } => {
1336 let session = self.session_mut(instance_id);
1337 session.pending_operation = None;
1338 ControlEventKind::SessionModelSet { model_id }
1339 }
1340 ControlObservation::SessionModelSetFailed { message } => {
1341 let session = self.session_mut(instance_id);
1342 session.pending_operation = None;
1343 ControlEventKind::SessionModelSetFailed { message }
1344 }
1345 ControlObservation::TerminalFrame { .. }
1346 | ControlObservation::TerminalStale { .. }
1347 | ControlObservation::ScreenState { .. }
1348 | ControlObservation::ProviderEvent { .. }
1349 | ControlObservation::ProviderGap { .. }
1350 | ControlObservation::CapabilitiesProbed { .. }
1351 | ControlObservation::CapabilityProbeFailed { .. }
1352 | ControlObservation::HistoryDiscovered { .. }
1353 | ControlObservation::HistoryLoaded { .. }
1354 | ControlObservation::HistoryFailed { .. }
1355 | ControlObservation::ResumeAuthorized { .. }
1356 | ControlObservation::ResumeDenied { .. }
1357 | ControlObservation::ResumeFailed { .. }
1358 | ControlObservation::InteractionResolutionCompleted { .. }
1359 | ControlObservation::InteractionResolutionFailed { .. } => {
1360 unreachable!("stream observations return before lifecycle event reduction")
1361 }
1362 };
1363 self.bump_revision();
1364 self.emit_event(None, instance_id, generation, event);
1365 self.emit_interaction_transitions(None, instance_id, generation, interaction_transitions);
1366 }
1367
1368 pub fn snapshot(&self) -> ControlSnapshot {
1369 ControlSnapshot {
1370 revision: self.revision,
1371 health: self.health(),
1372 sessions: self
1373 .sessions
1374 .values()
1375 .map(|state| state.snapshot.clone())
1376 .collect(),
1377 }
1378 }
1379
1380 pub fn health(&self) -> ControlHealth {
1381 ControlHealth {
1382 operation_id_exhausted: self.next_operation_id.is_none(),
1383 event_sequence_exhausted: self.next_event_sequence.is_none(),
1384 revision_exhausted: self.revision == u64::MAX,
1385 provider_sequence_exhausted_sessions: u32::try_from(
1386 self.sessions
1387 .values()
1388 .filter(|state| state.snapshot.provider.sequence == u64::MAX)
1389 .count(),
1390 )
1391 .expect("live session map is bounded below u32::MAX"),
1392 retained_instance_identities: u32::try_from(self.generation_watermarks.len())
1393 .expect("retained identity map is bounded below u32::MAX"),
1394 retained_instance_identity_capacity: CONTROL_INSTANCE_IDENTITIES_CAPACITY,
1395 }
1396 }
1397
1398 pub fn session_snapshot(&self, instance_id: AgentInstanceId) -> Option<&SessionSnapshot> {
1399 self.sessions.get(&instance_id).map(|state| &state.snapshot)
1400 }
1401
1402 pub fn session_instance_ids(&self) -> impl Iterator<Item = AgentInstanceId> + '_ {
1403 self.sessions.keys().copied()
1404 }
1405
1406 pub fn record_command_rejection(
1407 &mut self,
1408 command_id: CommandId,
1409 instance_id: AgentInstanceId,
1410 message: String,
1411 ) {
1412 if self.next_event_sequence.is_none() {
1413 return;
1414 }
1415 let generation = self
1416 .sessions
1417 .get(&instance_id)
1418 .map(|state| state.snapshot.generation)
1419 .or_else(|| self.generation_watermarks.get(&instance_id).copied())
1420 .unwrap_or_default();
1421 self.emit_event(
1422 Some(command_id),
1423 instance_id,
1424 generation,
1425 ControlEventKind::CommandRejected { message },
1426 );
1427 debug_assert!(self.counter_error.is_none());
1428 }
1429
1430 pub fn drain_effects(&mut self) -> Vec<EffectEnvelope> {
1431 std::mem::take(&mut self.effects)
1432 }
1433
1434 pub fn drain_events(&mut self) -> Vec<ControlEvent> {
1435 std::mem::take(&mut self.events)
1436 }
1437
1438 fn start(
1439 &mut self,
1440 command_id: CommandId,
1441 instance_id: AgentInstanceId,
1442 runtime_policy: ProviderRuntimePolicy,
1443 mut request: StartRequest,
1444 ) -> Result<(), ControlError> {
1445 runtime_policy
1446 .validate()
1447 .map_err(|error| ControlError::InvalidProviderRuntimePolicy { error })?;
1448 if !request.terminal_size.is_valid() {
1449 return Err(ControlError::InvalidTerminalSize);
1450 }
1451 if request.working_directory.is_empty()
1452 || request.working_directory.len() > WORKING_DIRECTORY_MAX_BYTES
1453 || request.working_directory.contains('\0')
1454 {
1455 return Err(ControlError::InvalidWorkingDirectory);
1456 }
1457 let session = self
1458 .sessions
1459 .get(&instance_id)
1460 .ok_or(ControlError::UnknownInstance { instance_id })?;
1461 let status = session.snapshot.status.clone();
1462 let transport = session.snapshot.transport;
1463 if transport == TransportKind::Pty {
1464 require_runtime_capability(runtime_policy, ProviderRuntimeCapability::RawPtyLifecycle)?;
1465 }
1466 if let Some(operation_id) = session.snapshot.pending_operation {
1467 return Err(ControlError::OperationPending {
1468 instance_id,
1469 operation_id,
1470 });
1471 }
1472 if !matches!(
1473 status,
1474 SessionStatus::Registered | SessionStatus::Exited { .. } | SessionStatus::Failed { .. }
1475 ) {
1476 return Err(ControlError::InvalidTransition {
1477 instance_id,
1478 action: "start".to_owned(),
1479 status,
1480 });
1481 }
1482
1483 request.initial_prompt = request
1484 .initial_prompt
1485 .as_deref()
1486 .map(normalize_semantic_prompt)
1487 .transpose()
1488 .map_err(|error| ControlError::InputRejected { error })?;
1489 if let Some(session_options) = &request.session_options {
1490 session_options
1491 .validate()
1492 .map_err(|error| ControlError::InvalidSessionOptions {
1493 message: error.to_string(),
1494 })?;
1495 }
1496 if transport == TransportKind::Pipe
1497 && request.initial_prompt.as_deref().is_none_or(str::is_empty)
1498 {
1499 return Err(ControlError::MissingInitialPrompt);
1500 }
1501 if request.initial_prompt.as_deref().is_some_and(|prompt| !prompt.is_empty()) {
1502 require_structured_prompt_policy(runtime_policy)?;
1503 }
1504
1505 let generation_watermark =
1506 self.generation_watermark(instance_id, session.snapshot.generation);
1507 let generation = checked_next_generation(generation_watermark).ok_or(
1508 ControlError::GenerationExhausted {
1509 instance_id,
1510 generation: generation_watermark,
1511 },
1512 )?;
1513 self.purge_generation_bound_history(instance_id);
1514 let operation_id = self.allocate_operation();
1515 let (agent_id, transport) = {
1516 let state = self
1517 .sessions
1518 .get_mut(&instance_id)
1519 .expect("validated session");
1520 state.pending_terminal_size = Some(request.terminal_size);
1521 state.pending_interrupt = false;
1522 state.pending_resume_identity = None;
1523 state.runtime_policy = runtime_policy;
1524 let session = &mut state.snapshot;
1525 session.history = HistorySnapshot::default();
1526 session.resume = ResumeSnapshot::default();
1527 session.generation = generation;
1528 session.status = SessionStatus::Starting;
1529 session.pending_operation = Some(operation_id);
1530 session.pending_input = None;
1531 session.process_id = None;
1532 session.terminal_size = None;
1533 session.terminal_frame = None;
1534 session.terminal_stale = None;
1535 session.session_options = request.session_options.clone();
1536 session.foreground = ForegroundSnapshot::default();
1537 session.screen_state = PtyScreenState::default();
1541 session.provider = ProviderSnapshot::default();
1542 (session.agent_id.clone(), session.transport)
1543 };
1544 self.generation_watermarks.insert(instance_id, generation);
1545 self.effects.push(EffectEnvelope {
1546 operation_id,
1547 instance_id,
1548 generation,
1549 effect: ControlEffect::Spawn {
1550 agent_id,
1551 transport,
1552 runtime_policy,
1553 request,
1554 },
1555 });
1556 self.bump_revision();
1557 self.emit_event(
1558 Some(command_id),
1559 instance_id,
1560 generation,
1561 ControlEventKind::StartRequested { operation_id },
1562 );
1563 Ok(())
1564 }
1565
1566 fn stop(
1567 &mut self,
1568 command_id: CommandId,
1569 instance_id: AgentInstanceId,
1570 force: bool,
1571 ) -> Result<(), ControlError> {
1572 let status = self
1573 .sessions
1574 .get(&instance_id)
1575 .ok_or(ControlError::UnknownInstance { instance_id })?
1576 .snapshot
1577 .status
1578 .clone();
1579 let pending_resolution = self
1580 .sessions
1581 .get(&instance_id)
1582 .and_then(|state| pending_interaction_resolution(&state.snapshot));
1583 if !matches!(status, SessionStatus::Starting | SessionStatus::Running) {
1584 return Err(ControlError::InvalidTransition {
1585 instance_id,
1586 action: "stop".to_owned(),
1587 status,
1588 });
1589 }
1590
1591 let operation_id = self.allocate_operation();
1592 let generation = {
1593 let state = self
1594 .sessions
1595 .get_mut(&instance_id)
1596 .expect("validated session");
1597 state.pending_interrupt = false;
1598 let session = &mut state.snapshot;
1599 session.status = SessionStatus::Stopping;
1600 session.pending_operation = Some(operation_id);
1601 session.pending_input = None;
1602 invalidate_foreground(session, "stop requested".to_owned());
1603 session.generation
1604 };
1605 if let Some((pending_operation_id, _)) = pending_resolution {
1606 self.effects.retain(|effect| {
1607 effect.instance_id != instance_id
1608 || effect.generation != generation
1609 || effect.operation_id != pending_operation_id
1610 });
1611 }
1612 self.effects.push(EffectEnvelope {
1613 operation_id,
1614 instance_id,
1615 generation,
1616 effect: ControlEffect::Stop { force },
1617 });
1618 self.bump_revision();
1619 self.emit_event(
1620 Some(command_id),
1621 instance_id,
1622 generation,
1623 ControlEventKind::StopRequested {
1624 operation_id,
1625 force,
1626 },
1627 );
1628 Ok(())
1629 }
1630
1631 fn send_input(
1632 &mut self,
1633 command_id: CommandId,
1634 instance_id: AgentInstanceId,
1635 action: InputAction,
1636 ) -> Result<(), ControlError> {
1637 let state = self
1638 .sessions
1639 .get(&instance_id)
1640 .ok_or(ControlError::UnknownInstance { instance_id })?;
1641 if state.snapshot.status != SessionStatus::Running {
1642 return Err(ControlError::InvalidTransition {
1643 instance_id,
1644 action: "send input".to_owned(),
1645 status: state.snapshot.status.clone(),
1646 });
1647 }
1648 if let Some(operation_id) = state.snapshot.pending_operation {
1649 return Err(ControlError::OperationPending {
1650 instance_id,
1651 operation_id,
1652 });
1653 }
1654
1655 let agent_id = state.snapshot.agent_id.clone();
1656 let transport = state.snapshot.transport;
1657 let runtime_policy = state.runtime_policy;
1658 if matches!(
1659 &action,
1660 InputAction::InsertDraft(_)
1661 | InputAction::SubmitPrompt(_)
1662 | InputAction::AgentCommand(_)
1663 ) {
1664 require_structured_prompt_policy(runtime_policy)?;
1665 }
1666 let interrupt_requested = matches!(
1667 &action,
1668 InputAction::TerminalControl(TerminalControl::Interrupt)
1669 );
1670 let (effect, input_kind) = match (transport, action) {
1671 (TransportKind::Pty, InputAction::AgentCommand(command)) => {
1672 let input = prepare_agent_command(command, &agent_id)
1673 .map_err(|error| ControlError::InputRejected { error })?;
1674 let input_kind = input.kind();
1675 (
1676 ControlEffect::WriteInput {
1677 input,
1678 required_foreground: ForegroundRequirement::Agent { agent_id },
1679 },
1680 input_kind,
1681 )
1682 }
1683 (TransportKind::Pty, InputAction::ShellCommand(command)) => {
1684 let input = prepare_shell_command(command)
1685 .map_err(|error| ControlError::InputRejected { error })?;
1686 let input_kind = input.kind();
1687 (
1688 ControlEffect::WriteInput {
1689 input,
1690 required_foreground: ForegroundRequirement::Shell,
1691 },
1692 input_kind,
1693 )
1694 }
1695 (TransportKind::Pty, action) => {
1696 let input =
1697 prepare_input(action).map_err(|error| ControlError::InputRejected { error })?;
1698 let input_kind = input.kind();
1699 let required_foreground = match input_kind {
1700 PreparedInputKind::InsertDraft | PreparedInputKind::SubmitPrompt => {
1701 ForegroundRequirement::Agent { agent_id }
1702 }
1703 PreparedInputKind::TerminalText
1704 | PreparedInputKind::TerminalBytes
1705 | PreparedInputKind::TerminalControl => {
1706 ForegroundRequirement::Any
1707 }
1708 PreparedInputKind::AgentCommand | PreparedInputKind::ShellCommand => {
1709 unreachable!("dispatcher-only input cannot be prepared generically")
1710 }
1711 };
1712 (
1713 ControlEffect::WriteInput {
1714 input,
1715 required_foreground,
1716 },
1717 input_kind,
1718 )
1719 }
1720 (TransportKind::Acp, InputAction::SubmitPrompt(prompt)) => {
1721 let prompt = normalize_semantic_prompt(&prompt.text)
1722 .map_err(|error| ControlError::InputRejected { error })?;
1723 (
1724 ControlEffect::SubmitPrompt { prompt },
1725 PreparedInputKind::SubmitPrompt,
1726 )
1727 }
1728 (TransportKind::Acp, InputAction::TerminalControl(TerminalControl::Interrupt)) => {
1729 (ControlEffect::Interrupt, PreparedInputKind::TerminalControl)
1730 }
1731 (transport, _) => {
1732 return Err(ControlError::UnsupportedTransportOperation {
1733 transport,
1734 action: "this input action".to_owned(),
1735 });
1736 }
1737 };
1738
1739 let operation_id = self.allocate_operation();
1740 let generation = {
1741 let state = self
1742 .sessions
1743 .get_mut(&instance_id)
1744 .expect("validated session");
1745 state.pending_interrupt = interrupt_requested;
1746 let session = &mut state.snapshot;
1747 session.pending_operation = Some(operation_id);
1748 session.pending_input = Some(input_kind);
1749 invalidate_foreground(session, "input requested".to_owned());
1750 session.generation
1751 };
1752 self.effects.push(EffectEnvelope {
1753 operation_id,
1754 instance_id,
1755 generation,
1756 effect,
1757 });
1758 self.bump_revision();
1759 self.emit_event(
1760 Some(command_id),
1761 instance_id,
1762 generation,
1763 ControlEventKind::InputRequested {
1764 operation_id,
1765 input_kind,
1766 },
1767 );
1768 Ok(())
1769 }
1770
1771 fn resize(
1772 &mut self,
1773 command_id: CommandId,
1774 instance_id: AgentInstanceId,
1775 size: TerminalSize,
1776 ) -> Result<(), ControlError> {
1777 if !size.is_valid() {
1778 return Err(ControlError::InvalidTerminalSize);
1779 }
1780 let state = self
1781 .sessions
1782 .get(&instance_id)
1783 .ok_or(ControlError::UnknownInstance { instance_id })?;
1784 if state.snapshot.transport != TransportKind::Pty {
1785 return Err(ControlError::UnsupportedTransportOperation {
1786 transport: state.snapshot.transport,
1787 action: "terminal resize".to_owned(),
1788 });
1789 }
1790 if state.snapshot.status != SessionStatus::Running {
1791 return Err(ControlError::InvalidTransition {
1792 instance_id,
1793 action: "resize".to_owned(),
1794 status: state.snapshot.status.clone(),
1795 });
1796 }
1797 if let Some(operation_id) = state.snapshot.pending_operation {
1798 return Err(ControlError::OperationPending {
1799 instance_id,
1800 operation_id,
1801 });
1802 }
1803
1804 let operation_id = self.allocate_operation();
1805 let generation = {
1806 let state = self
1807 .sessions
1808 .get_mut(&instance_id)
1809 .expect("validated session");
1810 state.snapshot.pending_operation = Some(operation_id);
1811 state.pending_terminal_size = Some(size);
1812 state.snapshot.generation
1813 };
1814 self.effects.push(EffectEnvelope {
1815 operation_id,
1816 instance_id,
1817 generation,
1818 effect: ControlEffect::Resize { size },
1819 });
1820 self.bump_revision();
1821 self.emit_event(
1822 Some(command_id),
1823 instance_id,
1824 generation,
1825 ControlEventKind::ResizeRequested { operation_id, size },
1826 );
1827 Ok(())
1828 }
1829
1830 fn refresh_foreground(
1831 &mut self,
1832 command_id: CommandId,
1833 instance_id: AgentInstanceId,
1834 ) -> Result<(), ControlError> {
1835 let state = self
1836 .sessions
1837 .get(&instance_id)
1838 .ok_or(ControlError::UnknownInstance { instance_id })?;
1839 if state.snapshot.transport != TransportKind::Pty {
1840 return Err(ControlError::UnsupportedTransportOperation {
1841 transport: state.snapshot.transport,
1842 action: "foreground refresh".to_owned(),
1843 });
1844 }
1845 if state.snapshot.status != SessionStatus::Running {
1846 return Err(ControlError::InvalidTransition {
1847 instance_id,
1848 action: "refresh foreground".to_owned(),
1849 status: state.snapshot.status.clone(),
1850 });
1851 }
1852 if let Some(operation_id) = state.snapshot.pending_operation {
1853 return Err(ControlError::OperationPending {
1854 instance_id,
1855 operation_id,
1856 });
1857 }
1858
1859 let operation_id = self.allocate_operation();
1860 let generation = {
1861 let session = self.session_mut(instance_id);
1862 session.pending_operation = Some(operation_id);
1863 invalidate_foreground(session, "foreground refresh pending".to_owned());
1864 session.generation
1865 };
1866 self.effects.push(EffectEnvelope {
1867 operation_id,
1868 instance_id,
1869 generation,
1870 effect: ControlEffect::ObserveForeground,
1871 });
1872 self.bump_revision();
1873 self.emit_event(
1874 Some(command_id),
1875 instance_id,
1876 generation,
1877 ControlEventKind::ForegroundRefreshRequested { operation_id },
1878 );
1879 Ok(())
1880 }
1881
1882 fn discover_history(
1883 &mut self,
1884 command_id: CommandId,
1885 instance_id: AgentInstanceId,
1886 query: HistoryQuery,
1887 ) -> Result<(), ControlError> {
1888 query
1889 .validate()
1890 .map_err(|error| ControlError::InvalidHistoryRequest {
1891 message: error.to_string(),
1892 })?;
1893 let state = self
1894 .sessions
1895 .get(&instance_id)
1896 .ok_or(ControlError::UnknownInstance { instance_id })?;
1897 if let Some(pending) = &state.snapshot.history.pending {
1898 return Err(ControlError::HistoryOperationPending {
1899 operation_id: pending.operation_id,
1900 });
1901 }
1902 let operation_id = self.allocate_operation();
1903 let (generation, agent_id, operation) = {
1904 let session = self.session_mut(instance_id);
1905 let operation = HistoryOperation::Discover {
1906 query: query.clone(),
1907 };
1908 session.history.pending = Some(PendingHistoryOperation {
1909 operation_id,
1910 operation: operation.clone(),
1911 });
1912 session.history.candidates.clear();
1913 session.history.loaded_candidate_id = None;
1914 session.history.loaded = None;
1915 session.history.last_error = None;
1916 (session.generation, session.agent_id.clone(), operation)
1917 };
1918 self.effects.push(EffectEnvelope {
1919 operation_id,
1920 instance_id,
1921 generation,
1922 effect: ControlEffect::DiscoverHistory { agent_id, query },
1923 });
1924 self.bump_revision();
1925 self.emit_event(
1926 Some(command_id),
1927 instance_id,
1928 generation,
1929 ControlEventKind::HistoryRequested {
1930 operation_id,
1931 operation,
1932 },
1933 );
1934 Ok(())
1935 }
1936
1937 fn probe_capabilities(
1938 &mut self,
1939 command_id: CommandId,
1940 instance_id: AgentInstanceId,
1941 request: CapabilityProbeRequest,
1942 ) -> Result<(), ControlError> {
1943 request
1944 .validate()
1945 .map_err(|error| ControlError::InvalidCapabilityProbeRequest {
1946 message: error.to_string(),
1947 })?;
1948 let state = self
1949 .sessions
1950 .get(&instance_id)
1951 .ok_or(ControlError::UnknownInstance { instance_id })?;
1952 if let Some(pending) = &state.snapshot.capabilities.pending {
1953 return Err(ControlError::CapabilityProbeOperationPending {
1954 operation_id: pending.operation_id,
1955 });
1956 }
1957 if state.snapshot.capabilities.settled {
1958 return Err(ControlError::CapabilityProbeSettled);
1959 }
1960
1961 let operation_id = self.allocate_operation();
1962 let (generation, agent_id) = {
1963 let session = self.session_mut(instance_id);
1964 let generation = session.generation;
1965 session.capabilities.pending = Some(PendingCapabilityProbe {
1966 operation_id,
1967 generation,
1968 request: request.clone(),
1969 });
1970 session.capabilities.session_option_models.clear();
1971 session.capabilities.last_failure = None;
1972 (generation, session.agent_id.clone())
1973 };
1974 self.effects.push(EffectEnvelope {
1975 operation_id,
1976 instance_id,
1977 generation,
1978 effect: ControlEffect::ProbeCapabilities { agent_id, request },
1979 });
1980 self.bump_revision();
1981 self.emit_event(
1982 Some(command_id),
1983 instance_id,
1984 generation,
1985 ControlEventKind::CapabilityProbeRequested { operation_id },
1986 );
1987 Ok(())
1988 }
1989
1990 fn load_history(
1991 &mut self,
1992 command_id: CommandId,
1993 instance_id: AgentInstanceId,
1994 candidate_id: String,
1995 ) -> Result<(), ControlError> {
1996 validate_candidate_id(&candidate_id).map_err(|error| {
1997 ControlError::InvalidHistoryRequest {
1998 message: error.to_string(),
1999 }
2000 })?;
2001 let state = self
2002 .sessions
2003 .get(&instance_id)
2004 .ok_or(ControlError::UnknownInstance { instance_id })?;
2005 if let Some(pending) = &state.snapshot.history.pending {
2006 return Err(ControlError::HistoryOperationPending {
2007 operation_id: pending.operation_id,
2008 });
2009 }
2010 if state.snapshot.history.candidate(&candidate_id).is_none() {
2011 return Err(ControlError::UnknownHistoryCandidate);
2012 }
2013 let operation_id = self.allocate_operation();
2014 let (generation, agent_id, operation) = {
2015 let session = self.session_mut(instance_id);
2016 let operation = HistoryOperation::Load {
2017 candidate_id: candidate_id.clone(),
2018 };
2019 session.history.pending = Some(PendingHistoryOperation {
2020 operation_id,
2021 operation: operation.clone(),
2022 });
2023 session.history.loaded = None;
2024 session.history.loaded_candidate_id = None;
2025 session.history.last_error = None;
2026 (session.generation, session.agent_id.clone(), operation)
2027 };
2028 self.effects.push(EffectEnvelope {
2029 operation_id,
2030 instance_id,
2031 generation,
2032 effect: ControlEffect::LoadHistory {
2033 agent_id,
2034 candidate_id,
2035 },
2036 });
2037 self.bump_revision();
2038 self.emit_event(
2039 Some(command_id),
2040 instance_id,
2041 generation,
2042 ControlEventKind::HistoryRequested {
2043 operation_id,
2044 operation,
2045 },
2046 );
2047 Ok(())
2048 }
2049
2050 fn resume(
2051 &mut self,
2052 command_id: CommandId,
2053 instance_id: AgentInstanceId,
2054 target: ResumeTarget,
2055 runtime_policy: ProviderRuntimePolicy,
2056 mut request: ResumeLaunchRequest,
2057 ) -> Result<(), ControlError> {
2058 runtime_policy
2059 .validate()
2060 .map_err(|error| ControlError::InvalidProviderRuntimePolicy { error })?;
2061 require_runtime_capability(
2062 runtime_policy,
2063 ProviderRuntimeCapability::RawPtyLifecycle,
2064 )?;
2065 let has_initial_prompt = request.initial_prompt.is_some();
2066 if has_initial_prompt {
2067 require_runtime_capability(
2068 runtime_policy,
2069 ProviderRuntimeCapability::ProviderSessionIdentity,
2070 )?;
2071 require_runtime_capability(
2072 runtime_policy,
2073 ProviderRuntimeCapability::SemanticResume,
2074 )?;
2075 }
2076 request.initial_prompt = request
2077 .initial_prompt
2078 .as_deref()
2079 .map(normalize_semantic_prompt)
2080 .transpose()
2081 .map_err(|error| ControlError::InvalidResumeRequest {
2082 message: error.to_string(),
2083 })?;
2084 request
2085 .validate()
2086 .map_err(|error| ControlError::InvalidResumeRequest {
2087 message: error.to_string(),
2088 })?;
2089 if request.initial_prompt.as_deref().is_some_and(|prompt| !prompt.is_empty()) {
2090 require_structured_prompt_policy(runtime_policy)?;
2091 }
2092 target
2093 .validate()
2094 .map_err(|error| ControlError::InvalidResumeRequest {
2095 message: error.to_string(),
2096 })?;
2097 let state = self
2098 .sessions
2099 .get(&instance_id)
2100 .ok_or(ControlError::UnknownInstance { instance_id })?;
2101 if !matches!(
2102 state.snapshot.transport,
2103 TransportKind::Pty | TransportKind::Pipe
2104 ) {
2105 return Err(ControlError::UnsupportedTransportOperation {
2106 transport: state.snapshot.transport,
2107 action: "resume".to_owned(),
2108 });
2109 }
2110 if state.snapshot.transport == TransportKind::Pipe
2111 && request.initial_prompt.as_deref().is_none_or(str::is_empty)
2112 {
2113 return Err(ControlError::MissingInitialPrompt);
2114 }
2115 let status = state.snapshot.status.clone();
2116 if !matches!(
2117 status,
2118 SessionStatus::Registered | SessionStatus::Exited { .. } | SessionStatus::Failed { .. }
2119 ) {
2120 return Err(ControlError::InvalidTransition {
2121 instance_id,
2122 action: "resume".to_owned(),
2123 status,
2124 });
2125 }
2126 if let Some(operation_id) = state.snapshot.pending_operation {
2127 return Err(ControlError::OperationPending {
2128 instance_id,
2129 operation_id,
2130 });
2131 }
2132
2133 let authority_target = match &target {
2134 ResumeTarget::CurrentProvider => {
2135 let identity = state
2136 .snapshot
2137 .provider
2138 .session
2139 .clone()
2140 .ok_or(ControlError::MissingProviderSession)?;
2141 identity
2142 .validate()
2143 .map_err(|error| ControlError::InvalidResumeRequest {
2144 message: error.to_string(),
2145 })?;
2146 ResumeAuthorityTarget::ProviderSession { identity }
2147 }
2148 ResumeTarget::ProviderSession { identity } => {
2149 identity
2150 .validate()
2151 .map_err(|error| ControlError::InvalidResumeRequest {
2152 message: error.to_string(),
2153 })?;
2154 ResumeAuthorityTarget::ProviderSession {
2155 identity: identity.clone(),
2156 }
2157 }
2158 ResumeTarget::HistoryCandidate { candidate_id } => {
2159 if state.snapshot.history.candidate(candidate_id).is_none()
2160 || state.snapshot.history.loaded_candidate_id.as_deref()
2161 != Some(candidate_id.as_str())
2162 || state.snapshot.history.loaded.is_none()
2163 {
2164 return Err(ControlError::HistoryCandidateNotLoaded);
2165 }
2166 ResumeAuthorityTarget::HistoryCandidate {
2167 candidate_id: candidate_id.clone(),
2168 }
2169 }
2170 };
2171
2172 let operation_id = self.allocate_operation();
2173 let (generation, agent_id) = {
2174 let state = self
2175 .sessions
2176 .get_mut(&instance_id)
2177 .expect("validated session");
2178 state.pending_resume_identity = None;
2179 state.runtime_policy = runtime_policy;
2180 let session = &mut state.snapshot;
2181 session.pending_operation = Some(operation_id);
2182 session.resume.pending = Some(PendingResumeOperation {
2183 operation_id,
2184 target: target.clone(),
2185 request: request.clone(),
2186 phase: ResumePhase::Authorizing,
2187 });
2188 session.resume.last_error = None;
2189 (session.generation, session.agent_id.clone())
2190 };
2191 self.effects.push(EffectEnvelope {
2192 operation_id,
2193 instance_id,
2194 generation,
2195 effect: ControlEffect::AuthorizeResume {
2196 agent_id,
2197 target: authority_target,
2198 request,
2199 },
2200 });
2201 self.bump_revision();
2202 self.emit_event(
2203 Some(command_id),
2204 instance_id,
2205 generation,
2206 ControlEventKind::ResumeRequested {
2207 operation_id,
2208 target,
2209 },
2210 );
2211 Ok(())
2212 }
2213
2214 fn ingest_provider(
2215 &mut self,
2216 command_id: CommandId,
2217 instance_id: AgentInstanceId,
2218 generation: SessionGeneration,
2219 source: ProviderSource,
2220 source_sequence: u64,
2221 events: Vec<ProviderEvent>,
2222 ) -> Result<(), ControlError> {
2223 source
2224 .binding
2225 .validate()
2226 .map_err(|error| ControlError::InvalidProviderEvent {
2227 message: error.to_string(),
2228 })?;
2229 if source_sequence == 0 || events.is_empty() || events.len() > PROVIDER_INGRESS_EVENTS_MAX {
2230 return Err(ControlError::InvalidProviderBatch {
2231 max: PROVIDER_INGRESS_EVENTS_MAX,
2232 });
2233 }
2234 for event in &events {
2235 event
2236 .validate_ingress()
2237 .map_err(|error| ControlError::InvalidProviderEvent {
2238 message: error.to_string(),
2239 })?;
2240 }
2241
2242 let state = self
2243 .sessions
2244 .get(&instance_id)
2245 .ok_or(ControlError::UnknownInstance { instance_id })?;
2246 if state.snapshot.generation != generation {
2247 return Err(ControlError::StaleProviderGeneration {
2248 expected: state.snapshot.generation,
2249 actual: generation,
2250 });
2251 }
2252 if !matches!(
2253 state.snapshot.status,
2254 SessionStatus::Starting | SessionStatus::Running | SessionStatus::Stopping
2255 ) {
2256 return Err(ControlError::InvalidTransition {
2257 instance_id,
2258 action: "ingest provider events".to_owned(),
2259 status: state.snapshot.status.clone(),
2260 });
2261 }
2262 let current_source_sequence = provider_source_sequence(&state.snapshot.provider, &source);
2263 if current_source_sequence == u64::MAX {
2264 return Err(ControlError::ProviderSourceSequenceExhausted {
2265 instance_id,
2266 generation,
2267 provider_source: source,
2268 });
2269 }
2270 if source_sequence <= current_source_sequence {
2271 return Err(ControlError::StaleProviderSequence);
2272 }
2273
2274 let missed = source_sequence
2275 .checked_sub(current_source_sequence)
2276 .and_then(|difference| difference.checked_sub(1))
2277 .expect("source sequence ordering was validated");
2278 let mut denied_capability = None;
2279 if missed > 0
2280 || events
2281 .iter()
2282 .any(|event| !matches!(event, ProviderEvent::SessionIdentityObserved { .. }))
2283 {
2284 let capability = ProviderRuntimeCapability::SemanticReadiness;
2289 if !state.runtime_policy.admits(capability) {
2290 denied_capability = Some(capability);
2291 }
2292 }
2293 if denied_capability.is_none()
2294 && events
2295 .iter()
2296 .any(provider_event_carries_session_identity)
2297 {
2298 let capability = ProviderRuntimeCapability::ProviderSessionIdentity;
2299 if !state.runtime_policy.admits(capability) {
2300 denied_capability = Some(capability);
2301 }
2302 }
2303 if let Some(capability) = denied_capability {
2304 let dropped = missed.saturating_add(events.len() as u64);
2327 if state.snapshot.provider.sequence.checked_add(1).is_none() {
2328 return Err(ControlError::ProviderSequenceExhausted {
2329 instance_id,
2330 generation,
2331 });
2332 }
2333 let canonical_sequence = {
2334 let snapshot = &mut self
2335 .sessions
2336 .get_mut(&instance_id)
2337 .expect("validated session")
2338 .snapshot
2339 .provider;
2340 reduce_provider_gap(snapshot, &source, source_sequence, dropped)
2341 };
2342 self.bump_revision();
2343 self.emit_event(
2344 Some(command_id),
2345 instance_id,
2346 generation,
2347 ControlEventKind::ProviderEvent {
2348 sequence: canonical_sequence,
2349 source: source.clone(),
2350 source_sequence,
2351 event: ProviderEvent::Error {
2352 message: provider_refusal_message(capability, dropped),
2353 },
2354 },
2355 );
2356 return Err(ControlError::ProviderRuntimePolicyDenied { capability });
2357 }
2358 let canonical_steps = events.len() as u64 + u64::from(missed > 0);
2359 if state
2360 .snapshot
2361 .provider
2362 .sequence
2363 .checked_add(canonical_steps)
2364 .is_none()
2365 {
2366 return Err(ControlError::ProviderSequenceExhausted {
2367 instance_id,
2368 generation,
2369 });
2370 }
2371 if missed > 0 {
2372 let canonical_sequence = {
2373 let snapshot = &mut self
2374 .sessions
2375 .get_mut(&instance_id)
2376 .expect("validated session")
2377 .snapshot
2378 .provider;
2379 reduce_provider_gap(snapshot, &source, source_sequence - 1, missed)
2380 };
2381 self.bump_revision();
2382 self.emit_event(
2383 Some(command_id),
2384 instance_id,
2385 generation,
2386 ControlEventKind::ProviderGap {
2387 sequence: canonical_sequence,
2388 source: source.clone(),
2389 source_sequence: source_sequence - 1,
2390 missed,
2391 },
2392 );
2393 }
2394
2395 for event in events {
2396 let (reduction, superseded_resolution) = {
2397 let state = self
2398 .sessions
2399 .get_mut(&instance_id)
2400 .expect("validated session");
2401 let pending_resolution = pending_interaction_resolution(&state.snapshot);
2402 let reduction = reduce_provider_event(
2403 &mut state.snapshot.provider,
2404 &source,
2405 source_sequence,
2406 event.clone(),
2407 );
2408 let superseded_resolution =
2409 pending_resolution.filter(|(operation_id, interaction_id)| {
2410 !interaction_resolution_is_pending(
2411 &state.snapshot,
2412 *operation_id,
2413 *interaction_id,
2414 )
2415 });
2416 if superseded_resolution.is_some() {
2417 state.snapshot.pending_operation = None;
2418 }
2419 (reduction, superseded_resolution)
2420 };
2421 if let Some((operation_id, _)) = superseded_resolution {
2422 self.effects.retain(|effect| {
2423 effect.instance_id != instance_id
2424 || effect.generation != generation
2425 || effect.operation_id != operation_id
2426 });
2427 }
2428 self.bump_revision();
2429 self.emit_event(
2430 Some(command_id),
2431 instance_id,
2432 generation,
2433 ControlEventKind::ProviderEvent {
2434 sequence: reduction.sequence,
2435 source: source.clone(),
2436 source_sequence,
2437 event,
2438 },
2439 );
2440 self.emit_interaction_transitions(
2441 Some(command_id),
2442 instance_id,
2443 generation,
2444 reduction.interaction_transitions,
2445 );
2446 }
2447 Ok(())
2448 }
2449
2450 fn resolve_interaction(
2451 &mut self,
2452 command_id: CommandId,
2453 instance_id: AgentInstanceId,
2454 generation: SessionGeneration,
2455 interaction_id: ProviderInteractionId,
2456 response: ProviderInteractionResponse,
2457 ) -> Result<(), ControlError> {
2458 let state = self
2459 .sessions
2460 .get(&instance_id)
2461 .ok_or(ControlError::UnknownInstance { instance_id })?;
2462 if state.snapshot.generation != generation {
2463 return Err(ControlError::StaleProviderInteractionGeneration {
2464 expected: state.snapshot.generation,
2465 actual: generation,
2466 });
2467 }
2468 if state.snapshot.status != SessionStatus::Running {
2469 return Err(ControlError::InvalidTransition {
2470 instance_id,
2471 action: "resolve provider interaction".to_owned(),
2472 status: state.snapshot.status.clone(),
2473 });
2474 }
2475 if let Some(operation_id) = state.snapshot.pending_operation {
2476 return Err(ControlError::OperationPending {
2477 instance_id,
2478 operation_id,
2479 });
2480 }
2481 let interaction = state
2482 .snapshot
2483 .provider
2484 .interactions
2485 .iter()
2486 .find(|interaction| interaction.id == interaction_id)
2487 .ok_or(ControlError::UnknownProviderInteraction { interaction_id })?;
2488 if interaction.status != ProviderInteractionStatus::Pending {
2489 return Err(ControlError::ProviderInteractionNotPending { interaction_id });
2490 }
2491 response
2492 .validate_for(interaction.interaction_kind)
2493 .map_err(|error| ControlError::InvalidProviderInteractionResponse {
2494 message: error.to_string(),
2495 })?;
2496 let response_kind = response.kind();
2497 let target = ProviderInteractionTarget {
2498 interaction_id,
2499 source: interaction.source.clone(),
2500 provider_request_id: interaction.provider_request_id.clone(),
2501 interaction_kind: interaction.interaction_kind,
2502 tool_name: interaction.tool_name.clone(),
2503 agent_id: interaction.agent_id.clone(),
2504 };
2505
2506 let operation_id = self.allocate_operation();
2507 let state = self
2508 .sessions
2509 .get_mut(&instance_id)
2510 .expect("validated session");
2511 state.snapshot.pending_operation = Some(operation_id);
2512 state
2513 .snapshot
2514 .provider
2515 .interactions
2516 .iter_mut()
2517 .find(|interaction| interaction.id == interaction_id)
2518 .expect("validated interaction")
2519 .status = ProviderInteractionStatus::Resolving {
2520 operation_id,
2521 response_kind,
2522 };
2523 self.effects.push(EffectEnvelope {
2524 operation_id,
2525 instance_id,
2526 generation,
2527 effect: ControlEffect::ResolveInteraction { target, response },
2528 });
2529 self.bump_revision();
2530 self.emit_event(
2531 Some(command_id),
2532 instance_id,
2533 generation,
2534 ControlEventKind::InteractionResolutionRequested {
2535 operation_id,
2536 interaction_id,
2537 response_kind,
2538 },
2539 );
2540 Ok(())
2541 }
2542
2543 fn set_session_mode(
2551 &mut self,
2552 command_id: CommandId,
2553 instance_id: AgentInstanceId,
2554 mode_id: String,
2555 ) -> Result<(), ControlError> {
2556 validate_session_control_id("session mode id", &mode_id).map_err(|error| {
2557 ControlError::InvalidSessionModeRequest {
2558 message: error.to_string(),
2559 }
2560 })?;
2561 let state = self
2562 .sessions
2563 .get(&instance_id)
2564 .ok_or(ControlError::UnknownInstance { instance_id })?;
2565 if state.snapshot.transport != TransportKind::Acp {
2566 return Err(ControlError::UnsupportedTransportOperation {
2567 transport: state.snapshot.transport,
2568 action: "session mode switch".to_owned(),
2569 });
2570 }
2571 if state.snapshot.status != SessionStatus::Running {
2572 return Err(ControlError::InvalidTransition {
2573 instance_id,
2574 action: "set session mode".to_owned(),
2575 status: state.snapshot.status.clone(),
2576 });
2577 }
2578 if let Some(operation_id) = state.snapshot.pending_operation {
2579 return Err(ControlError::OperationPending {
2580 instance_id,
2581 operation_id,
2582 });
2583 }
2584
2585 let operation_id = self.allocate_operation();
2586 let generation = {
2587 let state = self
2588 .sessions
2589 .get_mut(&instance_id)
2590 .expect("validated session");
2591 state.snapshot.pending_operation = Some(operation_id);
2592 state.snapshot.generation
2593 };
2594 self.effects.push(EffectEnvelope {
2595 operation_id,
2596 instance_id,
2597 generation,
2598 effect: ControlEffect::SetSessionMode {
2599 mode_id: mode_id.clone(),
2600 },
2601 });
2602 self.bump_revision();
2603 self.emit_event(
2604 Some(command_id),
2605 instance_id,
2606 generation,
2607 ControlEventKind::SessionModeSetRequested {
2608 operation_id,
2609 mode_id,
2610 },
2611 );
2612 Ok(())
2613 }
2614
2615 fn set_session_config_option(
2622 &mut self,
2623 command_id: CommandId,
2624 instance_id: AgentInstanceId,
2625 option_id: String,
2626 value_json: String,
2627 ) -> Result<(), ControlError> {
2628 validate_session_control_id("session config option id", &option_id).map_err(|error| {
2629 ControlError::InvalidSessionConfigOptionRequest {
2630 message: error.to_string(),
2631 }
2632 })?;
2633 validate_session_config_value_json(&value_json).map_err(|error| {
2634 ControlError::InvalidSessionConfigOptionRequest {
2635 message: error.to_string(),
2636 }
2637 })?;
2638 let state = self
2639 .sessions
2640 .get(&instance_id)
2641 .ok_or(ControlError::UnknownInstance { instance_id })?;
2642 if state.snapshot.transport != TransportKind::Acp {
2643 return Err(ControlError::UnsupportedTransportOperation {
2644 transport: state.snapshot.transport,
2645 action: "session config option switch".to_owned(),
2646 });
2647 }
2648 if state.snapshot.status != SessionStatus::Running {
2649 return Err(ControlError::InvalidTransition {
2650 instance_id,
2651 action: "set session config option".to_owned(),
2652 status: state.snapshot.status.clone(),
2653 });
2654 }
2655 if let Some(operation_id) = state.snapshot.pending_operation {
2656 return Err(ControlError::OperationPending {
2657 instance_id,
2658 operation_id,
2659 });
2660 }
2661
2662 let operation_id = self.allocate_operation();
2663 let generation = {
2664 let state = self
2665 .sessions
2666 .get_mut(&instance_id)
2667 .expect("validated session");
2668 state.snapshot.pending_operation = Some(operation_id);
2669 state.snapshot.generation
2670 };
2671 self.effects.push(EffectEnvelope {
2672 operation_id,
2673 instance_id,
2674 generation,
2675 effect: ControlEffect::SetSessionConfigOption {
2676 option_id: option_id.clone(),
2677 value_json,
2678 },
2679 });
2680 self.bump_revision();
2681 self.emit_event(
2682 Some(command_id),
2683 instance_id,
2684 generation,
2685 ControlEventKind::SessionConfigOptionSetRequested {
2686 operation_id,
2687 option_id,
2688 },
2689 );
2690 Ok(())
2691 }
2692
2693 fn set_session_model(
2699 &mut self,
2700 command_id: CommandId,
2701 instance_id: AgentInstanceId,
2702 model_id: String,
2703 ) -> Result<(), ControlError> {
2704 validate_session_control_id("session model id", &model_id).map_err(|error| {
2705 ControlError::InvalidSessionModelRequest {
2706 message: error.to_string(),
2707 }
2708 })?;
2709 let state = self
2710 .sessions
2711 .get(&instance_id)
2712 .ok_or(ControlError::UnknownInstance { instance_id })?;
2713 if state.snapshot.transport != TransportKind::Acp {
2714 return Err(ControlError::UnsupportedTransportOperation {
2715 transport: state.snapshot.transport,
2716 action: "session model switch".to_owned(),
2717 });
2718 }
2719 if state.snapshot.status != SessionStatus::Running {
2720 return Err(ControlError::InvalidTransition {
2721 instance_id,
2722 action: "set session model".to_owned(),
2723 status: state.snapshot.status.clone(),
2724 });
2725 }
2726 if let Some(operation_id) = state.snapshot.pending_operation {
2727 return Err(ControlError::OperationPending {
2728 instance_id,
2729 operation_id,
2730 });
2731 }
2732
2733 let operation_id = self.allocate_operation();
2734 let generation = {
2735 let state = self
2736 .sessions
2737 .get_mut(&instance_id)
2738 .expect("validated session");
2739 state.snapshot.pending_operation = Some(operation_id);
2740 state.snapshot.generation
2741 };
2742 self.effects.push(EffectEnvelope {
2743 operation_id,
2744 instance_id,
2745 generation,
2746 effect: ControlEffect::SetSessionModel {
2747 model_id: model_id.clone(),
2748 },
2749 });
2750 self.bump_revision();
2751 self.emit_event(
2752 Some(command_id),
2753 instance_id,
2754 generation,
2755 ControlEventKind::SessionModelSetRequested {
2756 operation_id,
2757 model_id,
2758 },
2759 );
2760 Ok(())
2761 }
2762
2763 fn remove(
2764 &mut self,
2765 command_id: CommandId,
2766 instance_id: AgentInstanceId,
2767 ) -> Result<(), ControlError> {
2768 let state = self
2769 .sessions
2770 .get(&instance_id)
2771 .ok_or(ControlError::UnknownInstance { instance_id })?;
2772 if !state.snapshot.status.allows_remove() {
2773 return Err(ControlError::InvalidTransition {
2774 instance_id,
2775 action: "remove".to_owned(),
2776 status: state.snapshot.status.clone(),
2777 });
2778 }
2779 if let Some(operation_id) = state.snapshot.pending_operation {
2780 return Err(ControlError::OperationPending {
2781 instance_id,
2782 operation_id,
2783 });
2784 }
2785 if let Some(pending) = &state.snapshot.capabilities.pending {
2786 return Err(ControlError::CapabilityProbeOperationPending {
2787 operation_id: pending.operation_id,
2788 });
2789 }
2790 if let Some(pending) = &state.snapshot.history.pending {
2791 return Err(ControlError::HistoryOperationPending {
2792 operation_id: pending.operation_id,
2793 });
2794 }
2795 let generation = state.snapshot.generation;
2796 self.sessions.remove(&instance_id);
2797 self.effects
2798 .retain(|effect| effect.instance_id != instance_id);
2799 self.bump_revision();
2800 self.emit_event(
2801 Some(command_id),
2802 instance_id,
2803 generation,
2804 ControlEventKind::Removed,
2805 );
2806 Ok(())
2807 }
2808
2809 fn generation_watermark(
2810 &self,
2811 instance_id: AgentInstanceId,
2812 current: SessionGeneration,
2813 ) -> SessionGeneration {
2814 self.generation_watermarks
2815 .get(&instance_id)
2816 .copied()
2817 .map_or(current, |watermark| watermark.max(current))
2818 }
2819
2820 fn has_counter_headroom(&self) -> bool {
2821 self.counter_error.is_none()
2822 && self.next_operation_id.is_some()
2823 && self
2824 .next_event_sequence
2825 .is_some_and(|sequence| sequence.checked_add(CONTROL_EVENT_HEADROOM - 1).is_some())
2826 && self
2827 .revision
2828 .checked_add(CONTROL_REVISION_HEADROOM)
2829 .is_some()
2830 }
2831
2832 fn observation_has_provider_headroom(&self, envelope: &ObservationEnvelope) -> bool {
2833 let Some(state) = self.sessions.get(&envelope.instance_id) else {
2834 return true;
2835 };
2836 if state.snapshot.provider.sequence == u64::MAX {
2837 return !matches!(
2838 envelope.observation,
2839 ControlObservation::ProviderEvent { .. }
2840 | ControlObservation::ProviderGap { .. }
2841 );
2842 }
2843 match &envelope.observation {
2844 ControlObservation::ProviderGap { source, missed, .. } => {
2845 provider_source_sequence(&state.snapshot.provider, source)
2846 .checked_add(*missed)
2847 .is_some()
2848 }
2849 _ => true,
2850 }
2851 }
2852
2853 fn retire_exhausted_counter(&mut self, error: &ControlError) {
2854 match error {
2855 ControlError::OperationIdExhausted => self.next_operation_id = None,
2856 ControlError::EventSequenceExhausted => self.next_event_sequence = None,
2857 ControlError::RevisionExhausted => self.revision = u64::MAX,
2858 _ => {}
2859 }
2860 }
2861
2862 fn purge_generation_bound_history(&mut self, instance_id: AgentInstanceId) {
2865 self.effects.retain(|effect| {
2866 effect.instance_id != instance_id
2867 || !matches!(
2868 effect.effect,
2869 ControlEffect::DiscoverHistory { .. } | ControlEffect::LoadHistory { .. }
2870 )
2871 });
2872 }
2873
2874 fn allocate_operation(&mut self) -> OperationId {
2875 if self.counter_error.is_some() {
2876 return OperationId(0);
2877 }
2878 let Some(next) = self.next_operation_id.take() else {
2879 self.counter_error = Some(ControlError::OperationIdExhausted);
2880 return OperationId(0);
2881 };
2882 self.next_operation_id = next.checked_add(1);
2883 OperationId(next)
2884 }
2885
2886 fn session_mut(&mut self, instance_id: AgentInstanceId) -> &mut SessionSnapshot {
2887 &mut self
2888 .sessions
2889 .get_mut(&instance_id)
2890 .expect("validated instance must remain registered")
2891 .snapshot
2892 }
2893
2894 fn bump_revision(&mut self) {
2895 if self.counter_error.is_some() {
2896 return;
2897 }
2898 let Some(revision) = self.revision.checked_add(1) else {
2899 self.counter_error = Some(ControlError::RevisionExhausted);
2900 return;
2901 };
2902 self.revision = revision;
2903 }
2904
2905 fn emit_ignored(
2906 &mut self,
2907 instance_id: AgentInstanceId,
2908 generation: SessionGeneration,
2909 reason: ObservationIgnoredReason,
2910 ) {
2911 self.emit_event(
2912 None,
2913 instance_id,
2914 generation,
2915 ControlEventKind::ObservationIgnored { reason },
2916 );
2917 }
2918
2919 fn emit_interaction_transitions(
2920 &mut self,
2921 command_id: Option<CommandId>,
2922 instance_id: AgentInstanceId,
2923 generation: SessionGeneration,
2924 transitions: Vec<ProviderInteractionTransition>,
2925 ) {
2926 for transition in transitions {
2927 let event = match transition {
2928 ProviderInteractionTransition::Requested(interaction) => {
2929 ControlEventKind::InteractionRequested { interaction }
2930 }
2931 ProviderInteractionTransition::Resolved {
2932 interaction_id,
2933 outcome,
2934 } => ControlEventKind::InteractionResolved {
2935 interaction_id,
2936 outcome,
2937 },
2938 };
2939 self.emit_event(command_id, instance_id, generation, event);
2940 }
2941 }
2942
2943 fn emit_event(
2944 &mut self,
2945 command_id: Option<CommandId>,
2946 instance_id: AgentInstanceId,
2947 generation: SessionGeneration,
2948 event: ControlEventKind,
2949 ) {
2950 if self.counter_error.is_some() {
2951 return;
2952 }
2953 let Some(sequence) = self.next_event_sequence.take() else {
2954 self.counter_error = Some(ControlError::EventSequenceExhausted);
2955 return;
2956 };
2957 self.next_event_sequence = sequence.checked_add(1);
2958 self.events.push(ControlEvent {
2959 sequence,
2960 command_id,
2961 instance_id,
2962 generation,
2963 event,
2964 });
2965 }
2966}
2967
2968fn checked_next_generation(current: SessionGeneration) -> Option<SessionGeneration> {
2969 current.0.checked_add(1).map(SessionGeneration)
2970}
2971
2972fn invalidate_foreground(snapshot: &mut SessionSnapshot, reason: String) {
2973 snapshot.foreground.authority = ForegroundAuthority::Stale;
2974 snapshot.foreground.stale_reason = Some(reason);
2975}
2976
2977fn is_history_observation(observation: &ControlObservation) -> bool {
2978 matches!(
2979 observation,
2980 ControlObservation::HistoryDiscovered { .. }
2981 | ControlObservation::HistoryLoaded { .. }
2982 | ControlObservation::HistoryFailed { .. }
2983 )
2984}
2985
2986fn is_capability_observation(observation: &ControlObservation) -> bool {
2987 matches!(
2988 observation,
2989 ControlObservation::CapabilitiesProbed { .. }
2990 | ControlObservation::CapabilityProbeFailed { .. }
2991 )
2992}
2993
2994fn is_resume_authority_observation(observation: &ControlObservation) -> bool {
2995 matches!(
2996 observation,
2997 ControlObservation::ResumeAuthorized { .. }
2998 | ControlObservation::ResumeDenied { .. }
2999 | ControlObservation::ResumeFailed { .. }
3000 )
3001}
3002
3003fn require_runtime_capability(
3004 runtime_policy: ProviderRuntimePolicy,
3005 capability: ProviderRuntimeCapability,
3006) -> Result<(), ControlError> {
3007 if runtime_policy.admits(capability) {
3008 Ok(())
3009 } else {
3010 Err(ControlError::ProviderRuntimePolicyDenied { capability })
3011 }
3012}
3013
3014fn require_structured_prompt_policy(
3015 runtime_policy: ProviderRuntimePolicy,
3016) -> Result<(), ControlError> {
3017 require_runtime_capability(
3018 runtime_policy,
3019 ProviderRuntimeCapability::SemanticReadiness,
3020 )?;
3021 require_runtime_capability(runtime_policy, ProviderRuntimeCapability::StructuredPrompt)
3022}
3023
3024fn denied_observation_capability(
3025 runtime_policy: ProviderRuntimePolicy,
3026 observation: &ControlObservation,
3027) -> Option<ProviderRuntimeCapability> {
3028 let capability = match observation {
3029 ControlObservation::ProviderGap { .. } => ProviderRuntimeCapability::SemanticReadiness,
3039 ControlObservation::ResumeAuthorized { .. }
3044 | ControlObservation::ResumeDenied { .. }
3045 | ControlObservation::ResumeFailed { .. } => return None,
3046 _ => return None,
3047 };
3048 (!runtime_policy.admits(capability)).then_some(capability)
3049}
3050
3051fn provider_event_denied_capability(
3060 runtime_policy: ProviderRuntimePolicy,
3061 event: &ProviderEvent,
3062) -> Option<ProviderRuntimeCapability> {
3063 let capability = if matches!(event, ProviderEvent::SessionIdentityObserved { .. }) {
3064 ProviderRuntimeCapability::ProviderSessionIdentity
3065 } else if !runtime_policy.admits(ProviderRuntimeCapability::SemanticReadiness) {
3066 ProviderRuntimeCapability::SemanticReadiness
3067 } else if provider_event_carries_session_identity(event)
3068 && !runtime_policy.admits(ProviderRuntimeCapability::ProviderSessionIdentity)
3069 {
3070 ProviderRuntimeCapability::ProviderSessionIdentity
3071 } else {
3072 return None;
3073 };
3074 (!runtime_policy.admits(capability)).then_some(capability)
3075}
3076
3077fn provider_refusal_message(capability: ProviderRuntimeCapability, dropped: u64) -> String {
3082 let noun = if dropped == 1 { "event" } else { "events" };
3083 format!("provider events rejected: capability {capability:?} not admitted ({dropped} {noun})")
3084}
3085
3086fn provider_event_carries_session_identity(event: &ProviderEvent) -> bool {
3087 matches!(
3088 event,
3089 ProviderEvent::SessionStarted { .. } | ProviderEvent::SessionIdentityObserved { .. }
3090 )
3091}
3092
3093fn is_interaction_resolution_observation(observation: &ControlObservation) -> bool {
3094 matches!(
3095 observation,
3096 ControlObservation::InteractionResolutionCompleted { .. }
3097 | ControlObservation::InteractionResolutionFailed { .. }
3098 )
3099}
3100
3101fn interaction_failure_is_valid(message: &str) -> bool {
3102 !message.trim().is_empty()
3103 && message.len() <= PROVIDER_INTERACTION_FAILURE_MAX_BYTES
3104 && !message
3105 .chars()
3106 .any(|character| character.is_control() && !matches!(character, '\n' | '\r' | '\t'))
3107}
3108
3109fn interaction_response_outcome(
3110 response_kind: ProviderInteractionResponseKind,
3111) -> ProviderInteractionOutcome {
3112 match response_kind {
3113 ProviderInteractionResponseKind::ApproveOnce => ProviderInteractionOutcome::Approved,
3114 ProviderInteractionResponseKind::Deny => ProviderInteractionOutcome::Denied,
3115 ProviderInteractionResponseKind::Answer => ProviderInteractionOutcome::Answered,
3116 }
3117}
3118
3119fn interaction_is_unresolved(interaction: &ProviderInteraction) -> bool {
3120 matches!(
3121 interaction.status,
3122 ProviderInteractionStatus::Pending | ProviderInteractionStatus::Resolving { .. }
3123 )
3124}
3125
3126fn pending_interaction_resolution(
3127 snapshot: &SessionSnapshot,
3128) -> Option<(OperationId, ProviderInteractionId)> {
3129 let operation_id = snapshot.pending_operation?;
3130 snapshot
3131 .provider
3132 .interactions
3133 .iter()
3134 .find_map(|interaction| {
3135 matches!(
3136 interaction.status,
3137 ProviderInteractionStatus::Resolving {
3138 operation_id: interaction_operation_id,
3139 ..
3140 } if interaction_operation_id == operation_id
3141 )
3142 .then_some((operation_id, interaction.id))
3143 })
3144}
3145
3146fn interaction_resolution_is_pending(
3147 snapshot: &SessionSnapshot,
3148 operation_id: OperationId,
3149 interaction_id: ProviderInteractionId,
3150) -> bool {
3151 snapshot.provider.interactions.iter().any(|interaction| {
3152 interaction.id == interaction_id
3153 && matches!(
3154 interaction.status,
3155 ProviderInteractionStatus::Resolving {
3156 operation_id: interaction_operation_id,
3157 ..
3158 } if interaction_operation_id == operation_id
3159 )
3160 })
3161}
3162
3163fn resume_identity_matches_target(
3164 snapshot: &SessionSnapshot,
3165 target: &ResumeTarget,
3166 identity: &ProviderSessionIdentity,
3167) -> bool {
3168 match target {
3169 ResumeTarget::CurrentProvider => snapshot.provider.session.as_ref() == Some(identity),
3170 ResumeTarget::ProviderSession { identity: expected } => expected == identity,
3171 ResumeTarget::HistoryCandidate { candidate_id } => {
3172 snapshot.history.loaded_candidate_id.as_deref() == Some(candidate_id.as_str())
3173 && snapshot
3174 .history
3175 .loaded
3176 .as_ref()
3177 .is_some_and(|session| session.session_id == identity.id)
3178 }
3179 }
3180}
3181
3182fn history_candidates_are_valid(
3183 candidates: &[gate4agent_types::HistoryCandidateSummary],
3184 limit: u16,
3185) -> bool {
3186 if candidates.len() > usize::from(limit)
3187 || candidates
3188 .iter()
3189 .any(|candidate| candidate.validate().is_err())
3190 {
3191 return false;
3192 }
3193 let unique = candidates
3194 .iter()
3195 .map(|candidate| candidate.id.as_str())
3196 .collect::<BTreeSet<_>>();
3197 unique.len() == candidates.len()
3198}
3199
3200#[derive(Clone, Debug, Eq, PartialEq)]
3201enum ProviderInteractionTransition {
3202 Requested(ProviderInteraction),
3203 Resolved {
3204 interaction_id: ProviderInteractionId,
3205 outcome: ProviderInteractionOutcome,
3206 },
3207}
3208
3209#[derive(Clone, Debug, Eq, PartialEq)]
3210struct ProviderReduction {
3211 sequence: u64,
3212 interaction_transitions: Vec<ProviderInteractionTransition>,
3213}
3214
3215fn provider_source_sequence(snapshot: &ProviderSnapshot, source: &ProviderSource) -> u64 {
3216 snapshot
3217 .sources
3218 .iter()
3219 .find(|cursor| cursor.source == *source)
3220 .map_or(0, |cursor| cursor.sequence)
3221}
3222
3223fn provider_source_cursor_mut<'a>(
3224 snapshot: &'a mut ProviderSnapshot,
3225 source: &ProviderSource,
3226) -> &'a mut ProviderSourceCursor {
3227 if let Some(index) = snapshot
3228 .sources
3229 .iter()
3230 .position(|cursor| cursor.source == *source)
3231 {
3232 return &mut snapshot.sources[index];
3233 }
3234 snapshot.sources.push(ProviderSourceCursor {
3235 source: source.clone(),
3236 sequence: 0,
3237 gap_count: 0,
3238 stale: false,
3239 });
3240 snapshot
3241 .sources
3242 .last_mut()
3243 .expect("provider source was just inserted")
3244}
3245
3246fn reduce_provider_event(
3247 snapshot: &mut ProviderSnapshot,
3248 source: &ProviderSource,
3249 source_sequence: u64,
3250 event: ProviderEvent,
3251) -> ProviderReduction {
3252 let canonical_sequence = snapshot
3253 .sequence
3254 .checked_add(1)
3255 .expect("provider sequence capacity must be preflighted");
3256 let mut interaction_transitions = Vec::new();
3257 match &event {
3258 ProviderEvent::SessionStarted {
3259 session_id,
3260 model,
3261 tools,
3262 } => {
3263 snapshot.session = Some(ProviderSessionIdentity {
3264 key: ProviderSessionKey::SessionId,
3265 id: session_id.clone(),
3266 transcript_path: None,
3267 });
3268 snapshot.model = (!model.is_empty()).then(|| model.clone());
3269 snapshot.tools = tools.clone();
3270 snapshot.lead_activity = ProviderActivity::Idle;
3271 snapshot.current_prompt = None;
3272 snapshot.active_tools.clear();
3273 remove_source_subagents(snapshot, source);
3274 interaction_transitions.extend(resolve_source_pending_interactions(
3275 snapshot,
3276 source,
3277 ProviderInteractionOutcome::TurnEnded,
3278 ));
3279 }
3280 ProviderEvent::SessionIdentityObserved { identity } => {
3281 snapshot.session = Some(identity.clone());
3282 }
3283 ProviderEvent::TurnStarted { prompt } => {
3284 interaction_transitions.extend(resolve_source_pending_interactions(
3285 snapshot,
3286 source,
3287 ProviderInteractionOutcome::Superseded,
3288 ));
3289 snapshot.lead_activity = ProviderActivity::Working;
3290 snapshot.current_prompt = prompt.clone();
3291 snapshot.active_tools.clear();
3292 }
3293 ProviderEvent::WorkingObserved => {
3294 interaction_transitions.extend(resolve_source_pending_interactions_by_kind(
3295 snapshot, source,
3296 ));
3297 snapshot.lead_activity = ProviderActivity::Working;
3298 }
3299 ProviderEvent::ToolStarted {
3300 id,
3301 name,
3302 input_json,
3303 agent_id,
3304 } => {
3305 let resume_activity =
3306 matching_interaction_resume_activity(snapshot, source, id, agent_id.as_deref());
3307 let resolved =
3308 resolve_matching_provider_interactions(snapshot, source, id, agent_id.as_deref());
3309 if !resolved.is_empty() {
3310 snapshot.lead_activity = resume_activity.unwrap_or(ProviderActivity::Working);
3311 } else if agent_id.is_none() {
3312 snapshot.lead_activity = ProviderActivity::Working;
3313 }
3314 interaction_transitions.extend(resolved);
3315 if agent_id.is_none() {
3316 if let Some(active) = snapshot.active_tools.iter_mut().find(|tool| tool.id == *id) {
3317 active.name = name.clone();
3318 active.input_json = input_json.clone();
3319 } else {
3320 snapshot.active_tools.push(ActiveProviderTool {
3321 id: id.clone(),
3322 name: name.clone(),
3323 input_json: input_json.clone(),
3324 });
3325 }
3326 }
3327 }
3328 ProviderEvent::ToolCompleted { id, agent_id, .. } => {
3329 let resume_activity =
3330 matching_interaction_resume_activity(snapshot, source, id, agent_id.as_deref());
3331 let resolved =
3332 resolve_matching_provider_interactions(snapshot, source, id, agent_id.as_deref());
3333 if !resolved.is_empty() {
3334 snapshot.lead_activity = resume_activity.unwrap_or(ProviderActivity::Working);
3335 }
3336 interaction_transitions.extend(resolved);
3337 if agent_id.is_none() {
3338 snapshot.active_tools.retain(|tool| tool.id != *id);
3339 }
3340 }
3341 ProviderEvent::TurnCompleted {
3342 usage,
3343 is_cumulative,
3344 } => {
3345 snapshot.completed_turns = snapshot.completed_turns.saturating_add(1);
3346 if *is_cumulative {
3347 snapshot.usage = usage.clone();
3348 } else {
3349 add_token_usage(&mut snapshot.usage, usage);
3350 }
3351 snapshot.lead_activity = ProviderActivity::Idle;
3352 snapshot.current_prompt = None;
3353 snapshot.active_tools.clear();
3354 interaction_transitions.extend(resolve_source_pending_interactions(
3355 snapshot,
3356 source,
3357 ProviderInteractionOutcome::TurnEnded,
3358 ));
3359 }
3360 ProviderEvent::TurnInterrupted => {
3361 snapshot.lead_activity = ProviderActivity::Idle;
3362 snapshot.current_prompt = None;
3363 snapshot.active_tools.clear();
3364 interaction_transitions.extend(resolve_source_pending_interactions(
3365 snapshot,
3366 source,
3367 ProviderInteractionOutcome::Interrupted,
3368 ));
3369 }
3370 ProviderEvent::InteractionRequested {
3371 request_id,
3372 interaction_kind,
3373 tool_name,
3374 prompt,
3375 agent_id,
3376 ..
3377 } => {
3378 let inherited_resume_activity = request_id.as_deref().and_then(|request_id| {
3379 matching_interaction_resume_activity(
3380 snapshot,
3381 source,
3382 request_id,
3383 agent_id.as_deref(),
3384 )
3385 });
3386 if let Some(request_id) = request_id {
3387 interaction_transitions.extend(resolve_provider_request_interactions(
3388 snapshot,
3389 source,
3390 request_id,
3391 agent_id.as_deref(),
3392 ProviderInteractionOutcome::Superseded,
3393 ));
3394 }
3395 let interaction = ProviderInteraction {
3396 id: ProviderInteractionId(canonical_sequence),
3397 source: source.clone(),
3398 provider_request_id: request_id.clone(),
3399 interaction_kind: *interaction_kind,
3400 tool_name: tool_name.clone(),
3401 prompt: prompt.clone(),
3402 agent_id: agent_id.clone(),
3403 resume_lead_activity: agent_id
3404 .is_some()
3405 .then_some(inherited_resume_activity.unwrap_or(snapshot.lead_activity)),
3406 status: ProviderInteractionStatus::Pending,
3407 };
3408 push_provider_interaction(snapshot, interaction.clone(), &mut interaction_transitions);
3409 interaction_transitions.push(ProviderInteractionTransition::Requested(interaction));
3410 snapshot.lead_activity = ProviderActivity::WaitingForInput;
3411 }
3412 ProviderEvent::InteractionResolved {
3413 request_id,
3414 outcome,
3415 } => {
3416 let resolved = resolve_provider_reported_interactions(
3417 snapshot,
3418 source,
3419 request_id,
3420 *outcome,
3421 );
3422 if !resolved.is_empty() {
3423 snapshot.lead_activity = if snapshot
3424 .interactions
3425 .iter()
3426 .any(interaction_is_unresolved)
3427 {
3428 ProviderActivity::WaitingForInput
3429 } else {
3430 ProviderActivity::Working
3431 };
3432 }
3433 interaction_transitions.extend(resolved);
3434 }
3435 ProviderEvent::RateLimited { .. } | ProviderEvent::Error { .. } => {
3436 snapshot.lead_activity = ProviderActivity::Blocked;
3437 }
3438 ProviderEvent::Ready => {
3439 snapshot.lead_activity = if snapshot.interactions.iter().any(interaction_is_unresolved)
3440 {
3441 ProviderActivity::WaitingForInput
3442 } else {
3443 ProviderActivity::Idle
3444 };
3445 }
3446 ProviderEvent::SessionEnded { .. } => {
3447 snapshot.lead_activity = ProviderActivity::Idle;
3448 snapshot.current_prompt = None;
3449 snapshot.active_tools.clear();
3450 remove_source_subagents(snapshot, source);
3451 interaction_transitions.extend(resolve_source_pending_interactions(
3452 snapshot,
3453 source,
3454 ProviderInteractionOutcome::TurnEnded,
3455 ));
3456 }
3457 ProviderEvent::SubagentStarted {
3458 agent_id,
3459 agent_type,
3460 description,
3461 } => {
3462 if let Some(existing) = snapshot.subagents.iter_mut().find(|subagent| {
3463 subagent.source == *source && subagent.provider_agent_id == *agent_id
3464 }) {
3465 existing.agent_type = agent_type.clone().or(existing.agent_type.take());
3466 existing.description = description.clone().or(existing.description.take());
3467 } else if snapshot.subagents.len() < PROVIDER_SUBAGENTS_MAX {
3468 snapshot.subagents.push(ProviderSubagent {
3469 source: source.clone(),
3470 provider_agent_id: agent_id.clone(),
3471 agent_type: agent_type.clone(),
3472 description: description.clone(),
3473 });
3474 }
3475 }
3476 ProviderEvent::SubagentStopped { agent_id } => {
3477 let resume_activity = subagent_interaction_resume_activity(snapshot, source, agent_id);
3478 interaction_transitions.extend(resolve_subagent_pending_interactions(
3479 snapshot,
3480 source,
3481 agent_id,
3482 ProviderInteractionOutcome::TurnEnded,
3483 ));
3484 if let Some(resume_activity) = resume_activity {
3485 snapshot.lead_activity = resume_activity;
3486 }
3487 snapshot.subagents.retain(|subagent| {
3488 subagent.source != *source || subagent.provider_agent_id != *agent_id
3489 });
3490 }
3491 ProviderEvent::Text { .. }
3499 | ProviderEvent::Thinking { .. }
3500 | ProviderEvent::ContextWindowUsage { .. }
3501 | ProviderEvent::HostRequestObserved { .. }
3502 | ProviderEvent::UnrecognizedNotification { .. }
3503 | ProviderEvent::UserMessage { .. }
3504 | ProviderEvent::Plan { .. }
3505 | ProviderEvent::AvailableCommandsUpdated { .. }
3506 | ProviderEvent::ModeChanged { .. }
3507 | ProviderEvent::SessionInfoUpdated { .. }
3508 | ProviderEvent::UsageUpdated { .. }
3509 | ProviderEvent::ConfigOptionsUpdated { .. } => {}
3510 }
3511 refresh_provider_activity(snapshot);
3512 provider_source_cursor_mut(snapshot, source).sequence = source_sequence;
3513 provider_source_cursor_mut(snapshot, source).stale = false;
3514 snapshot.sequence = canonical_sequence;
3515 snapshot.last_event = Some(event);
3516 snapshot.stale = snapshot.sources.iter().any(|cursor| cursor.stale);
3517 ProviderReduction {
3518 sequence: canonical_sequence,
3519 interaction_transitions,
3520 }
3521}
3522
3523fn refresh_provider_activity(snapshot: &mut ProviderSnapshot) {
3524 snapshot.activity =
3525 if snapshot.lead_activity == ProviderActivity::Idle && !snapshot.subagents.is_empty() {
3526 ProviderActivity::Working
3527 } else {
3528 snapshot.lead_activity
3529 };
3530}
3531
3532fn remove_source_subagents(snapshot: &mut ProviderSnapshot, source: &ProviderSource) {
3533 snapshot
3534 .subagents
3535 .retain(|subagent| subagent.source != *source);
3536}
3537
3538fn push_provider_interaction(
3539 snapshot: &mut ProviderSnapshot,
3540 interaction: ProviderInteraction,
3541 transitions: &mut Vec<ProviderInteractionTransition>,
3542) {
3543 if snapshot.interactions.len() >= PROVIDER_INTERACTIONS_MAX {
3544 let remove_index = snapshot
3545 .interactions
3546 .iter()
3547 .position(|existing| {
3548 matches!(existing.status, ProviderInteractionStatus::Resolved { .. })
3549 })
3550 .unwrap_or(0);
3551 let removed = snapshot.interactions.remove(remove_index);
3552 if interaction_is_unresolved(&removed) {
3553 transitions.push(ProviderInteractionTransition::Resolved {
3554 interaction_id: removed.id,
3555 outcome: ProviderInteractionOutcome::Superseded,
3556 });
3557 }
3558 }
3559 snapshot.interactions.push(interaction);
3560}
3561
3562fn resolve_matching_provider_interactions(
3563 snapshot: &mut ProviderSnapshot,
3564 source: &ProviderSource,
3565 provider_request_id: &str,
3566 agent_id: Option<&str>,
3567) -> Vec<ProviderInteractionTransition> {
3568 let matching: Vec<_> = snapshot
3569 .interactions
3570 .iter()
3571 .filter_map(|interaction| {
3572 (interaction.source == *source
3573 && interaction.provider_request_id.as_deref() == Some(provider_request_id)
3574 && interaction.agent_id.as_deref() == agent_id)
3575 .then(|| {
3576 provider_progress_outcome(interaction).map(|outcome| (interaction.id, outcome))
3577 })
3578 .flatten()
3579 })
3580 .collect();
3581 resolve_interaction_ids(snapshot, matching)
3582}
3583
3584fn matching_interaction_resume_activity(
3585 snapshot: &ProviderSnapshot,
3586 source: &ProviderSource,
3587 provider_request_id: &str,
3588 agent_id: Option<&str>,
3589) -> Option<ProviderActivity> {
3590 snapshot.interactions.iter().rev().find_map(|interaction| {
3591 (interaction.source == *source
3592 && interaction.provider_request_id.as_deref() == Some(provider_request_id)
3593 && interaction.agent_id.as_deref() == agent_id
3594 && interaction_is_unresolved(interaction))
3595 .then_some(interaction.resume_lead_activity)
3596 .flatten()
3597 })
3598}
3599
3600fn resolve_provider_request_interactions(
3601 snapshot: &mut ProviderSnapshot,
3602 source: &ProviderSource,
3603 provider_request_id: &str,
3604 agent_id: Option<&str>,
3605 outcome: ProviderInteractionOutcome,
3606) -> Vec<ProviderInteractionTransition> {
3607 let matching = snapshot
3608 .interactions
3609 .iter()
3610 .filter(|interaction| {
3611 interaction.source == *source
3612 && interaction.provider_request_id.as_deref() == Some(provider_request_id)
3613 && interaction.agent_id.as_deref() == agent_id
3614 && interaction_is_unresolved(interaction)
3615 })
3616 .map(|interaction| (interaction.id, outcome))
3617 .collect();
3618 resolve_interaction_ids(snapshot, matching)
3619}
3620
3621fn resolve_provider_reported_interactions(
3622 snapshot: &mut ProviderSnapshot,
3623 source: &ProviderSource,
3624 provider_request_id: &str,
3625 outcome: ProviderInteractionOutcome,
3626) -> Vec<ProviderInteractionTransition> {
3627 let matching = snapshot
3628 .interactions
3629 .iter()
3630 .filter(|interaction| {
3631 interaction.source == *source
3632 && interaction.provider_request_id.as_deref() == Some(provider_request_id)
3633 && interaction_is_unresolved(interaction)
3634 })
3635 .map(|interaction| (interaction.id, outcome))
3636 .collect();
3637 resolve_interaction_ids(snapshot, matching)
3638}
3639
3640fn resolve_source_pending_interactions(
3641 snapshot: &mut ProviderSnapshot,
3642 source: &ProviderSource,
3643 outcome: ProviderInteractionOutcome,
3644) -> Vec<ProviderInteractionTransition> {
3645 let matching = snapshot
3646 .interactions
3647 .iter()
3648 .filter(|interaction| {
3649 interaction.source == *source && interaction_is_unresolved(interaction)
3650 })
3651 .map(|interaction| (interaction.id, outcome))
3652 .collect();
3653 resolve_interaction_ids(snapshot, matching)
3654}
3655
3656fn resolve_source_pending_interactions_by_kind(
3657 snapshot: &mut ProviderSnapshot,
3658 source: &ProviderSource,
3659) -> Vec<ProviderInteractionTransition> {
3660 let matching = snapshot
3661 .interactions
3662 .iter()
3663 .filter_map(|interaction| {
3664 (interaction.source == *source)
3665 .then(|| {
3666 provider_progress_outcome(interaction).map(|outcome| (interaction.id, outcome))
3667 })
3668 .flatten()
3669 })
3670 .collect::<Vec<_>>();
3671 resolve_interaction_ids(snapshot, matching)
3672}
3673
3674fn provider_progress_outcome(
3675 interaction: &ProviderInteraction,
3676) -> Option<ProviderInteractionOutcome> {
3677 match interaction.status {
3678 ProviderInteractionStatus::Pending => Some(match interaction.interaction_kind {
3679 ProviderInteractionKind::Approval => ProviderInteractionOutcome::Approved,
3680 ProviderInteractionKind::Question => ProviderInteractionOutcome::Answered,
3681 }),
3682 ProviderInteractionStatus::Resolving { response_kind, .. } => {
3683 Some(interaction_response_outcome(response_kind))
3684 }
3685 ProviderInteractionStatus::Resolved { .. } => None,
3686 }
3687}
3688
3689fn resolve_subagent_pending_interactions(
3690 snapshot: &mut ProviderSnapshot,
3691 source: &ProviderSource,
3692 agent_id: &str,
3693 outcome: ProviderInteractionOutcome,
3694) -> Vec<ProviderInteractionTransition> {
3695 let matching = snapshot
3696 .interactions
3697 .iter()
3698 .filter(|interaction| {
3699 interaction.source == *source
3700 && interaction.agent_id.as_deref() == Some(agent_id)
3701 && interaction_is_unresolved(interaction)
3702 })
3703 .map(|interaction| (interaction.id, outcome))
3704 .collect();
3705 resolve_interaction_ids(snapshot, matching)
3706}
3707
3708fn subagent_interaction_resume_activity(
3709 snapshot: &ProviderSnapshot,
3710 source: &ProviderSource,
3711 agent_id: &str,
3712) -> Option<ProviderActivity> {
3713 snapshot.interactions.iter().find_map(|interaction| {
3714 (interaction.source == *source
3715 && interaction.agent_id.as_deref() == Some(agent_id)
3716 && interaction_is_unresolved(interaction))
3717 .then_some(interaction.resume_lead_activity)
3718 .flatten()
3719 })
3720}
3721
3722fn resolve_all_pending_interactions(
3723 snapshot: &mut ProviderSnapshot,
3724 outcome: ProviderInteractionOutcome,
3725) -> Vec<ProviderInteractionTransition> {
3726 let matching = snapshot
3727 .interactions
3728 .iter()
3729 .filter(|interaction| interaction_is_unresolved(interaction))
3730 .map(|interaction| (interaction.id, outcome))
3731 .collect();
3732 resolve_interaction_ids(snapshot, matching)
3733}
3734
3735fn resolve_interaction_ids(
3736 snapshot: &mut ProviderSnapshot,
3737 matching: Vec<(ProviderInteractionId, ProviderInteractionOutcome)>,
3738) -> Vec<ProviderInteractionTransition> {
3739 let mut transitions = Vec::with_capacity(matching.len());
3740 for (interaction_id, outcome) in matching {
3741 if let Some(interaction) = snapshot
3742 .interactions
3743 .iter_mut()
3744 .find(|interaction| interaction.id == interaction_id)
3745 {
3746 interaction.status = ProviderInteractionStatus::Resolved { outcome };
3747 transitions.push(ProviderInteractionTransition::Resolved {
3748 interaction_id,
3749 outcome,
3750 });
3751 }
3752 }
3753 transitions
3754}
3755
3756fn reduce_provider_gap(
3757 snapshot: &mut ProviderSnapshot,
3758 source: &ProviderSource,
3759 source_sequence: u64,
3760 missed: u64,
3761) -> u64 {
3762 let cursor = provider_source_cursor_mut(snapshot, source);
3763 cursor.sequence = source_sequence;
3764 cursor.gap_count = cursor.gap_count.saturating_add(missed);
3765 cursor.stale = true;
3766 snapshot.gap_count = snapshot.gap_count.saturating_add(missed);
3767 snapshot.stale = true;
3768 snapshot.sequence = snapshot
3769 .sequence
3770 .checked_add(1)
3771 .expect("provider sequence capacity must be preflighted");
3772 snapshot.sequence
3773}
3774
3775fn add_token_usage(total: &mut TokenUsage, delta: &TokenUsage) {
3776 total.input_tokens = total.input_tokens.saturating_add(delta.input_tokens);
3777 total.output_tokens = total.output_tokens.saturating_add(delta.output_tokens);
3778 total.cache_read_tokens = total
3779 .cache_read_tokens
3780 .saturating_add(delta.cache_read_tokens);
3781 total.cache_write_tokens = total
3782 .cache_write_tokens
3783 .saturating_add(delta.cache_write_tokens);
3784 total.reasoning_tokens = total
3785 .reasoning_tokens
3786 .saturating_add(delta.reasoning_tokens);
3787 if delta.context_window.is_some() {
3788 total.context_window = delta.context_window;
3789 }
3790}
3791
3792impl Default for Gate4AgentEngine {
3793 fn default() -> Self {
3794 Self::new()
3795 }
3796}
3797
3798#[cfg(test)]
3799mod tests {
3800 use super::*;
3801 use gate4agent_types::{
3802 AdapterBinding, AdapterFamily, AdapterId, AdapterVerification, AgentId, ApprovalLevel,
3803 CapabilityModelSummary, ForegroundAuthority, ForegroundProcess, ForegroundProcessKind,
3804 HistoryCandidateSummary, HistoryMessageRecord, HistoryMessageRole, HistoryQuery,
3805 HistorySessionRecord, InputAction, PreparedInputKind, PromptFraming, PromptPayload,
3806 SessionOptionSelection, ShellCommand, TerminalFrame, TerminalText, TransportKind,
3807 };
3808
3809 fn instance() -> AgentInstanceId {
3810 AgentInstanceId(7)
3811 }
3812
3813 fn provider_source() -> ProviderSource {
3814 ProviderSource {
3815 family: AdapterFamily::PtySemantic,
3816 binding: AdapterBinding::new(
3817 AdapterId::new("codex").unwrap(),
3818 "test/v1",
3819 AdapterVerification::SyntheticFixture,
3820 )
3821 .unwrap(),
3822 }
3823 }
3824
3825 fn hook_source() -> ProviderSource {
3826 ProviderSource {
3827 family: AdapterFamily::Hook,
3828 binding: AdapterBinding::new(
3829 AdapterId::new("grok").unwrap(),
3830 "test/v1",
3831 AdapterVerification::SyntheticFixture,
3832 )
3833 .unwrap(),
3834 }
3835 }
3836
3837 fn register_instance(command_id: u64, instance_id: AgentInstanceId) -> CommandEnvelope {
3838 CommandEnvelope {
3839 id: CommandId(command_id),
3840 command: ControlCommand::Register {
3841 instance_id,
3842 agent_id: AgentId::new("claude").unwrap(),
3843 transport: TransportKind::Pty,
3844 },
3845 }
3846 }
3847
3848 fn register(command_id: u64) -> CommandEnvelope {
3849 register_instance(command_id, instance())
3850 }
3851
3852 fn register_acp(command_id: u64) -> CommandEnvelope {
3861 CommandEnvelope {
3862 id: CommandId(command_id),
3863 command: ControlCommand::Register {
3864 instance_id: instance(),
3865 agent_id: AgentId::new("claude").unwrap(),
3866 transport: TransportKind::Acp,
3867 },
3868 }
3869 }
3870
3871 fn verified_runtime_policy() -> ProviderRuntimePolicy {
3872 ProviderRuntimePolicy::new(true, true, true, true, true, true).unwrap()
3873 }
3874
3875 fn start(command_id: u64) -> CommandEnvelope {
3876 start_with_policy(command_id, verified_runtime_policy())
3877 }
3878
3879 fn start_with_policy(
3880 command_id: u64,
3881 runtime_policy: ProviderRuntimePolicy,
3882 ) -> CommandEnvelope {
3883 CommandEnvelope {
3884 id: CommandId(command_id),
3885 command: ControlCommand::Start {
3886 instance_id: instance(),
3887 runtime_policy,
3888 request: StartRequest {
3889 working_directory: ".".to_owned(),
3890 terminal_size: TerminalSize {
3891 rows: 24,
3892 columns: 80,
3893 },
3894 initial_prompt: None,
3895 session_options: None,
3896 approval_level: ApprovalLevel::default(),
3897 },
3898 },
3899 }
3900 }
3901
3902 fn terminal_text(command_id: u64, text: &str) -> CommandEnvelope {
3903 CommandEnvelope {
3904 id: CommandId(command_id),
3905 command: ControlCommand::SendInput {
3906 instance_id: instance(),
3907 action: InputAction::TerminalText(TerminalText {
3908 text: text.to_owned(),
3909 }),
3910 },
3911 }
3912 }
3913
3914 fn running_engine() -> (Gate4AgentEngine, EffectEnvelope) {
3915 running_engine_with_policy(verified_runtime_policy())
3916 }
3917
3918 fn running_engine_with_policy(
3919 runtime_policy: ProviderRuntimePolicy,
3920 ) -> (Gate4AgentEngine, EffectEnvelope) {
3921 let mut engine = Gate4AgentEngine::new();
3922 engine.apply_command(register(1)).unwrap();
3923 engine.drain_events();
3924 engine.apply_command(start_with_policy(2, runtime_policy)).unwrap();
3925 let effect = engine.drain_effects().pop().unwrap();
3926 engine.apply_observation(ObservationEnvelope {
3927 operation_id: Some(effect.operation_id),
3928 instance_id: effect.instance_id,
3929 generation: effect.generation,
3930 observation: ControlObservation::Spawned {
3931 process_id: Some(42),
3932 },
3933 });
3934 (engine, effect)
3935 }
3936
3937 fn running_acp_engine_with_policy(
3943 runtime_policy: ProviderRuntimePolicy,
3944 ) -> (Gate4AgentEngine, EffectEnvelope) {
3945 let mut engine = Gate4AgentEngine::new();
3946 engine.apply_command(register_acp(1)).unwrap();
3947 engine.drain_events();
3948 engine.apply_command(start_with_policy(2, runtime_policy)).unwrap();
3949 let effect = engine.drain_effects().pop().unwrap();
3950 engine.apply_observation(ObservationEnvelope {
3951 operation_id: Some(effect.operation_id),
3952 instance_id: effect.instance_id,
3953 generation: effect.generation,
3954 observation: ControlObservation::Spawned {
3955 process_id: Some(42),
3956 },
3957 });
3958 (engine, effect)
3959 }
3960
3961 #[test]
3962 fn runtime_policy_keeps_raw_input_and_rejects_structured_prompt() {
3963 let (mut engine, spawn) = running_engine_with_policy(ProviderRuntimePolicy::raw_pty());
3964 assert!(matches!(
3965 &spawn.effect,
3966 ControlEffect::Spawn { runtime_policy, .. }
3967 if *runtime_policy == ProviderRuntimePolicy::raw_pty()
3968 ));
3969 engine
3970 .apply_command(CommandEnvelope {
3971 id: CommandId(3),
3972 command: ControlCommand::SendInput {
3973 instance_id: instance(),
3974 action: InputAction::TerminalBytes(vec![0x1b, b'[', b'A']),
3975 },
3976 })
3977 .unwrap();
3978 let raw_input = engine.drain_effects().pop().unwrap();
3979 assert!(matches!(raw_input.effect, ControlEffect::WriteInput { .. }));
3980 engine.apply_observation(ObservationEnvelope {
3981 operation_id: Some(raw_input.operation_id),
3982 instance_id: instance(),
3983 generation: spawn.generation,
3984 observation: ControlObservation::InputCompleted,
3985 });
3986
3987 assert_eq!(
3988 engine.apply_command(CommandEnvelope {
3989 id: CommandId(4),
3990 command: ControlCommand::SendInput {
3991 instance_id: instance(),
3992 action: InputAction::SubmitPrompt(PromptPayload {
3993 text: "semantic prompt".to_owned(),
3994 framing: PromptFraming::Literal,
3995 }),
3996 },
3997 }),
3998 Err(ControlError::ProviderRuntimePolicyDenied {
3999 capability: ProviderRuntimeCapability::SemanticReadiness,
4000 })
4001 );
4002 assert!(engine.drain_effects().is_empty());
4003 }
4004
4005 #[test]
4019 fn runtime_policy_fail_closed_provider_observations_and_ingress() {
4020 let (mut engine, spawn) = running_engine_with_policy(ProviderRuntimePolicy::raw_pty());
4021 let mut sequence = 0u64;
4022 for (event, expected_capability) in [
4023 (
4024 ProviderEvent::WorkingObserved,
4025 ProviderRuntimeCapability::SemanticReadiness,
4026 ),
4027 (
4028 ProviderEvent::SessionIdentityObserved {
4029 identity: ProviderSessionIdentity {
4030 key: ProviderSessionKey::SessionId,
4031 id: "provider-session".to_owned(),
4032 transcript_path: None,
4033 },
4034 },
4035 ProviderRuntimeCapability::ProviderSessionIdentity,
4036 ),
4037 ] {
4038 sequence += 1;
4039 engine.apply_observation(ObservationEnvelope {
4040 operation_id: None,
4041 instance_id: instance(),
4042 generation: spawn.generation,
4043 observation: ControlObservation::ProviderEvent {
4044 source: provider_source(),
4045 sequence,
4046 event,
4047 },
4048 });
4049 assert!(engine.drain_events().iter().any(|event| matches!(
4050 &event.event,
4051 ControlEventKind::ProviderEvent {
4052 event: ProviderEvent::Error { message },
4053 ..
4054 } if message.contains(&format!("{expected_capability:?}"))
4055 )));
4056 }
4057 assert_eq!(engine.snapshot().sessions[0].provider.sequence, sequence);
4058
4059 sequence += 1;
4060 assert_eq!(
4061 engine.apply_command(CommandEnvelope {
4062 id: CommandId(5),
4063 command: ControlCommand::IngestProvider {
4064 instance_id: instance(),
4065 generation: spawn.generation,
4066 source: provider_source(),
4067 source_sequence: sequence,
4068 events: vec![ProviderEvent::WorkingObserved],
4069 },
4070 }),
4071 Err(ControlError::ProviderRuntimePolicyDenied {
4072 capability: ProviderRuntimeCapability::SemanticReadiness,
4073 })
4074 );
4075 sequence += 1;
4076 assert_eq!(
4077 engine.apply_command(CommandEnvelope {
4078 id: CommandId(6),
4079 command: ControlCommand::IngestProvider {
4080 instance_id: instance(),
4081 generation: spawn.generation,
4082 source: provider_source(),
4083 source_sequence: sequence,
4084 events: vec![ProviderEvent::SessionIdentityObserved {
4085 identity: ProviderSessionIdentity {
4086 key: ProviderSessionKey::SessionId,
4087 id: "provider-session".to_owned(),
4088 transcript_path: None,
4089 },
4090 }],
4091 },
4092 }),
4093 Err(ControlError::ProviderRuntimePolicyDenied {
4094 capability: ProviderRuntimeCapability::ProviderSessionIdentity,
4095 })
4096 );
4097 assert_eq!(engine.snapshot().sessions[0].provider.sequence, sequence);
4098 }
4099
4100 #[test]
4112 fn acp_runtime_policy_admits_structured_prompt_and_refuses_raw_pty_input_by_name() {
4113 let acp_policy = ProviderRuntimePolicy::new(false, true, true, true, false, false)
4114 .expect("ACP-shaped runtime policy is internally valid");
4115 let (mut engine, _spawn) = running_acp_engine_with_policy(acp_policy);
4116
4117 engine
4118 .apply_command(CommandEnvelope {
4119 id: CommandId(3),
4120 command: ControlCommand::SendInput {
4121 instance_id: instance(),
4122 action: InputAction::SubmitPrompt(PromptPayload {
4123 text: "acp prompt".to_owned(),
4124 framing: PromptFraming::Literal,
4125 }),
4126 },
4127 })
4128 .expect("ACP structured prompt must be admitted by an ACP-shaped runtime policy");
4129 let submit = engine.drain_effects().pop().unwrap();
4130 assert!(matches!(submit.effect, ControlEffect::SubmitPrompt { .. }));
4131 engine.apply_observation(ObservationEnvelope {
4132 operation_id: Some(submit.operation_id),
4133 instance_id: instance(),
4134 generation: submit.generation,
4135 observation: ControlObservation::InputCompleted,
4136 });
4137
4138 assert_eq!(
4139 engine.apply_command(CommandEnvelope {
4140 id: CommandId(4),
4141 command: ControlCommand::SendInput {
4142 instance_id: instance(),
4143 action: InputAction::TerminalBytes(vec![0x1b, b'[', b'A']),
4144 },
4145 }),
4146 Err(ControlError::UnsupportedTransportOperation {
4147 transport: TransportKind::Acp,
4148 action: "this input action".to_owned(),
4149 }),
4150 );
4151 assert!(engine.drain_effects().is_empty());
4152 }
4153
4154 #[test]
4161 fn acp_runtime_policy_admits_a_non_identity_provider_event() {
4162 let acp_policy = ProviderRuntimePolicy::new(false, true, true, true, false, false)
4163 .expect("ACP-shaped runtime policy is internally valid");
4164 let (mut engine, spawn) = running_acp_engine_with_policy(acp_policy);
4165
4166 engine.apply_observation(ObservationEnvelope {
4167 operation_id: None,
4168 instance_id: instance(),
4169 generation: spawn.generation,
4170 observation: ControlObservation::ProviderEvent {
4171 source: provider_source(),
4172 sequence: 1,
4173 event: ProviderEvent::WorkingObserved,
4174 },
4175 });
4176 assert!(!engine.drain_events().iter().any(|event| matches!(
4177 event.event,
4178 ControlEventKind::ObservationIgnored {
4179 reason: ObservationIgnoredReason::ProviderRuntimePolicyDenied { .. }
4180 }
4181 )));
4182 assert_eq!(engine.snapshot().sessions[0].provider.sequence, 1);
4183 }
4184
4185 #[test]
4198 fn acp_runtime_policy_survives_a_capability_refusal_and_surfaces_it() {
4199 let acp_policy = ProviderRuntimePolicy::new(false, true, true, false, false, false)
4203 .expect("ACP-shaped runtime policy with identity withheld is internally valid");
4204 let (mut engine, spawn) = running_acp_engine_with_policy(acp_policy);
4205
4206 engine.apply_observation(ObservationEnvelope {
4207 operation_id: None,
4208 instance_id: instance(),
4209 generation: spawn.generation,
4210 observation: ControlObservation::ProviderEvent {
4211 source: provider_source(),
4212 sequence: 1,
4213 event: ProviderEvent::SessionIdentityObserved {
4214 identity: ProviderSessionIdentity {
4215 key: ProviderSessionKey::SessionId,
4216 id: "acp-session".to_owned(),
4217 transcript_path: None,
4218 },
4219 },
4220 },
4221 });
4222 let refusal_events = engine.drain_events();
4223 assert!(
4224 refusal_events.iter().any(|event| matches!(
4225 &event.event,
4226 ControlEventKind::ProviderEvent {
4227 event: ProviderEvent::Error { message },
4228 ..
4229 } if message.contains("ProviderSessionIdentity") && message.contains("1 event")
4230 )),
4231 "a denied observation must mint a visible refusal, not silence: {refusal_events:?}",
4232 );
4233 assert!(!refusal_events.iter().any(|event| matches!(
4234 event.event,
4235 ControlEventKind::ObservationIgnored {
4236 reason: ObservationIgnoredReason::ProviderRuntimePolicyDenied { .. }
4237 }
4238 )));
4239 assert_eq!(
4240 engine.snapshot().sessions[0].provider.sequence,
4241 1,
4242 "the refusal still advances the canonical provider sequence",
4243 );
4244
4245 engine.apply_observation(ObservationEnvelope {
4249 operation_id: None,
4250 instance_id: instance(),
4251 generation: spawn.generation,
4252 observation: ControlObservation::ProviderEvent {
4253 source: provider_source(),
4254 sequence: 2,
4255 event: ProviderEvent::WorkingObserved,
4256 },
4257 });
4258 let follow_up_events = engine.drain_events();
4259 assert!(
4260 !follow_up_events.iter().any(|event| matches!(
4261 event.event,
4262 ControlEventKind::ObservationIgnored {
4263 reason: ObservationIgnoredReason::StaleProviderEvent
4264 | ObservationIgnoredReason::ProviderRuntimePolicyDenied { .. },
4265 }
4266 )),
4267 "a later, capability-clean observation must not be dropped either: {follow_up_events:?}",
4268 );
4269 assert_eq!(engine.snapshot().sessions[0].provider.sequence, 2);
4270 }
4271
4272 #[test]
4273 fn runtime_policy_rejects_invalid_start_contract() {
4274 let mut engine = Gate4AgentEngine::new();
4275 engine.apply_command(register(1)).unwrap();
4276 let invalid = ProviderRuntimePolicy {
4290 raw_pty_lifecycle: false,
4291 semantic_readiness: true,
4292 structured_prompt: false,
4293 provider_session_identity: true,
4294 semantic_resume: false,
4295 hook_semantics: true,
4296 };
4297 assert_eq!(
4298 engine.apply_command(start_with_policy(2, invalid)),
4299 Err(ControlError::InvalidProviderRuntimePolicy {
4300 error: gate4agent_types::ProviderRuntimePolicyError::SemanticCapabilityRequiresRawPty,
4301 })
4302 );
4303 assert!(engine.drain_effects().is_empty());
4304 }
4305
4306 #[test]
4307 fn raw_pty_policy_admits_prompt_free_native_resume_only() {
4308 let mut engine = Gate4AgentEngine::new();
4309 engine.apply_command(register(1)).unwrap();
4310 engine.drain_events();
4311 let identity = ProviderSessionIdentity {
4312 key: ProviderSessionKey::SessionId,
4313 id: "provider-session".to_owned(),
4314 transcript_path: None,
4315 };
4316 let request = ResumeLaunchRequest {
4317 working_directory: ".".to_owned(),
4318 terminal_size: TerminalSize {
4319 rows: 24,
4320 columns: 80,
4321 },
4322 initial_prompt: None,
4323 };
4324 engine
4325 .apply_command(CommandEnvelope {
4326 id: CommandId(2),
4327 command: ControlCommand::Resume {
4328 instance_id: instance(),
4329 target: ResumeTarget::ProviderSession {
4330 identity: identity.clone(),
4331 },
4332 runtime_policy: ProviderRuntimePolicy::raw_pty(),
4333 request,
4334 },
4335 })
4336 .unwrap();
4337 let authorize = engine.drain_effects().pop().unwrap();
4338 assert!(matches!(
4339 &authorize.effect,
4340 ControlEffect::AuthorizeResume {
4341 target: ResumeAuthorityTarget::ProviderSession { identity: requested },
4342 ..
4343 } if requested == &identity
4344 ));
4345 engine.apply_observation(ObservationEnvelope {
4346 operation_id: Some(authorize.operation_id),
4347 instance_id: instance(),
4348 generation: authorize.generation,
4349 observation: ControlObservation::ResumeAuthorized {
4350 provider_session: identity.clone(),
4351 },
4352 });
4353 assert!(matches!(
4354 engine.drain_effects().pop().unwrap().effect,
4355 ControlEffect::SpawnResume {
4356 provider_session,
4357 runtime_policy,
4358 request: ResumeLaunchRequest { initial_prompt: None, .. },
4359 ..
4360 } if provider_session == identity
4361 && runtime_policy == ProviderRuntimePolicy::raw_pty()
4362 ));
4363
4364 let mut prompted = Gate4AgentEngine::new();
4365 prompted.apply_command(register(1)).unwrap();
4366 assert_eq!(
4367 prompted.apply_command(CommandEnvelope {
4368 id: CommandId(2),
4369 command: ControlCommand::Resume {
4370 instance_id: instance(),
4371 target: ResumeTarget::ProviderSession { identity },
4372 runtime_policy: ProviderRuntimePolicy::raw_pty(),
4373 request: ResumeLaunchRequest {
4374 working_directory: ".".to_owned(),
4375 terminal_size: TerminalSize {
4376 rows: 24,
4377 columns: 80,
4378 },
4379 initial_prompt: Some("continue".to_owned()),
4380 },
4381 },
4382 }),
4383 Err(ControlError::ProviderRuntimePolicyDenied {
4384 capability: ProviderRuntimeCapability::ProviderSessionIdentity,
4385 })
4386 );
4387 assert!(prompted.drain_effects().is_empty());
4388 }
4389
4390 fn resolving_interaction(
4391 interaction_kind: ProviderInteractionKind,
4392 response: ProviderInteractionResponse,
4393 ) -> (
4394 Gate4AgentEngine,
4395 EffectEnvelope,
4396 OperationId,
4397 ProviderInteractionId,
4398 ) {
4399 let (mut engine, spawn) = running_engine();
4400 let interaction_id = ProviderInteractionId(1);
4401 engine.apply_observation(ObservationEnvelope {
4402 operation_id: None,
4403 instance_id: spawn.instance_id,
4404 generation: spawn.generation,
4405 observation: ControlObservation::ProviderEvent {
4406 source: provider_source(),
4407 sequence: 1,
4408 event: ProviderEvent::InteractionRequested {
4409 request_id: Some("request-1".to_owned()),
4410 interaction_kind,
4411 tool_name: "fixture-tool".to_owned(),
4412 title: None,
4413 prompt: "continue?".to_owned(),
4414 options: Vec::new(),
4415 agent_id: None,
4416 },
4417 },
4418 });
4419 engine
4420 .apply_command(CommandEnvelope {
4421 id: CommandId(90),
4422 command: ControlCommand::ResolveInteraction {
4423 instance_id: instance(),
4424 generation: spawn.generation,
4425 interaction_id,
4426 response,
4427 },
4428 })
4429 .unwrap();
4430 let operation_id = engine.snapshot().sessions[0]
4431 .pending_operation
4432 .expect("interaction resolution operation");
4433 (engine, spawn, operation_id, interaction_id)
4434 }
4435
4436 fn assert_late_interaction_failure_is_ignored(
4437 engine: &mut Gate4AgentEngine,
4438 spawn: &EffectEnvelope,
4439 operation_id: OperationId,
4440 interaction_id: ProviderInteractionId,
4441 expected_outcome: ProviderInteractionOutcome,
4442 ) {
4443 engine.drain_events();
4444 engine.apply_observation(ObservationEnvelope {
4445 operation_id: Some(operation_id),
4446 instance_id: spawn.instance_id,
4447 generation: spawn.generation,
4448 observation: ControlObservation::InteractionResolutionFailed {
4449 interaction_id,
4450 message: "late executor failure".to_owned(),
4451 },
4452 });
4453 assert_eq!(
4454 engine.snapshot().sessions[0].provider.interactions[0].status,
4455 ProviderInteractionStatus::Resolved {
4456 outcome: expected_outcome,
4457 }
4458 );
4459 assert!(engine.drain_events().iter().any(|event| matches!(
4460 event.event,
4461 ControlEventKind::ObservationIgnored {
4462 reason: ObservationIgnoredReason::OperationMismatch,
4463 }
4464 )));
4465 }
4466
4467 fn inactive_engine_with_provider_session() -> (Gate4AgentEngine, ProviderSessionIdentity) {
4468 let (mut engine, spawn) = running_engine();
4469 let identity = ProviderSessionIdentity {
4470 key: ProviderSessionKey::SessionId,
4471 id: "provider-session-1".to_owned(),
4472 transcript_path: None,
4473 };
4474 engine.apply_observation(ObservationEnvelope {
4475 operation_id: None,
4476 instance_id: instance(),
4477 generation: spawn.generation,
4478 observation: ControlObservation::ProviderEvent {
4479 source: provider_source(),
4480 sequence: 1,
4481 event: ProviderEvent::SessionIdentityObserved {
4482 identity: identity.clone(),
4483 },
4484 },
4485 });
4486 engine.apply_observation(ObservationEnvelope {
4487 operation_id: None,
4488 instance_id: instance(),
4489 generation: spawn.generation,
4490 observation: ControlObservation::ProcessExited {
4491 exit_code: Some(0),
4492 final_terminal: None,
4493 },
4494 });
4495 engine.drain_events();
4496 (engine, identity)
4497 }
4498
4499 #[test]
4500 fn spawn_is_not_reported_running_before_observation() {
4501 let mut engine = Gate4AgentEngine::new();
4502 engine.apply_command(register(1)).unwrap();
4503 engine.apply_command(start(2)).unwrap();
4504
4505 let snapshot = engine.snapshot();
4506 assert_eq!(snapshot.sessions[0].status, SessionStatus::Starting);
4507 assert_eq!(snapshot.sessions[0].process_id, None);
4508 assert_eq!(engine.drain_effects().len(), 1);
4509 }
4510
4511 #[test]
4512 fn capability_probe_settles_across_session_generation_without_blocking_start() {
4513 let mut engine = Gate4AgentEngine::new();
4514 engine.apply_command(register(1)).unwrap();
4515 engine
4516 .apply_command(CommandEnvelope {
4517 id: CommandId(2),
4518 command: ControlCommand::ProbeCapabilities {
4519 instance_id: instance(),
4520 request: CapabilityProbeRequest {
4521 working_directory: ".".to_owned(),
4522 },
4523 },
4524 })
4525 .unwrap();
4526 let probe = engine.drain_effects().pop().unwrap();
4527
4528 engine.apply_command(start(3)).unwrap();
4529 let spawn = engine.drain_effects().pop().unwrap();
4530 assert_ne!(probe.generation, spawn.generation);
4531 assert_eq!(
4532 engine.snapshot().sessions[0].status,
4533 SessionStatus::Starting
4534 );
4535
4536 engine.apply_observation(ObservationEnvelope {
4537 operation_id: Some(probe.operation_id),
4538 instance_id: probe.instance_id,
4539 generation: probe.generation,
4540 observation: ControlObservation::CapabilitiesProbed {
4541 session_option_models: vec![
4542 CapabilityModelSummary {
4543 id: "duplicate".to_owned(),
4544 label: "Duplicate".to_owned(),
4545 },
4546 CapabilityModelSummary {
4547 id: "duplicate".to_owned(),
4548 label: "Duplicate again".to_owned(),
4549 },
4550 ],
4551 },
4552 });
4553 assert!(engine.snapshot().sessions[0].capabilities.pending.is_some());
4554
4555 engine.apply_observation(ObservationEnvelope {
4556 operation_id: Some(probe.operation_id),
4557 instance_id: probe.instance_id,
4558 generation: probe.generation,
4559 observation: ControlObservation::CapabilitiesProbed {
4560 session_option_models: vec![CapabilityModelSummary {
4561 id: "account-model".to_owned(),
4562 label: "Account model".to_owned(),
4563 }],
4564 },
4565 });
4566 let snapshot = engine.snapshot();
4567 assert_eq!(snapshot.sessions[0].status, SessionStatus::Starting);
4568 assert_eq!(
4569 snapshot.sessions[0].pending_operation,
4570 Some(spawn.operation_id)
4571 );
4572 assert!(snapshot.sessions[0].capabilities.settled);
4573 assert_eq!(
4574 snapshot.sessions[0].capabilities.session_option_models[0].id,
4575 "account-model"
4576 );
4577 assert_eq!(
4578 engine
4579 .apply_command(CommandEnvelope {
4580 id: CommandId(4),
4581 command: ControlCommand::ProbeCapabilities {
4582 instance_id: instance(),
4583 request: CapabilityProbeRequest {
4584 working_directory: ".".to_owned(),
4585 },
4586 },
4587 })
4588 .unwrap_err(),
4589 ControlError::CapabilityProbeSettled
4590 );
4591 }
4592
4593 #[test]
4594 fn unbounded_terminal_geometry_is_rejected_before_effect_creation() {
4595 let mut engine = Gate4AgentEngine::new();
4596 engine.apply_command(register(1)).unwrap();
4597 let error = engine
4598 .apply_command(CommandEnvelope {
4599 id: CommandId(2),
4600 command: ControlCommand::Start {
4601 instance_id: instance(),
4602 runtime_policy: verified_runtime_policy(),
4603 request: StartRequest {
4604 working_directory: ".".to_owned(),
4605 terminal_size: TerminalSize {
4606 rows: 1_001,
4607 columns: 80,
4608 },
4609 initial_prompt: None,
4610 session_options: None,
4611 approval_level: ApprovalLevel::default(),
4612 },
4613 },
4614 })
4615 .unwrap_err();
4616 assert_eq!(error, ControlError::InvalidTerminalSize);
4617 assert!(engine.drain_effects().is_empty());
4618 }
4619
4620 #[test]
4621 fn invalid_session_options_are_rejected_before_effect_creation() {
4622 let mut engine = Gate4AgentEngine::new();
4623 engine.apply_command(register(1)).unwrap();
4624 let mut request = match start(2).command {
4625 ControlCommand::Start { request, .. } => request,
4626 _ => unreachable!(),
4627 };
4628 request.session_options = Some(SessionOptionSelection::new("bad\nmodel"));
4629 let error = engine
4630 .apply_command(CommandEnvelope {
4631 id: CommandId(2),
4632 command: ControlCommand::Start {
4633 instance_id: instance(),
4634 runtime_policy: verified_runtime_policy(),
4635 request,
4636 },
4637 })
4638 .unwrap_err();
4639 assert!(matches!(error, ControlError::InvalidSessionOptions { .. }));
4640 assert!(engine.drain_effects().is_empty());
4641 }
4642
4643 #[test]
4644 fn direct_provider_gap_advances_source_cursor_and_accepts_next_event() {
4645 let (mut engine, spawn) = running_engine();
4646 engine.apply_observation(ObservationEnvelope {
4647 operation_id: None,
4648 instance_id: spawn.instance_id,
4649 generation: spawn.generation,
4650 observation: ControlObservation::ProviderEvent {
4651 source: provider_source(),
4652 sequence: 1,
4653 event: ProviderEvent::TurnCompleted {
4654 usage: TokenUsage {
4655 input_tokens: 3,
4656 output_tokens: 5,
4657 context_window: Some(100),
4658 ..TokenUsage::default()
4659 },
4660 is_cumulative: false,
4661 },
4662 },
4663 });
4664 engine.apply_observation(ObservationEnvelope {
4665 operation_id: None,
4666 instance_id: spawn.instance_id,
4667 generation: spawn.generation,
4668 observation: ControlObservation::ProviderGap {
4669 source: provider_source(),
4670 source_sequence: 3,
4671 missed: 2,
4672 },
4673 });
4674 engine.apply_observation(ObservationEnvelope {
4675 operation_id: None,
4676 instance_id: spawn.instance_id,
4677 generation: spawn.generation,
4678 observation: ControlObservation::ProviderEvent {
4679 source: provider_source(),
4680 sequence: 4,
4681 event: ProviderEvent::WorkingObserved,
4682 },
4683 });
4684 engine.apply_observation(ObservationEnvelope {
4685 operation_id: None,
4686 instance_id: spawn.instance_id,
4687 generation: spawn.generation,
4688 observation: ControlObservation::ProviderEvent {
4689 source: provider_source(),
4690 sequence: 1,
4691 event: ProviderEvent::TurnCompleted {
4692 usage: TokenUsage {
4693 input_tokens: 99,
4694 ..TokenUsage::default()
4695 },
4696 is_cumulative: false,
4697 },
4698 },
4699 });
4700
4701 let snapshot = engine.snapshot();
4702 let provider = &snapshot.sessions[0].provider;
4703 assert_eq!(provider.completed_turns, 1);
4704 assert_eq!(provider.usage.input_tokens, 3);
4705 assert_eq!(provider.usage.output_tokens, 5);
4706 assert_eq!(provider.gap_count, 2);
4707 assert!(!provider.stale);
4708 assert_eq!(provider.sources[0].sequence, 4);
4709 let events = engine.drain_events();
4710 assert!(events.iter().any(|event| matches!(
4711 event.event,
4712 ControlEventKind::ProviderGap {
4713 source_sequence: 3,
4714 missed: 2,
4715 ..
4716 }
4717 )));
4718 assert!(events.iter().any(|event| matches!(
4719 event.event,
4720 ControlEventKind::ObservationIgnored {
4721 reason: ObservationIgnoredReason::StaleProviderEvent
4722 }
4723 )));
4724 }
4725
4726 #[test]
4727 fn provider_activity_tracks_turn_tools_and_attention_without_late_text_resurrection() {
4728 let (mut engine, spawn) = running_engine();
4729 let events = [
4730 ProviderEvent::TurnStarted {
4731 prompt: Some("fix tests".to_owned()),
4732 },
4733 ProviderEvent::ToolStarted {
4734 id: "tool-1".to_owned(),
4735 name: "shell".to_owned(),
4736 input_json: "{\"command\":\"cargo test\"}".to_owned(),
4737 agent_id: None,
4738 },
4739 ProviderEvent::InteractionRequested {
4740 request_id: Some("tool-1".to_owned()),
4741 interaction_kind: ProviderInteractionKind::Approval,
4742 tool_name: "shell".to_owned(),
4743 title: None,
4744 prompt: "approve".to_owned(),
4745 options: Vec::new(),
4746 agent_id: None,
4747 },
4748 ProviderEvent::ToolCompleted {
4749 id: "tool-1".to_owned(),
4750 output: "ok".to_owned(),
4751 is_error: false,
4752 duration_ms: None,
4753 agent_id: None,
4754 non_execution_kind: None,
4755 },
4756 ProviderEvent::TurnCompleted {
4757 usage: TokenUsage::default(),
4758 is_cumulative: false,
4759 },
4760 ProviderEvent::Text {
4761 text: "late final text".to_owned(),
4762 is_delta: false,
4763 },
4764 ];
4765
4766 for (index, event) in events.into_iter().enumerate() {
4767 engine.apply_observation(ObservationEnvelope {
4768 operation_id: None,
4769 instance_id: spawn.instance_id,
4770 generation: spawn.generation,
4771 observation: ControlObservation::ProviderEvent {
4772 source: provider_source(),
4773 sequence: u64::try_from(index + 1).unwrap(),
4774 event,
4775 },
4776 });
4777 let provider = &engine.snapshot().sessions[0].provider;
4778 match index {
4779 0 => {
4780 assert_eq!(provider.activity, ProviderActivity::Working);
4781 assert_eq!(provider.current_prompt.as_deref(), Some("fix tests"));
4782 }
4783 1 => assert_eq!(provider.active_tools.len(), 1),
4784 2 => assert_eq!(provider.activity, ProviderActivity::WaitingForInput),
4785 3 => {
4786 assert!(provider.active_tools.is_empty());
4787 assert_eq!(provider.activity, ProviderActivity::Working);
4788 assert_eq!(
4789 provider.interactions[0].status,
4790 ProviderInteractionStatus::Resolved {
4791 outcome: ProviderInteractionOutcome::Approved
4792 }
4793 );
4794 }
4795 4 | 5 => assert_eq!(provider.activity, ProviderActivity::Idle),
4796 _ => unreachable!(),
4797 }
4798 }
4799 }
4800
4801 #[test]
4802 fn provider_question_has_canonical_identity_and_matching_tool_resolution() {
4803 let (mut engine, spawn) = running_engine();
4804 engine.drain_events();
4805 let source = provider_source();
4806 engine.apply_observation(ObservationEnvelope {
4807 operation_id: None,
4808 instance_id: spawn.instance_id,
4809 generation: spawn.generation,
4810 observation: ControlObservation::ProviderEvent {
4811 source: source.clone(),
4812 sequence: 1,
4813 event: ProviderEvent::InteractionRequested {
4814 request_id: Some("question-7".to_owned()),
4815 interaction_kind: ProviderInteractionKind::Question,
4816 tool_name: "AskUserQuestion".to_owned(),
4817 title: None,
4818 prompt: "{\"question\":\"Continue?\"}".to_owned(),
4819 options: Vec::new(),
4820 agent_id: None,
4821 },
4822 },
4823 });
4824
4825 let interaction = engine.snapshot().sessions[0].provider.interactions[0].clone();
4826 assert_eq!(interaction.id, ProviderInteractionId(1));
4827 assert_eq!(interaction.source, source);
4828 assert_eq!(interaction.status, ProviderInteractionStatus::Pending);
4829 assert!(engine.drain_events().iter().any(|event| matches!(
4830 &event.event,
4831 ControlEventKind::InteractionRequested { interaction: observed }
4832 if observed.id == interaction.id
4833 )));
4834
4835 engine.apply_observation(ObservationEnvelope {
4836 operation_id: None,
4837 instance_id: spawn.instance_id,
4838 generation: spawn.generation,
4839 observation: ControlObservation::ProviderEvent {
4840 source: provider_source(),
4841 sequence: 2,
4842 event: ProviderEvent::ToolStarted {
4843 id: "question-7".to_owned(),
4844 name: "AskUserQuestion".to_owned(),
4845 input_json: "{}".to_owned(),
4846 agent_id: None,
4847 },
4848 },
4849 });
4850
4851 let interaction = &engine.snapshot().sessions[0].provider.interactions[0];
4852 assert_eq!(
4853 interaction.status,
4854 ProviderInteractionStatus::Resolved {
4855 outcome: ProviderInteractionOutcome::Answered
4856 }
4857 );
4858 assert!(engine.drain_events().iter().any(|event| matches!(
4859 event.event,
4860 ControlEventKind::InteractionResolved {
4861 interaction_id: ProviderInteractionId(1),
4862 outcome: ProviderInteractionOutcome::Answered,
4863 }
4864 )));
4865 }
4866
4867 #[test]
4868 fn provider_reported_interaction_resolution_is_exact_and_immediate() {
4869 for outcome in [
4870 ProviderInteractionOutcome::Approved,
4871 ProviderInteractionOutcome::Denied,
4872 ] {
4873 let (mut engine, spawn) = running_engine();
4874 let source = provider_source();
4875 engine.apply_observation(ObservationEnvelope {
4876 operation_id: None,
4877 instance_id: spawn.instance_id,
4878 generation: spawn.generation,
4879 observation: ControlObservation::ProviderEvent {
4880 source: source.clone(),
4881 sequence: 1,
4882 event: ProviderEvent::InteractionRequested {
4883 request_id: Some("approval-1".to_owned()),
4884 interaction_kind: ProviderInteractionKind::Approval,
4885 tool_name: "shell".to_owned(),
4886 title: None,
4887 prompt: String::new(),
4888 options: Vec::new(),
4889 agent_id: None,
4890 },
4891 },
4892 });
4893 engine.drain_events();
4894 engine.apply_observation(ObservationEnvelope {
4895 operation_id: None,
4896 instance_id: spawn.instance_id,
4897 generation: spawn.generation,
4898 observation: ControlObservation::ProviderEvent {
4899 source,
4900 sequence: 2,
4901 event: ProviderEvent::InteractionResolved {
4902 request_id: "approval-1".to_owned(),
4903 outcome,
4904 },
4905 },
4906 });
4907
4908 assert_eq!(
4909 engine.snapshot().sessions[0].provider.interactions[0].status,
4910 ProviderInteractionStatus::Resolved { outcome }
4911 );
4912 assert_eq!(
4913 engine.snapshot().sessions[0].provider.lead_activity,
4914 ProviderActivity::Working
4915 );
4916 assert_eq!(
4917 engine
4918 .drain_events()
4919 .iter()
4920 .filter(|event| matches!(
4921 event.event,
4922 ControlEventKind::InteractionResolved {
4923 interaction_id: ProviderInteractionId(1),
4924 outcome: observed,
4925 } if observed == outcome
4926 ))
4927 .count(),
4928 1
4929 );
4930 }
4931 }
4932
4933 #[test]
4934 fn provider_reported_resolution_is_fail_closed_for_orphan_duplicate_and_lifecycle() {
4935 let (mut engine, spawn) = running_engine();
4936 let source = provider_source();
4937 engine.apply_observation(ObservationEnvelope {
4938 operation_id: None,
4939 instance_id: spawn.instance_id,
4940 generation: spawn.generation,
4941 observation: ControlObservation::ProviderEvent {
4942 source: source.clone(),
4943 sequence: 1,
4944 event: ProviderEvent::InteractionResolved {
4945 request_id: "orphan".to_owned(),
4946 outcome: ProviderInteractionOutcome::Approved,
4947 },
4948 },
4949 });
4950 assert!(engine.snapshot().sessions[0].provider.interactions.is_empty());
4951 assert!(!engine.drain_events().iter().any(|event| matches!(
4952 event.event,
4953 ControlEventKind::InteractionResolved { .. }
4954 )));
4955
4956 engine.apply_observation(ObservationEnvelope {
4957 operation_id: None,
4958 instance_id: spawn.instance_id,
4959 generation: spawn.generation,
4960 observation: ControlObservation::ProviderEvent {
4961 source: source.clone(),
4962 sequence: 2,
4963 event: ProviderEvent::InteractionRequested {
4964 request_id: Some("approval-2".to_owned()),
4965 interaction_kind: ProviderInteractionKind::Approval,
4966 tool_name: "shell".to_owned(),
4967 title: None,
4968 prompt: String::new(),
4969 options: Vec::new(),
4970 agent_id: None,
4971 },
4972 },
4973 });
4974 engine.drain_events();
4975 for sequence in [3, 4] {
4976 engine.apply_observation(ObservationEnvelope {
4977 operation_id: None,
4978 instance_id: spawn.instance_id,
4979 generation: spawn.generation,
4980 observation: ControlObservation::ProviderEvent {
4981 source: source.clone(),
4982 sequence,
4983 event: ProviderEvent::InteractionResolved {
4984 request_id: "approval-2".to_owned(),
4985 outcome: ProviderInteractionOutcome::Denied,
4986 },
4987 },
4988 });
4989 }
4990 let events = engine.drain_events();
4991 assert_eq!(
4992 events
4993 .iter()
4994 .filter(|event| matches!(
4995 event.event,
4996 ControlEventKind::InteractionResolved {
4997 interaction_id: ProviderInteractionId(2),
4998 outcome: ProviderInteractionOutcome::Denied,
4999 }
5000 ))
5001 .count(),
5002 1
5003 );
5004
5005 engine.apply_observation(ObservationEnvelope {
5006 operation_id: None,
5007 instance_id: spawn.instance_id,
5008 generation: spawn.generation,
5009 observation: ControlObservation::ProviderEvent {
5010 source,
5011 sequence: 5,
5012 event: ProviderEvent::TurnCompleted {
5013 usage: TokenUsage::default(),
5014 is_cumulative: false,
5015 },
5016 },
5017 });
5018 assert_eq!(
5019 engine.snapshot().sessions[0].provider.interactions[0].status,
5020 ProviderInteractionStatus::Resolved {
5021 outcome: ProviderInteractionOutcome::Denied,
5022 }
5023 );
5024 assert!(!engine.drain_events().iter().any(|event| matches!(
5025 event.event,
5026 ControlEventKind::InteractionResolved { .. }
5027 )));
5028 }
5029
5030 #[test]
5031 fn interrupt_resolves_pending_interactions_only_after_effect_completion() {
5032 let (mut engine, spawn) = running_engine();
5033 engine.apply_observation(ObservationEnvelope {
5034 operation_id: None,
5035 instance_id: spawn.instance_id,
5036 generation: spawn.generation,
5037 observation: ControlObservation::ProviderEvent {
5038 source: provider_source(),
5039 sequence: 1,
5040 event: ProviderEvent::InteractionRequested {
5041 request_id: Some("approval-1".to_owned()),
5042 interaction_kind: ProviderInteractionKind::Approval,
5043 tool_name: "shell".to_owned(),
5044 title: None,
5045 prompt: "cargo test".to_owned(),
5046 options: Vec::new(),
5047 agent_id: None,
5048 },
5049 },
5050 });
5051 engine.drain_events();
5052
5053 let send_interrupt = |id| CommandEnvelope {
5054 id: CommandId(id),
5055 command: ControlCommand::SendInput {
5056 instance_id: instance(),
5057 action: InputAction::TerminalControl(TerminalControl::Interrupt),
5058 },
5059 };
5060 engine.apply_command(send_interrupt(80)).unwrap();
5061 let failed = engine.drain_effects().pop().unwrap();
5062 engine.apply_observation(ObservationEnvelope {
5063 operation_id: Some(failed.operation_id),
5064 instance_id: failed.instance_id,
5065 generation: failed.generation,
5066 observation: ControlObservation::InputFailed {
5067 message: "write rejected".to_owned(),
5068 },
5069 });
5070 assert_eq!(
5071 engine.snapshot().sessions[0].provider.interactions[0].status,
5072 ProviderInteractionStatus::Pending
5073 );
5074
5075 engine.apply_command(send_interrupt(81)).unwrap();
5076 let completed = engine.drain_effects().pop().unwrap();
5077 engine.apply_observation(ObservationEnvelope {
5078 operation_id: Some(completed.operation_id),
5079 instance_id: completed.instance_id,
5080 generation: completed.generation,
5081 observation: ControlObservation::InputCompleted,
5082 });
5083 assert_eq!(
5084 engine.snapshot().sessions[0].provider.interactions[0].status,
5085 ProviderInteractionStatus::Resolved {
5086 outcome: ProviderInteractionOutcome::Interrupted
5087 }
5088 );
5089 assert!(engine.drain_events().iter().any(|event| matches!(
5090 event.event,
5091 ControlEventKind::InteractionResolved {
5092 interaction_id: ProviderInteractionId(1),
5093 outcome: ProviderInteractionOutcome::Interrupted,
5094 }
5095 )));
5096 }
5097
5098 #[test]
5099 fn canonical_interaction_resolution_is_generation_checked_and_fail_closed() {
5100 let (mut engine, spawn) = running_engine();
5101 engine.apply_observation(ObservationEnvelope {
5102 operation_id: None,
5103 instance_id: spawn.instance_id,
5104 generation: spawn.generation,
5105 observation: ControlObservation::ProviderEvent {
5106 source: provider_source(),
5107 sequence: 1,
5108 event: ProviderEvent::InteractionRequested {
5109 request_id: Some("question-1".to_owned()),
5110 interaction_kind: ProviderInteractionKind::Question,
5111 tool_name: "AskUserQuestion".to_owned(),
5112 title: None,
5113 prompt: "continue?".to_owned(),
5114 options: Vec::new(),
5115 agent_id: None,
5116 },
5117 },
5118 });
5119 engine.drain_events();
5120 let interaction_id = ProviderInteractionId(1);
5121 let resolve = |id, generation, response| CommandEnvelope {
5122 id: CommandId(id),
5123 command: ControlCommand::ResolveInteraction {
5124 instance_id: instance(),
5125 generation,
5126 interaction_id,
5127 response,
5128 },
5129 };
5130
5131 assert_eq!(
5132 engine
5133 .apply_command(resolve(
5134 90,
5135 SessionGeneration(spawn.generation.0.saturating_add(1)),
5136 ProviderInteractionResponse::Answer {
5137 text: "yes".to_owned(),
5138 },
5139 ))
5140 .unwrap_err(),
5141 ControlError::StaleProviderInteractionGeneration {
5142 expected: spawn.generation,
5143 actual: SessionGeneration(spawn.generation.0.saturating_add(1)),
5144 }
5145 );
5146 assert!(engine.drain_effects().is_empty());
5147 assert!(matches!(
5148 engine
5149 .apply_command(resolve(
5150 91,
5151 spawn.generation,
5152 ProviderInteractionResponse::ApproveOnce,
5153 ))
5154 .unwrap_err(),
5155 ControlError::InvalidProviderInteractionResponse { .. }
5156 ));
5157
5158 engine
5159 .apply_command(resolve(
5160 92,
5161 spawn.generation,
5162 ProviderInteractionResponse::Answer {
5163 text: "yes".to_owned(),
5164 },
5165 ))
5166 .unwrap();
5167 let effect = engine.drain_effects().pop().unwrap();
5168 assert!(matches!(
5169 &effect.effect,
5170 ControlEffect::ResolveInteraction {
5171 target,
5172 response: ProviderInteractionResponse::Answer { text },
5173 } if target.interaction_id == interaction_id
5174 && target.provider_request_id.as_deref() == Some("question-1")
5175 && text == "yes"
5176 ));
5177 assert_eq!(
5178 engine.snapshot().sessions[0].provider.interactions[0].status,
5179 ProviderInteractionStatus::Resolving {
5180 operation_id: effect.operation_id,
5181 response_kind: ProviderInteractionResponseKind::Answer,
5182 }
5183 );
5184 assert!(engine.drain_events().iter().any(|event| matches!(
5185 event.event,
5186 ControlEventKind::InteractionResolutionRequested {
5187 operation_id,
5188 interaction_id: ProviderInteractionId(1),
5189 response_kind: ProviderInteractionResponseKind::Answer,
5190 } if operation_id == effect.operation_id
5191 )));
5192
5193 engine.apply_observation(ObservationEnvelope {
5194 operation_id: Some(effect.operation_id),
5195 instance_id: effect.instance_id,
5196 generation: effect.generation,
5197 observation: ControlObservation::InteractionResolutionCompleted {
5198 interaction_id: ProviderInteractionId(999),
5199 },
5200 });
5201 assert!(matches!(
5202 engine.snapshot().sessions[0].provider.interactions[0].status,
5203 ProviderInteractionStatus::Resolving { .. }
5204 ));
5205 assert!(engine.drain_events().iter().any(|event| matches!(
5206 event.event,
5207 ControlEventKind::ObservationIgnored {
5208 reason: ObservationIgnoredReason::InvalidInteractionObservation,
5209 }
5210 )));
5211
5212 engine.apply_observation(ObservationEnvelope {
5213 operation_id: Some(effect.operation_id),
5214 instance_id: effect.instance_id,
5215 generation: effect.generation,
5216 observation: ControlObservation::InteractionResolutionFailed {
5217 interaction_id,
5218 message: "provider rejected response".to_owned(),
5219 },
5220 });
5221 assert_eq!(
5222 engine.snapshot().sessions[0].provider.interactions[0].status,
5223 ProviderInteractionStatus::Pending
5224 );
5225 assert_eq!(engine.snapshot().sessions[0].pending_operation, None);
5226
5227 engine
5228 .apply_command(resolve(
5229 93,
5230 spawn.generation,
5231 ProviderInteractionResponse::Answer {
5232 text: "yes".to_owned(),
5233 },
5234 ))
5235 .unwrap();
5236 let effect = engine.drain_effects().pop().unwrap();
5237 engine.apply_observation(ObservationEnvelope {
5238 operation_id: Some(effect.operation_id),
5239 instance_id: effect.instance_id,
5240 generation: effect.generation,
5241 observation: ControlObservation::InteractionResolutionCompleted { interaction_id },
5242 });
5243 assert_eq!(
5244 engine.snapshot().sessions[0].provider.interactions[0].status,
5245 ProviderInteractionStatus::Resolved {
5246 outcome: ProviderInteractionOutcome::Answered,
5247 }
5248 );
5249 }
5250
5251 #[test]
5252 fn working_progress_settles_denied_resolution_before_late_failure() {
5253 let (mut engine, spawn, operation_id, interaction_id) = resolving_interaction(
5254 ProviderInteractionKind::Question,
5255 ProviderInteractionResponse::Deny,
5256 );
5257
5258 engine.apply_observation(ObservationEnvelope {
5259 operation_id: None,
5260 instance_id: spawn.instance_id,
5261 generation: spawn.generation,
5262 observation: ControlObservation::ProviderEvent {
5263 source: provider_source(),
5264 sequence: 2,
5265 event: ProviderEvent::WorkingObserved,
5266 },
5267 });
5268
5269 assert_eq!(engine.snapshot().sessions[0].pending_operation, None);
5270 assert!(engine.drain_effects().is_empty());
5271 assert_eq!(
5272 engine.snapshot().sessions[0].provider.interactions[0].status,
5273 ProviderInteractionStatus::Resolved {
5274 outcome: ProviderInteractionOutcome::Denied,
5275 }
5276 );
5277 assert_late_interaction_failure_is_ignored(
5278 &mut engine,
5279 &spawn,
5280 operation_id,
5281 interaction_id,
5282 ProviderInteractionOutcome::Denied,
5283 );
5284 }
5285
5286 #[test]
5287 fn exact_tool_progress_settles_resolution_before_late_failure() {
5288 for progress in [
5289 ProviderEvent::ToolStarted {
5290 id: "request-1".to_owned(),
5291 name: "fixture-tool".to_owned(),
5292 input_json: "{}".to_owned(),
5293 agent_id: None,
5294 },
5295 ProviderEvent::ToolCompleted {
5296 id: "request-1".to_owned(),
5297 output: "done".to_owned(),
5298 is_error: false,
5299 duration_ms: Some(1),
5300 agent_id: None,
5301 non_execution_kind: None,
5302 },
5303 ] {
5304 let (mut engine, spawn, operation_id, interaction_id) = resolving_interaction(
5305 ProviderInteractionKind::Approval,
5306 ProviderInteractionResponse::ApproveOnce,
5307 );
5308 engine.apply_observation(ObservationEnvelope {
5309 operation_id: None,
5310 instance_id: spawn.instance_id,
5311 generation: spawn.generation,
5312 observation: ControlObservation::ProviderEvent {
5313 source: provider_source(),
5314 sequence: 2,
5315 event: progress,
5316 },
5317 });
5318
5319 assert_eq!(engine.snapshot().sessions[0].pending_operation, None);
5320 assert!(engine.drain_effects().is_empty());
5321 assert_eq!(
5322 engine.snapshot().sessions[0].provider.interactions[0].status,
5323 ProviderInteractionStatus::Resolved {
5324 outcome: ProviderInteractionOutcome::Approved,
5325 }
5326 );
5327 assert_late_interaction_failure_is_ignored(
5328 &mut engine,
5329 &spawn,
5330 operation_id,
5331 interaction_id,
5332 ProviderInteractionOutcome::Approved,
5333 );
5334 }
5335 }
5336
5337 #[test]
5338 fn terminal_transitions_supersede_in_flight_interaction_resolution() {
5339 for terminate_with_process_exit in [false, true] {
5340 let (mut engine, spawn) = running_engine();
5341 engine.apply_observation(ObservationEnvelope {
5342 operation_id: None,
5343 instance_id: spawn.instance_id,
5344 generation: spawn.generation,
5345 observation: ControlObservation::ProviderEvent {
5346 source: provider_source(),
5347 sequence: 1,
5348 event: ProviderEvent::InteractionRequested {
5349 request_id: Some("approval-1".to_owned()),
5350 interaction_kind: ProviderInteractionKind::Approval,
5351 tool_name: "shell".to_owned(),
5352 title: None,
5353 prompt: "approve".to_owned(),
5354 options: Vec::new(),
5355 agent_id: None,
5356 },
5357 },
5358 });
5359 engine
5360 .apply_command(CommandEnvelope {
5361 id: CommandId(94),
5362 command: ControlCommand::ResolveInteraction {
5363 instance_id: instance(),
5364 generation: spawn.generation,
5365 interaction_id: ProviderInteractionId(1),
5366 response: ProviderInteractionResponse::ApproveOnce,
5367 },
5368 })
5369 .unwrap();
5370 let operation_id = engine.snapshot().sessions[0].pending_operation.unwrap();
5371
5372 if terminate_with_process_exit {
5373 engine.apply_observation(ObservationEnvelope {
5374 operation_id: None,
5375 instance_id: spawn.instance_id,
5376 generation: spawn.generation,
5377 observation: ControlObservation::ProcessExited {
5378 exit_code: Some(0),
5379 final_terminal: None,
5380 },
5381 });
5382 } else {
5383 engine.apply_observation(ObservationEnvelope {
5384 operation_id: None,
5385 instance_id: spawn.instance_id,
5386 generation: spawn.generation,
5387 observation: ControlObservation::ProviderEvent {
5388 source: provider_source(),
5389 sequence: 2,
5390 event: ProviderEvent::TurnCompleted {
5391 usage: TokenUsage::default(),
5392 is_cumulative: false,
5393 },
5394 },
5395 });
5396 }
5397
5398 assert_eq!(engine.snapshot().sessions[0].pending_operation, None);
5399 assert!(engine.drain_effects().is_empty());
5400 assert_eq!(
5401 engine.snapshot().sessions[0].provider.interactions[0].status,
5402 ProviderInteractionStatus::Resolved {
5403 outcome: ProviderInteractionOutcome::TurnEnded,
5404 }
5405 );
5406
5407 engine.drain_events();
5408 engine.apply_observation(ObservationEnvelope {
5409 operation_id: Some(operation_id),
5410 instance_id: spawn.instance_id,
5411 generation: spawn.generation,
5412 observation: ControlObservation::InteractionResolutionCompleted {
5413 interaction_id: ProviderInteractionId(1),
5414 },
5415 });
5416 assert!(engine.drain_events().iter().any(|event| matches!(
5417 event.event,
5418 ControlEventKind::ObservationIgnored {
5419 reason: ObservationIgnoredReason::OperationMismatch,
5420 }
5421 )));
5422 }
5423
5424 let (mut engine, spawn) = running_engine();
5425 engine.apply_observation(ObservationEnvelope {
5426 operation_id: None,
5427 instance_id: spawn.instance_id,
5428 generation: spawn.generation,
5429 observation: ControlObservation::ProviderEvent {
5430 source: provider_source(),
5431 sequence: 1,
5432 event: ProviderEvent::InteractionRequested {
5433 request_id: Some("approval-stop".to_owned()),
5434 interaction_kind: ProviderInteractionKind::Approval,
5435 tool_name: "shell".to_owned(),
5436 title: None,
5437 prompt: "approve".to_owned(),
5438 options: Vec::new(),
5439 agent_id: None,
5440 },
5441 },
5442 });
5443 engine
5444 .apply_command(CommandEnvelope {
5445 id: CommandId(95),
5446 command: ControlCommand::ResolveInteraction {
5447 instance_id: instance(),
5448 generation: spawn.generation,
5449 interaction_id: ProviderInteractionId(1),
5450 response: ProviderInteractionResponse::ApproveOnce,
5451 },
5452 })
5453 .unwrap();
5454 engine
5455 .apply_command(CommandEnvelope {
5456 id: CommandId(96),
5457 command: ControlCommand::Stop {
5458 instance_id: instance(),
5459 force: false,
5460 },
5461 })
5462 .unwrap();
5463 let effects = engine.drain_effects();
5464 assert_eq!(effects.len(), 1);
5465 assert!(matches!(
5466 effects[0].effect,
5467 ControlEffect::Stop { force: false }
5468 ));
5469 engine.apply_observation(ObservationEnvelope {
5470 operation_id: Some(effects[0].operation_id),
5471 instance_id: effects[0].instance_id,
5472 generation: effects[0].generation,
5473 observation: ControlObservation::StopCompleted {
5474 forced: false,
5475 exit_code: Some(0),
5476 final_terminal: None,
5477 },
5478 });
5479 assert_eq!(
5480 engine.snapshot().sessions[0].provider.interactions[0].status,
5481 ProviderInteractionStatus::Resolved {
5482 outcome: ProviderInteractionOutcome::TurnEnded,
5483 }
5484 );
5485 }
5486
5487 #[test]
5488 fn generic_input_completion_never_resolves_provider_interactions() {
5489 let (mut engine, spawn) = running_engine();
5490 for (sequence, event) in [
5491 ProviderEvent::InteractionRequested {
5492 request_id: Some("approval-1".to_owned()),
5493 interaction_kind: ProviderInteractionKind::Approval,
5494 tool_name: "shell".to_owned(),
5495 title: None,
5496 prompt: "cargo test".to_owned(),
5497 options: Vec::new(),
5498 agent_id: None,
5499 },
5500 ProviderEvent::InteractionRequested {
5501 request_id: Some("question-1".to_owned()),
5502 interaction_kind: ProviderInteractionKind::Question,
5503 tool_name: "AskUserQuestion".to_owned(),
5504 title: None,
5505 prompt: "{\"question\":\"Continue?\"}".to_owned(),
5506 options: Vec::new(),
5507 agent_id: None,
5508 },
5509 ProviderEvent::InteractionRequested {
5510 request_id: Some("question-2".to_owned()),
5511 interaction_kind: ProviderInteractionKind::Question,
5512 tool_name: "AskUserQuestion".to_owned(),
5513 title: None,
5514 prompt: "{\"question\":\"Use the latest answer?\"}".to_owned(),
5515 options: Vec::new(),
5516 agent_id: None,
5517 },
5518 ]
5519 .into_iter()
5520 .enumerate()
5521 {
5522 engine.apply_observation(ObservationEnvelope {
5523 operation_id: None,
5524 instance_id: spawn.instance_id,
5525 generation: spawn.generation,
5526 observation: ControlObservation::ProviderEvent {
5527 source: provider_source(),
5528 sequence: u64::try_from(sequence + 1).unwrap(),
5529 event,
5530 },
5531 });
5532 }
5533 engine.apply_observation(ObservationEnvelope {
5534 operation_id: None,
5535 instance_id: spawn.instance_id,
5536 generation: spawn.generation,
5537 observation: ControlObservation::ProviderEvent {
5538 source: provider_source(),
5539 sequence: 4,
5540 event: ProviderEvent::Ready,
5541 },
5542 });
5543 assert_eq!(
5544 engine.snapshot().sessions[0].provider.activity,
5545 ProviderActivity::WaitingForInput
5546 );
5547
5548 engine
5549 .apply_command(CommandEnvelope {
5550 id: CommandId(82),
5551 command: ControlCommand::SendInput {
5552 instance_id: instance(),
5553 action: InputAction::TerminalControl(TerminalControl::Enter),
5554 },
5555 })
5556 .unwrap();
5557 let effect = engine.drain_effects().pop().unwrap();
5558 engine.apply_observation(ObservationEnvelope {
5559 operation_id: Some(effect.operation_id),
5560 instance_id: effect.instance_id,
5561 generation: effect.generation,
5562 observation: ControlObservation::InputCompleted,
5563 });
5564
5565 engine
5566 .apply_command(CommandEnvelope {
5567 id: CommandId(83),
5568 command: ControlCommand::SendInput {
5569 instance_id: instance(),
5570 action: InputAction::SubmitPrompt(PromptPayload {
5571 text: "continue".to_owned(),
5572 framing: PromptFraming::Literal,
5573 }),
5574 },
5575 })
5576 .unwrap();
5577 let effect = engine.drain_effects().pop().unwrap();
5578 engine.apply_observation(ObservationEnvelope {
5579 operation_id: Some(effect.operation_id),
5580 instance_id: effect.instance_id,
5581 generation: effect.generation,
5582 observation: ControlObservation::InputCompleted,
5583 });
5584
5585 let snapshot = engine.snapshot();
5586 let interactions = &snapshot.sessions[0].provider.interactions;
5587 assert_eq!(interactions[0].status, ProviderInteractionStatus::Pending);
5588 assert_eq!(interactions[1].status, ProviderInteractionStatus::Pending);
5589 assert_eq!(interactions[2].status, ProviderInteractionStatus::Pending);
5590 assert_eq!(
5591 snapshot.sessions[0].provider.activity,
5592 ProviderActivity::WaitingForInput
5593 );
5594 }
5595
5596 #[test]
5597 fn child_owned_interaction_restores_lead_state_without_adopting_child_activity() {
5598 let (mut engine, spawn) = running_engine();
5599 let source = provider_source();
5600 let events = [
5601 ProviderEvent::TurnCompleted {
5602 usage: TokenUsage::default(),
5603 is_cumulative: false,
5604 },
5605 ProviderEvent::SubagentStarted {
5606 agent_id: "child-1".to_owned(),
5607 agent_type: Some("reviewer".to_owned()),
5608 description: None,
5609 },
5610 ProviderEvent::InteractionRequested {
5611 request_id: Some("question-1".to_owned()),
5612 interaction_kind: ProviderInteractionKind::Question,
5613 tool_name: "AskUserQuestion".to_owned(),
5614 title: None,
5615 prompt: "{\"question\":\"Continue?\"}".to_owned(),
5616 options: Vec::new(),
5617 agent_id: Some("child-1".to_owned()),
5618 },
5619 ];
5620 for (index, event) in events.into_iter().enumerate() {
5621 engine.apply_observation(ObservationEnvelope {
5622 operation_id: None,
5623 instance_id: spawn.instance_id,
5624 generation: spawn.generation,
5625 observation: ControlObservation::ProviderEvent {
5626 source: source.clone(),
5627 sequence: u64::try_from(index + 1).unwrap(),
5628 event,
5629 },
5630 });
5631 }
5632 let provider = &engine.snapshot().sessions[0].provider;
5633 assert_eq!(provider.lead_activity, ProviderActivity::WaitingForInput);
5634 assert_eq!(
5635 provider.interactions[0].resume_lead_activity,
5636 Some(ProviderActivity::Idle)
5637 );
5638
5639 engine.apply_observation(ObservationEnvelope {
5640 operation_id: None,
5641 instance_id: spawn.instance_id,
5642 generation: spawn.generation,
5643 observation: ControlObservation::ProviderEvent {
5644 source: source.clone(),
5645 sequence: 4,
5646 event: ProviderEvent::InteractionRequested {
5647 request_id: Some("question-1".to_owned()),
5648 interaction_kind: ProviderInteractionKind::Question,
5649 tool_name: "AskUserQuestion".to_owned(),
5650 title: None,
5651 prompt: "{\"question\":\"Continue?\"}".to_owned(),
5652 options: Vec::new(),
5653 agent_id: Some("child-1".to_owned()),
5654 },
5655 },
5656 });
5657 let provider = &engine.snapshot().sessions[0].provider;
5658 assert_eq!(
5659 provider.interactions[0].status,
5660 ProviderInteractionStatus::Resolved {
5661 outcome: ProviderInteractionOutcome::Superseded
5662 }
5663 );
5664 assert_eq!(
5665 provider.interactions[1].resume_lead_activity,
5666 Some(ProviderActivity::Idle)
5667 );
5668
5669 engine.apply_observation(ObservationEnvelope {
5670 operation_id: None,
5671 instance_id: spawn.instance_id,
5672 generation: spawn.generation,
5673 observation: ControlObservation::ProviderEvent {
5674 source: source.clone(),
5675 sequence: 5,
5676 event: ProviderEvent::ToolStarted {
5677 id: "question-1".to_owned(),
5678 name: "AskUserQuestion".to_owned(),
5679 input_json: "{}".to_owned(),
5680 agent_id: Some("child-1".to_owned()),
5681 },
5682 },
5683 });
5684 let provider = &engine.snapshot().sessions[0].provider;
5685 assert_eq!(provider.lead_activity, ProviderActivity::Idle);
5686 assert_eq!(provider.activity, ProviderActivity::Working);
5687 assert!(provider.active_tools.is_empty());
5688 assert_eq!(
5689 provider.interactions[1].status,
5690 ProviderInteractionStatus::Resolved {
5691 outcome: ProviderInteractionOutcome::Answered
5692 }
5693 );
5694
5695 engine.apply_observation(ObservationEnvelope {
5696 operation_id: None,
5697 instance_id: spawn.instance_id,
5698 generation: spawn.generation,
5699 observation: ControlObservation::ProviderEvent {
5700 source: source.clone(),
5701 sequence: 6,
5702 event: ProviderEvent::InteractionRequested {
5703 request_id: Some("approval-2".to_owned()),
5704 interaction_kind: ProviderInteractionKind::Approval,
5705 tool_name: "shell".to_owned(),
5706 title: None,
5707 prompt: "approve".to_owned(),
5708 options: Vec::new(),
5709 agent_id: Some("child-1".to_owned()),
5710 },
5711 },
5712 });
5713
5714 engine.apply_observation(ObservationEnvelope {
5715 operation_id: None,
5716 instance_id: spawn.instance_id,
5717 generation: spawn.generation,
5718 observation: ControlObservation::ProviderEvent {
5719 source: source.clone(),
5720 sequence: 7,
5721 event: ProviderEvent::InteractionRequested {
5722 request_id: Some("question-3".to_owned()),
5723 interaction_kind: ProviderInteractionKind::Question,
5724 tool_name: "AskUserQuestion".to_owned(),
5725 title: None,
5726 prompt: "{\"question\":\"Another?\"}".to_owned(),
5727 options: Vec::new(),
5728 agent_id: Some("child-1".to_owned()),
5729 },
5730 },
5731 });
5732
5733 engine.apply_observation(ObservationEnvelope {
5734 operation_id: None,
5735 instance_id: spawn.instance_id,
5736 generation: spawn.generation,
5737 observation: ControlObservation::ProviderEvent {
5738 source,
5739 sequence: 8,
5740 event: ProviderEvent::SubagentStopped {
5741 agent_id: "child-1".to_owned(),
5742 },
5743 },
5744 });
5745 assert_eq!(
5746 engine.snapshot().sessions[0].provider.activity,
5747 ProviderActivity::Idle
5748 );
5749 assert_eq!(
5750 engine.snapshot().sessions[0].provider.interactions[2].status,
5751 ProviderInteractionStatus::Resolved {
5752 outcome: ProviderInteractionOutcome::TurnEnded
5753 }
5754 );
5755 assert_eq!(
5756 engine.snapshot().sessions[0].provider.interactions[3].status,
5757 ProviderInteractionStatus::Resolved {
5758 outcome: ProviderInteractionOutcome::TurnEnded
5759 }
5760 );
5761 }
5762
5763 #[test]
5764 fn provider_session_identity_observation_does_not_reset_live_turn_state() {
5765 let (mut engine, spawn) = running_engine();
5766 for (sequence, event) in [
5767 ProviderEvent::SubagentStarted {
5768 agent_id: "child-1".to_owned(),
5769 agent_type: None,
5770 description: None,
5771 },
5772 ProviderEvent::InteractionRequested {
5773 request_id: Some("approval-1".to_owned()),
5774 interaction_kind: ProviderInteractionKind::Approval,
5775 tool_name: "shell".to_owned(),
5776 title: None,
5777 prompt: "approve".to_owned(),
5778 options: Vec::new(),
5779 agent_id: Some("child-1".to_owned()),
5780 },
5781 ProviderEvent::SessionIdentityObserved {
5782 identity: ProviderSessionIdentity {
5783 key: ProviderSessionKey::SessionId,
5784 id: "provider-session-1".to_owned(),
5785 transcript_path: Some("C:/sessions/provider-session-1.jsonl".to_owned()),
5786 },
5787 },
5788 ]
5789 .into_iter()
5790 .enumerate()
5791 {
5792 engine.apply_observation(ObservationEnvelope {
5793 operation_id: None,
5794 instance_id: spawn.instance_id,
5795 generation: spawn.generation,
5796 observation: ControlObservation::ProviderEvent {
5797 source: provider_source(),
5798 sequence: u64::try_from(sequence + 1).unwrap(),
5799 event,
5800 },
5801 });
5802 }
5803 let provider = &engine.snapshot().sessions[0].provider;
5804 assert_eq!(
5805 provider.session,
5806 Some(ProviderSessionIdentity {
5807 key: ProviderSessionKey::SessionId,
5808 id: "provider-session-1".to_owned(),
5809 transcript_path: Some("C:/sessions/provider-session-1.jsonl".to_owned()),
5810 })
5811 );
5812 assert_eq!(provider.activity, ProviderActivity::WaitingForInput);
5813 assert_eq!(provider.subagents.len(), 1);
5814 assert_eq!(provider.interactions.len(), 1);
5815 assert_eq!(
5816 provider.interactions[0].status,
5817 ProviderInteractionStatus::Pending
5818 );
5819 }
5820
5821 #[test]
5822 fn working_observation_resolves_input_and_preserves_live_turn_context() {
5823 let (mut engine, spawn) = running_engine();
5824 for (sequence, event) in [
5825 ProviderEvent::TurnStarted {
5826 prompt: Some("ship the fix".to_owned()),
5827 },
5828 ProviderEvent::ToolStarted {
5829 id: "tool-1".to_owned(),
5830 name: "bash".to_owned(),
5831 input_json: "{\"command\":\"cargo check\"}".to_owned(),
5832 agent_id: None,
5833 },
5834 ProviderEvent::InteractionRequested {
5835 request_id: Some("approval-1".to_owned()),
5836 interaction_kind: ProviderInteractionKind::Approval,
5837 tool_name: "bash".to_owned(),
5838 title: None,
5839 prompt: "approve cargo check".to_owned(),
5840 options: Vec::new(),
5841 agent_id: None,
5842 },
5843 ProviderEvent::InteractionRequested {
5844 request_id: Some("question-1".to_owned()),
5845 interaction_kind: ProviderInteractionKind::Question,
5846 tool_name: "AskUserQuestion".to_owned(),
5847 title: None,
5848 prompt: "continue?".to_owned(),
5849 options: Vec::new(),
5850 agent_id: None,
5851 },
5852 ProviderEvent::WorkingObserved,
5853 ]
5854 .into_iter()
5855 .enumerate()
5856 {
5857 engine.apply_observation(ObservationEnvelope {
5858 operation_id: None,
5859 instance_id: spawn.instance_id,
5860 generation: spawn.generation,
5861 observation: ControlObservation::ProviderEvent {
5862 source: provider_source(),
5863 sequence: u64::try_from(sequence + 1).unwrap(),
5864 event,
5865 },
5866 });
5867 }
5868
5869 let provider = &engine.snapshot().sessions[0].provider;
5870 assert_eq!(provider.activity, ProviderActivity::Working);
5871 assert_eq!(provider.current_prompt.as_deref(), Some("ship the fix"));
5872 assert_eq!(provider.active_tools.len(), 1);
5873 assert_eq!(
5874 provider.interactions[0].status,
5875 ProviderInteractionStatus::Resolved {
5876 outcome: ProviderInteractionOutcome::Approved
5877 }
5878 );
5879 assert_eq!(
5880 provider.interactions[1].status,
5881 ProviderInteractionStatus::Resolved {
5882 outcome: ProviderInteractionOutcome::Answered
5883 }
5884 );
5885 }
5886
5887 #[test]
5888 fn provider_interruption_ends_the_turn_without_counting_completion() {
5889 let (mut engine, spawn) = running_engine();
5890 for (sequence, event) in [
5891 ProviderEvent::TurnStarted {
5892 prompt: Some("cancel me".to_owned()),
5893 },
5894 ProviderEvent::InteractionRequested {
5895 request_id: Some("approval-1".to_owned()),
5896 interaction_kind: ProviderInteractionKind::Approval,
5897 tool_name: "bash".to_owned(),
5898 title: None,
5899 prompt: "approve".to_owned(),
5900 options: Vec::new(),
5901 agent_id: None,
5902 },
5903 ProviderEvent::TurnInterrupted,
5904 ]
5905 .into_iter()
5906 .enumerate()
5907 {
5908 engine.apply_observation(ObservationEnvelope {
5909 operation_id: None,
5910 instance_id: spawn.instance_id,
5911 generation: spawn.generation,
5912 observation: ControlObservation::ProviderEvent {
5913 source: provider_source(),
5914 sequence: u64::try_from(sequence + 1).unwrap(),
5915 event,
5916 },
5917 });
5918 }
5919 let provider = &engine.snapshot().sessions[0].provider;
5920 assert_eq!(provider.activity, ProviderActivity::Idle);
5921 assert_eq!(provider.completed_turns, 0);
5922 assert_eq!(provider.current_prompt, None);
5923 assert_eq!(
5924 provider.interactions[0].status,
5925 ProviderInteractionStatus::Resolved {
5926 outcome: ProviderInteractionOutcome::Interrupted
5927 }
5928 );
5929 }
5930
5931 #[test]
5932 fn live_subagent_gates_lead_completion_until_exact_stop() {
5933 let (mut engine, spawn) = running_engine();
5934 for (sequence, event) in [
5935 ProviderEvent::SubagentStarted {
5936 agent_id: "child-1".to_owned(),
5937 agent_type: Some("reviewer".to_owned()),
5938 description: Some("review the reducer".to_owned()),
5939 },
5940 ProviderEvent::TurnCompleted {
5941 usage: TokenUsage::default(),
5942 is_cumulative: false,
5943 },
5944 ]
5945 .into_iter()
5946 .enumerate()
5947 {
5948 engine.apply_observation(ObservationEnvelope {
5949 operation_id: None,
5950 instance_id: spawn.instance_id,
5951 generation: spawn.generation,
5952 observation: ControlObservation::ProviderEvent {
5953 source: provider_source(),
5954 sequence: u64::try_from(sequence + 1).unwrap(),
5955 event,
5956 },
5957 });
5958 }
5959 let provider = &engine.snapshot().sessions[0].provider;
5960 assert_eq!(provider.lead_activity, ProviderActivity::Idle);
5961 assert_eq!(provider.activity, ProviderActivity::Working);
5962 assert_eq!(provider.subagents.len(), 1);
5963
5964 engine.apply_observation(ObservationEnvelope {
5965 operation_id: None,
5966 instance_id: spawn.instance_id,
5967 generation: spawn.generation,
5968 observation: ControlObservation::ProviderEvent {
5969 source: provider_source(),
5970 sequence: 3,
5971 event: ProviderEvent::SubagentStopped {
5972 agent_id: "child-1".to_owned(),
5973 },
5974 },
5975 });
5976 let provider = &engine.snapshot().sessions[0].provider;
5977 assert!(provider.subagents.is_empty());
5978 assert_eq!(provider.activity, ProviderActivity::Idle);
5979 }
5980
5981 #[test]
5982 fn pending_interaction_roster_is_bounded_with_explicit_supersession() {
5983 let (mut engine, spawn) = running_engine();
5984 engine.drain_events();
5985 for index in 0..=PROVIDER_INTERACTIONS_MAX {
5986 engine.apply_observation(ObservationEnvelope {
5987 operation_id: None,
5988 instance_id: spawn.instance_id,
5989 generation: spawn.generation,
5990 observation: ControlObservation::ProviderEvent {
5991 source: provider_source(),
5992 sequence: u64::try_from(index + 1).unwrap(),
5993 event: ProviderEvent::InteractionRequested {
5994 request_id: Some(format!("approval-{index}")),
5995 interaction_kind: ProviderInteractionKind::Approval,
5996 tool_name: "shell".to_owned(),
5997 title: None,
5998 prompt: String::new(),
5999 options: Vec::new(),
6000 agent_id: None,
6001 },
6002 },
6003 });
6004 }
6005
6006 let snapshot = engine.snapshot();
6007 assert_eq!(
6008 snapshot.sessions[0].provider.interactions.len(),
6009 PROVIDER_INTERACTIONS_MAX
6010 );
6011 assert_eq!(
6012 snapshot.sessions[0].provider.interactions[0].id,
6013 ProviderInteractionId(2)
6014 );
6015 assert!(engine.drain_events().iter().any(|event| matches!(
6016 event.event,
6017 ControlEventKind::InteractionResolved {
6018 interaction_id: ProviderInteractionId(1),
6019 outcome: ProviderInteractionOutcome::Superseded,
6020 }
6021 )));
6022 }
6023
6024 #[test]
6025 fn ingest_provider_jump_emits_exact_gap_and_accepts_next_sequence() {
6026 let (mut engine, spawn) = running_engine();
6027 engine.apply_observation(ObservationEnvelope {
6028 operation_id: None,
6029 instance_id: spawn.instance_id,
6030 generation: spawn.generation,
6031 observation: ControlObservation::ProviderEvent {
6032 source: provider_source(),
6033 sequence: 1,
6034 event: ProviderEvent::Ready,
6035 },
6036 });
6037
6038 engine
6039 .apply_command(CommandEnvelope {
6040 id: CommandId(30),
6041 command: ControlCommand::IngestProvider {
6042 instance_id: spawn.instance_id,
6043 generation: spawn.generation,
6044 source: hook_source(),
6045 source_sequence: 1,
6046 events: vec![
6047 ProviderEvent::TurnStarted {
6048 prompt: Some("first".to_owned()),
6049 },
6050 ProviderEvent::ToolStarted {
6051 id: "tool-1".to_owned(),
6052 name: "shell".to_owned(),
6053 input_json: "{}".to_owned(),
6054 agent_id: None,
6055 },
6056 ],
6057 },
6058 })
6059 .unwrap();
6060 engine
6061 .apply_command(CommandEnvelope {
6062 id: CommandId(31),
6063 command: ControlCommand::IngestProvider {
6064 instance_id: spawn.instance_id,
6065 generation: spawn.generation,
6066 source: hook_source(),
6067 source_sequence: 3,
6068 events: vec![ProviderEvent::TurnCompleted {
6069 usage: TokenUsage::default(),
6070 is_cumulative: false,
6071 }],
6072 },
6073 })
6074 .unwrap();
6075
6076 let provider = &engine.snapshot().sessions[0].provider;
6077 assert_eq!(provider.sequence, 5);
6078 assert_eq!(provider.sources.len(), 2);
6079 assert_eq!(provider.gap_count, 1);
6080 assert!(!provider.stale);
6081 assert_eq!(provider.completed_turns, 1);
6082 assert!(engine.drain_effects().is_empty());
6083 assert!(engine.drain_events().iter().any(|event| matches!(
6084 event.event,
6085 ControlEventKind::ProviderGap {
6086 sequence: 4,
6087 source_sequence: 2,
6088 missed: 1,
6089 ..
6090 }
6091 )));
6092
6093 engine
6094 .apply_command(CommandEnvelope {
6095 id: CommandId(32),
6096 command: ControlCommand::IngestProvider {
6097 instance_id: spawn.instance_id,
6098 generation: spawn.generation,
6099 source: hook_source(),
6100 source_sequence: 4,
6101 events: vec![ProviderEvent::WorkingObserved],
6102 },
6103 })
6104 .unwrap();
6105 assert_eq!(
6106 provider_source_sequence(
6107 &engine.snapshot().sessions[0].provider,
6108 &hook_source(),
6109 ),
6110 4,
6111 );
6112
6113 let stale = engine
6114 .apply_command(CommandEnvelope {
6115 id: CommandId(33),
6116 command: ControlCommand::IngestProvider {
6117 instance_id: spawn.instance_id,
6118 generation: spawn.generation,
6119 source: hook_source(),
6120 source_sequence: 4,
6121 events: vec![ProviderEvent::Ready],
6122 },
6123 })
6124 .unwrap_err();
6125 assert_eq!(stale, ControlError::StaleProviderSequence);
6126 }
6127
6128 #[test]
6129 fn external_ingress_rejects_stale_generation_empty_batches_and_oversized_events() {
6130 let (mut engine, spawn) = running_engine();
6131 let command = |id, generation, events| CommandEnvelope {
6132 id: CommandId(id),
6133 command: ControlCommand::IngestProvider {
6134 instance_id: spawn.instance_id,
6135 generation,
6136 source: hook_source(),
6137 source_sequence: 1,
6138 events,
6139 },
6140 };
6141
6142 assert!(matches!(
6143 engine.apply_command(command(
6144 40,
6145 SessionGeneration(0),
6146 vec![ProviderEvent::Ready]
6147 )),
6148 Err(ControlError::StaleProviderGeneration { .. })
6149 ));
6150 assert_eq!(
6151 engine.apply_command(command(41, spawn.generation, Vec::new())),
6152 Err(ControlError::InvalidProviderBatch {
6153 max: PROVIDER_INGRESS_EVENTS_MAX
6154 })
6155 );
6156 assert!(matches!(
6157 engine.apply_command(command(
6158 42,
6159 spawn.generation,
6160 vec![ProviderEvent::Text {
6161 text: "x".repeat(gate4agent_types::PROVIDER_EVENT_TEXT_MAX_BYTES + 1),
6162 is_delta: false,
6163 }]
6164 )),
6165 Err(ControlError::InvalidProviderEvent { .. })
6166 ));
6167 }
6168
6169 #[test]
6170 fn pipe_requires_initial_prompt_and_rejects_followup_input() {
6171 let mut engine = Gate4AgentEngine::new();
6172 let mut register = register(1);
6173 if let ControlCommand::Register { transport, .. } = &mut register.command {
6174 *transport = TransportKind::Pipe;
6175 }
6176 engine.apply_command(register).unwrap();
6177 assert_eq!(
6178 engine.apply_command(start(2)).unwrap_err(),
6179 ControlError::MissingInitialPrompt
6180 );
6181
6182 let mut request = start(3);
6183 if let ControlCommand::Start { request, .. } = &mut request.command {
6184 request.initial_prompt = Some("hello".to_owned());
6185 }
6186 engine.apply_command(request).unwrap();
6187 let spawn = engine.drain_effects().pop().unwrap();
6188 engine.apply_observation(ObservationEnvelope {
6189 operation_id: Some(spawn.operation_id),
6190 instance_id: spawn.instance_id,
6191 generation: spawn.generation,
6192 observation: ControlObservation::Spawned {
6193 process_id: Some(1),
6194 },
6195 });
6196 let error = engine
6197 .apply_command(CommandEnvelope {
6198 id: CommandId(4),
6199 command: ControlCommand::SendInput {
6200 instance_id: instance(),
6201 action: InputAction::SubmitPrompt(PromptPayload {
6202 text: "again".to_owned(),
6203 framing: PromptFraming::Literal,
6204 }),
6205 },
6206 })
6207 .unwrap_err();
6208 assert!(matches!(
6209 error,
6210 ControlError::UnsupportedTransportOperation {
6211 transport: TransportKind::Pipe,
6212 ..
6213 }
6214 ));
6215 }
6216
6217 #[test]
6218 fn matching_spawn_observation_confirms_running() {
6219 let (engine, _) = running_engine();
6220 let session = &engine.snapshot().sessions[0];
6221 assert_eq!(session.status, SessionStatus::Running);
6222 assert_eq!(session.process_id, Some(42));
6223 assert_eq!(session.pending_operation, None);
6224 assert_eq!(
6225 session.terminal_size,
6226 Some(TerminalSize {
6227 rows: 24,
6228 columns: 80,
6229 })
6230 );
6231 }
6232
6233 #[test]
6234 fn resize_is_pending_until_observed() {
6235 let (mut engine, _) = running_engine();
6236 let size = TerminalSize {
6237 rows: 40,
6238 columns: 120,
6239 };
6240 engine
6241 .apply_command(CommandEnvelope {
6242 id: CommandId(9),
6243 command: ControlCommand::Resize {
6244 instance_id: instance(),
6245 size,
6246 },
6247 })
6248 .unwrap();
6249 let effect = engine.drain_effects().pop().unwrap();
6250 assert_eq!(
6251 engine.snapshot().sessions[0].terminal_size,
6252 Some(TerminalSize {
6253 rows: 24,
6254 columns: 80,
6255 })
6256 );
6257
6258 engine.apply_observation(ObservationEnvelope {
6259 operation_id: Some(effect.operation_id),
6260 instance_id: effect.instance_id,
6261 generation: effect.generation,
6262 observation: ControlObservation::ResizeCompleted { size },
6263 });
6264 assert_eq!(engine.snapshot().sessions[0].terminal_size, Some(size));
6265 }
6266
6267 #[test]
6268 fn foreground_refresh_is_generation_bound_replaceable_authority() {
6269 let (mut engine, _) = running_engine();
6270 engine.drain_events();
6271 engine
6272 .apply_command(CommandEnvelope {
6273 id: CommandId(10),
6274 command: ControlCommand::RefreshForeground {
6275 instance_id: instance(),
6276 },
6277 })
6278 .unwrap();
6279 let effect = engine.drain_effects().pop().unwrap();
6280 assert_eq!(effect.effect, ControlEffect::ObserveForeground);
6281 assert_eq!(
6282 engine.snapshot().sessions[0].foreground.authority,
6283 ForegroundAuthority::Stale
6284 );
6285
6286 let process = ForegroundProcess {
6287 root_process_id: 42,
6288 process_id: 84,
6289 process_name: "claude".to_owned(),
6290 kind: ForegroundProcessKind::Agent {
6291 agent_id: AgentId::new("claude").unwrap(),
6292 },
6293 };
6294 engine.apply_observation(ObservationEnvelope {
6295 operation_id: Some(effect.operation_id),
6296 instance_id: effect.instance_id,
6297 generation: effect.generation,
6298 observation: ControlObservation::ForegroundObserved {
6299 process: process.clone(),
6300 },
6301 });
6302
6303 let session = &engine.snapshot().sessions[0];
6304 assert_eq!(session.pending_operation, None);
6305 assert_eq!(session.foreground.authority, ForegroundAuthority::Confirmed);
6306 assert_eq!(session.foreground.process.as_ref(), Some(&process));
6307 assert_eq!(session.foreground.stale_reason, None);
6308 assert!(engine.drain_events().iter().any(|event| matches!(
6309 &event.event,
6310 ControlEventKind::ForegroundObserved { process: observed } if observed == &process
6311 )));
6312 }
6313
6314 #[test]
6315 fn foreground_refresh_rejects_another_root_process() {
6316 let (mut engine, _) = running_engine();
6317 engine.drain_events();
6318 engine
6319 .apply_command(CommandEnvelope {
6320 id: CommandId(11),
6321 command: ControlCommand::RefreshForeground {
6322 instance_id: instance(),
6323 },
6324 })
6325 .unwrap();
6326 let effect = engine.drain_effects().pop().unwrap();
6327 engine.apply_observation(ObservationEnvelope {
6328 operation_id: Some(effect.operation_id),
6329 instance_id: effect.instance_id,
6330 generation: effect.generation,
6331 observation: ControlObservation::ForegroundObserved {
6332 process: ForegroundProcess {
6333 root_process_id: 999,
6334 process_id: 1_000,
6335 process_name: "claude".to_owned(),
6336 kind: ForegroundProcessKind::Agent {
6337 agent_id: AgentId::new("claude").unwrap(),
6338 },
6339 },
6340 },
6341 });
6342
6343 assert_eq!(
6344 engine.snapshot().sessions[0].foreground.authority,
6345 ForegroundAuthority::Stale
6346 );
6347 assert!(engine.drain_events().iter().any(|event| matches!(
6348 event.event,
6349 ControlEventKind::ObservationIgnored {
6350 reason: ObservationIgnoredReason::InvalidForegroundObservation
6351 }
6352 )));
6353 }
6354
6355 #[test]
6356 fn terminal_frames_are_replaceable_and_never_move_backwards() {
6357 let (mut engine, spawn) = running_engine();
6358 engine.drain_events();
6359 let frame = TerminalFrame {
6360 sequence: 10,
6361 size: TerminalSize {
6362 rows: 24,
6363 columns: 80,
6364 },
6365 cursor_row: 1,
6366 cursor_column: 2,
6367 contents: "current".to_owned(),
6368 formatted: b"current".to_vec(),
6369 scrollback_formatted: Vec::new(),
6370 alternate_screen: false,
6371 mouse_protocol_enabled: false,
6372 mouse_protocol_encoding: Default::default(),
6373 produced_at_unix_ms: 0,
6374 screen_state: PtyScreenState::default(),
6375 bracketed_paste: None,
6376 };
6377 engine.apply_observation(ObservationEnvelope {
6378 operation_id: None,
6379 instance_id: instance(),
6380 generation: spawn.generation,
6381 observation: ControlObservation::TerminalFrame {
6382 frame: frame.clone(),
6383 },
6384 });
6385 assert_eq!(engine.snapshot().sessions[0].terminal_frame, Some(frame));
6386 assert!(engine.drain_events().is_empty());
6387
6388 engine.apply_observation(ObservationEnvelope {
6389 operation_id: None,
6390 instance_id: instance(),
6391 generation: spawn.generation,
6392 observation: ControlObservation::TerminalFrame {
6393 frame: TerminalFrame {
6394 sequence: 9,
6395 size: TerminalSize {
6396 rows: 24,
6397 columns: 80,
6398 },
6399 cursor_row: 0,
6400 cursor_column: 0,
6401 contents: "stale".to_owned(),
6402 formatted: Vec::new(),
6403 scrollback_formatted: Vec::new(),
6404 alternate_screen: false,
6405 mouse_protocol_enabled: false,
6406 mouse_protocol_encoding: Default::default(),
6407 produced_at_unix_ms: 0,
6408 screen_state: PtyScreenState::default(),
6409 bracketed_paste: None,
6410 },
6411 },
6412 });
6413 assert_eq!(
6414 engine.snapshot().sessions[0]
6415 .terminal_frame
6416 .as_ref()
6417 .unwrap()
6418 .sequence,
6419 10
6420 );
6421 assert!(engine.drain_events().iter().any(|event| matches!(
6422 event.event,
6423 ControlEventKind::ObservationIgnored {
6424 reason: ObservationIgnoredReason::StaleTerminalFrame
6425 }
6426 )));
6427 }
6428
6429 #[test]
6430 fn stop_remains_pending_until_observed() {
6431 let (mut engine, _) = running_engine();
6432 engine.drain_events();
6433 engine
6434 .apply_command(CommandEnvelope {
6435 id: CommandId(3),
6436 command: ControlCommand::Stop {
6437 instance_id: instance(),
6438 force: true,
6439 },
6440 })
6441 .unwrap();
6442 let effect = engine.drain_effects().pop().unwrap();
6443 assert_eq!(
6444 engine.snapshot().sessions[0].status,
6445 SessionStatus::Stopping
6446 );
6447
6448 engine.apply_observation(ObservationEnvelope {
6449 operation_id: Some(effect.operation_id),
6450 instance_id: effect.instance_id,
6451 generation: effect.generation,
6452 observation: ControlObservation::StopCompleted {
6453 forced: true,
6454 exit_code: None,
6455 final_terminal: None,
6456 },
6457 });
6458 assert_eq!(
6459 engine.snapshot().sessions[0].status,
6460 SessionStatus::Exited { exit_code: None }
6461 );
6462 }
6463
6464 #[test]
6465 fn typed_input_is_an_effect_until_executor_confirms_it() {
6466 let (mut engine, _) = running_engine();
6467 engine.drain_events();
6468 engine
6469 .apply_command(CommandEnvelope {
6470 id: CommandId(7),
6471 command: ControlCommand::SendInput {
6472 instance_id: instance(),
6473 action: InputAction::SubmitPrompt(PromptPayload {
6474 text: "inspect the lifecycle".to_owned(),
6475 framing: PromptFraming::BracketedPaste,
6476 }),
6477 },
6478 })
6479 .unwrap();
6480 let effect = engine.drain_effects().pop().unwrap();
6481 assert!(matches!(
6482 &effect.effect,
6483 ControlEffect::WriteInput {
6484 required_foreground: ForegroundRequirement::Agent { agent_id },
6485 ..
6486 } if agent_id.as_str() == "claude"
6487 ));
6488 assert_eq!(
6489 engine.snapshot().sessions[0].pending_input,
6490 Some(PreparedInputKind::SubmitPrompt)
6491 );
6492 assert_eq!(
6493 engine.snapshot().sessions[0].foreground.authority,
6494 ForegroundAuthority::Stale
6495 );
6496
6497 engine.apply_observation(ObservationEnvelope {
6498 operation_id: Some(effect.operation_id),
6499 instance_id: effect.instance_id,
6500 generation: effect.generation,
6501 observation: ControlObservation::InputCompleted,
6502 });
6503 let session = &engine.snapshot().sessions[0];
6504 assert_eq!(session.status, SessionStatus::Running);
6505 assert_eq!(session.pending_operation, None);
6506 assert_eq!(session.pending_input, None);
6507 assert!(engine.drain_events().iter().any(|event| matches!(
6508 event.event,
6509 ControlEventKind::InputCompleted {
6510 input_kind: PreparedInputKind::SubmitPrompt
6511 }
6512 )));
6513 }
6514
6515 #[test]
6516 fn shell_command_effect_requires_fresh_shell_routing() {
6517 let (mut engine, _) = running_engine();
6518 engine
6519 .apply_command(CommandEnvelope {
6520 id: CommandId(70),
6521 command: ControlCommand::SendInput {
6522 instance_id: instance(),
6523 action: InputAction::ShellCommand(ShellCommand {
6524 text: "git status --short".to_owned(),
6525 }),
6526 },
6527 })
6528 .unwrap();
6529
6530 let effect = engine.drain_effects().pop().unwrap();
6531 assert!(matches!(
6532 effect.effect,
6533 ControlEffect::WriteInput {
6534 input,
6535 required_foreground: ForegroundRequirement::Shell,
6536 } if input.kind() == PreparedInputKind::ShellCommand
6537 ));
6538 assert_eq!(
6539 engine.snapshot().sessions[0].pending_input,
6540 Some(PreparedInputKind::ShellCommand)
6541 );
6542 }
6543
6544 #[test]
6545 fn agent_command_effect_is_bound_to_the_session_agent_route() {
6546 let (mut engine, _) = running_engine();
6547 engine
6548 .apply_command(CommandEnvelope {
6549 id: CommandId(71),
6550 command: ControlCommand::SendInput {
6551 instance_id: instance(),
6552 action: InputAction::AgentCommand(gate4agent_types::AgentCommand {
6553 agent_id: AgentId::new("claude").unwrap(),
6554 name: "review".to_owned(),
6555 arguments: vec!["routing".to_owned()],
6556 }),
6557 },
6558 })
6559 .unwrap();
6560
6561 let effect = engine.drain_effects().pop().unwrap();
6562 assert!(matches!(
6563 effect.effect,
6564 ControlEffect::WriteInput {
6565 input,
6566 required_foreground: ForegroundRequirement::Agent { agent_id },
6567 } if input.kind() == PreparedInputKind::AgentCommand && agent_id.as_str() == "claude"
6568 ));
6569 }
6570
6571 #[test]
6572 fn provider_command_cannot_target_another_running_agent() {
6573 let (mut engine, _) = running_engine();
6574 let error = engine
6575 .apply_command(CommandEnvelope {
6576 id: CommandId(8),
6577 command: ControlCommand::SendInput {
6578 instance_id: instance(),
6579 action: InputAction::AgentCommand(gate4agent_types::AgentCommand {
6580 agent_id: AgentId::new("kimi").unwrap(),
6581 name: "help".to_owned(),
6582 arguments: Vec::new(),
6583 }),
6584 },
6585 })
6586 .unwrap_err();
6587
6588 assert!(matches!(error, ControlError::InputRejected { .. }));
6589 assert!(engine.drain_effects().is_empty());
6590 assert_eq!(engine.snapshot().sessions[0].pending_operation, None);
6591 }
6592
6593 #[test]
6594 fn remove_and_reregister_strictly_advance_the_generation_watermark() {
6595 let (mut engine, first_spawn) = running_engine();
6596 engine.apply_observation(ObservationEnvelope {
6597 operation_id: None,
6598 instance_id: instance(),
6599 generation: first_spawn.generation,
6600 observation: ControlObservation::ProcessExited {
6601 exit_code: Some(0),
6602 final_terminal: None,
6603 },
6604 });
6605 engine
6606 .apply_command(CommandEnvelope {
6607 id: CommandId(3),
6608 command: ControlCommand::Remove {
6609 instance_id: instance(),
6610 },
6611 })
6612 .unwrap();
6613 engine.apply_command(register(4)).unwrap();
6614
6615 let generation = engine.snapshot().sessions[0].generation;
6616 assert_eq!(generation.0, first_spawn.generation.0 + 1);
6617 assert_eq!(
6618 engine.generation_watermarks.get(&instance()),
6619 Some(&generation)
6620 );
6621 }
6622
6623 #[test]
6624 fn reregister_generation_exhaustion_is_atomic() {
6625 let mut engine = Gate4AgentEngine::new();
6626 let exhausted = SessionGeneration(u64::MAX);
6627 engine.generation_watermarks.insert(instance(), exhausted);
6628 let before = engine.clone();
6629
6630 let error = engine.apply_command(register(1)).unwrap_err();
6631
6632 assert_eq!(
6633 error,
6634 ControlError::GenerationExhausted {
6635 instance_id: instance(),
6636 generation: exhausted,
6637 }
6638 );
6639 assert_eq!(engine, before);
6640 }
6641
6642 #[test]
6643 fn live_session_capacity_rejection_is_atomic() {
6644 let mut engine = Gate4AgentEngine::new();
6645 let first_instance = 10_000_u64;
6646 for offset in 0..CONTROL_SESSIONS_MAX {
6647 engine
6648 .apply_command(register_instance(
6649 offset as u64 + 1,
6650 AgentInstanceId(first_instance + offset as u64),
6651 ))
6652 .unwrap();
6653 }
6654 engine.drain_events();
6655 let before = engine.clone();
6656 let rejected_instance = AgentInstanceId(first_instance + CONTROL_SESSIONS_MAX as u64);
6657
6658 let error = engine
6659 .apply_command(register_instance(10_000, rejected_instance))
6660 .unwrap_err();
6661
6662 assert_eq!(
6663 error,
6664 ControlError::SessionCapacityExceeded {
6665 instance_id: rejected_instance,
6666 max: CONTROL_SESSIONS_MAX,
6667 }
6668 );
6669 assert_eq!(engine, before);
6670 }
6671
6672 #[test]
6673 fn retained_identity_capacity_is_bounded_without_weakening_stale_fence() {
6674 let mut engine = Gate4AgentEngine::new();
6675 let first_instance = 20_000_u64;
6676 for offset in 0..CONTROL_INSTANCE_IDENTITIES_MAX {
6677 engine.generation_watermarks.insert(
6678 AgentInstanceId(first_instance + offset as u64),
6679 SessionGeneration::default(),
6680 );
6681 }
6682 let health = engine.snapshot().health;
6683 assert_eq!(
6684 health.retained_instance_identities,
6685 u32::try_from(CONTROL_INSTANCE_IDENTITIES_MAX).unwrap()
6686 );
6687 assert_eq!(
6688 health.retained_instance_identity_capacity,
6689 u32::try_from(CONTROL_INSTANCE_IDENTITIES_MAX).unwrap()
6690 );
6691 let known_instance = AgentInstanceId(first_instance);
6692 let rejected_instance =
6693 AgentInstanceId(first_instance + CONTROL_INSTANCE_IDENTITIES_MAX as u64);
6694 let before = engine.clone();
6695
6696 let error = engine
6697 .apply_command(register_instance(1, rejected_instance))
6698 .unwrap_err();
6699
6700 assert_eq!(
6701 error,
6702 ControlError::InstanceIdentityCapacityExceeded {
6703 instance_id: rejected_instance,
6704 max: CONTROL_INSTANCE_IDENTITIES_MAX,
6705 }
6706 );
6707 assert_eq!(engine, before);
6708
6709 engine
6710 .apply_command(register_instance(2, known_instance))
6711 .unwrap();
6712 let current = engine
6713 .session_snapshot(known_instance)
6714 .expect("known retained identity must remain reusable")
6715 .clone();
6716 assert_eq!(current.generation, SessionGeneration(1));
6717 engine.drain_events();
6718
6719 engine.apply_observation(ObservationEnvelope {
6720 operation_id: None,
6721 instance_id: known_instance,
6722 generation: SessionGeneration::default(),
6723 observation: ControlObservation::ProcessExited {
6724 exit_code: Some(9),
6725 final_terminal: None,
6726 },
6727 });
6728
6729 assert_eq!(engine.session_snapshot(known_instance), Some(¤t));
6730 assert!(engine.drain_events().iter().any(|event| matches!(
6731 event.event,
6732 ControlEventKind::ObservationIgnored {
6733 reason: ObservationIgnoredReason::StaleGeneration
6734 }
6735 )));
6736 }
6737
6738 #[test]
6739 fn operation_id_max_is_issued_once_and_later_commands_fail_atomically() {
6740 let (mut engine, spawn) = running_engine();
6741 engine.drain_events();
6742 engine.next_operation_id = Some(u64::MAX);
6743 engine.apply_command(terminal_text(3, "first")).unwrap();
6744 let effect = engine.drain_effects().pop().unwrap();
6745 assert_eq!(effect.operation_id, OperationId(u64::MAX));
6746 engine
6747 .try_apply_observation(ObservationEnvelope {
6748 operation_id: Some(effect.operation_id),
6749 instance_id: instance(),
6750 generation: spawn.generation,
6751 observation: ControlObservation::InputCompleted,
6752 })
6753 .unwrap();
6754 let before = engine.clone();
6755
6756 let error = engine
6757 .apply_command(terminal_text(4, "second"))
6758 .unwrap_err();
6759
6760 assert_eq!(error, ControlError::OperationIdExhausted);
6761 assert_eq!(engine, before);
6762 assert!(engine.snapshot().health.operation_id_exhausted);
6763 }
6764
6765 #[test]
6766 fn event_sequence_max_is_emitted_once_and_observation_failure_is_atomic() {
6767 let (mut engine, spawn) = running_engine();
6768 engine.drain_events();
6769 engine.next_event_sequence = Some(u64::MAX);
6770 engine.apply_command(terminal_text(3, "first")).unwrap();
6771 let effect = engine.drain_effects().pop().unwrap();
6772 let before = engine.clone();
6773
6774 let error = engine
6775 .try_apply_observation(ObservationEnvelope {
6776 operation_id: Some(effect.operation_id),
6777 instance_id: instance(),
6778 generation: spawn.generation,
6779 observation: ControlObservation::InputCompleted,
6780 })
6781 .unwrap_err();
6782
6783 assert_eq!(error, ControlError::EventSequenceExhausted);
6784 assert_eq!(engine, before);
6785 assert!(engine.snapshot().health.event_sequence_exhausted);
6786 let events = engine.drain_events();
6787 assert_eq!(events.len(), 1);
6788 assert_eq!(events[0].sequence, u64::MAX);
6789 }
6790
6791 #[test]
6792 fn revision_max_is_published_once_and_later_observation_is_atomic() {
6793 let (mut engine, spawn) = running_engine();
6794 engine.drain_events();
6795 engine.revision = u64::MAX - 1;
6796 engine.apply_command(terminal_text(3, "first")).unwrap();
6797 let effect = engine.drain_effects().pop().unwrap();
6798 assert_eq!(engine.snapshot().revision, u64::MAX);
6799 let before = engine.clone();
6800
6801 let error = engine
6802 .try_apply_observation(ObservationEnvelope {
6803 operation_id: Some(effect.operation_id),
6804 instance_id: instance(),
6805 generation: spawn.generation,
6806 observation: ControlObservation::InputCompleted,
6807 })
6808 .unwrap_err();
6809
6810 assert_eq!(error, ControlError::RevisionExhausted);
6811 assert_eq!(engine, before);
6812 assert!(engine.snapshot().health.revision_exhausted);
6813 }
6814
6815 #[test]
6816 fn provider_sequence_capacity_is_preflighted_for_the_whole_ingress_batch() {
6817 let (mut engine, spawn) = running_engine();
6818 engine.session_mut(instance()).provider.sequence = u64::MAX - 1;
6819 let before = engine.clone();
6820
6821 let error = engine
6822 .apply_command(CommandEnvelope {
6823 id: CommandId(3),
6824 command: ControlCommand::IngestProvider {
6825 instance_id: instance(),
6826 generation: spawn.generation,
6827 source: provider_source(),
6828 source_sequence: 2,
6829 events: vec![ProviderEvent::Text {
6830 text: "bounded".to_owned(),
6831 is_delta: false,
6832 }],
6833 },
6834 })
6835 .unwrap_err();
6836
6837 assert_eq!(
6838 error,
6839 ControlError::ProviderSequenceExhausted {
6840 instance_id: instance(),
6841 generation: spawn.generation,
6842 }
6843 );
6844 assert_eq!(engine, before);
6845 }
6846
6847 #[test]
6848 fn provider_sequence_terminal_state_is_observable_and_blocks_observations() {
6849 let (mut engine, spawn) = running_engine();
6850 engine.session_mut(instance()).provider.sequence = u64::MAX;
6851 let before = engine.clone();
6852 assert_eq!(
6853 engine
6854 .snapshot()
6855 .health
6856 .provider_sequence_exhausted_sessions,
6857 1
6858 );
6859
6860 let error = engine
6861 .try_apply_observation(ObservationEnvelope {
6862 operation_id: None,
6863 instance_id: instance(),
6864 generation: spawn.generation,
6865 observation: ControlObservation::ProviderGap {
6866 source: provider_source(),
6867 source_sequence: 1,
6868 missed: 1,
6869 },
6870 })
6871 .unwrap_err();
6872
6873 assert_eq!(
6874 error,
6875 ControlError::ProviderSequenceExhausted {
6876 instance_id: instance(),
6877 generation: spawn.generation,
6878 }
6879 );
6880 assert_eq!(engine, before);
6881 }
6882
6883 #[test]
6884 fn provider_source_sequence_exhaustion_is_typed_and_atomic() {
6885 let (mut engine, spawn) = running_engine();
6886 let source = provider_source();
6887 engine
6888 .session_mut(instance())
6889 .provider
6890 .sources
6891 .push(ProviderSourceCursor {
6892 source: source.clone(),
6893 sequence: u64::MAX,
6894 gap_count: 0,
6895 stale: false,
6896 });
6897 let before = engine.clone();
6898
6899 let error = engine
6900 .apply_command(CommandEnvelope {
6901 id: CommandId(3),
6902 command: ControlCommand::IngestProvider {
6903 instance_id: instance(),
6904 generation: spawn.generation,
6905 source: source.clone(),
6906 source_sequence: u64::MAX,
6907 events: vec![ProviderEvent::WorkingObserved],
6908 },
6909 })
6910 .unwrap_err();
6911
6912 assert_eq!(
6913 error,
6914 ControlError::ProviderSourceSequenceExhausted {
6915 instance_id: instance(),
6916 generation: spawn.generation,
6917 provider_source: source,
6918 }
6919 );
6920 assert_eq!(engine, before);
6921 }
6922
6923 #[test]
6924 fn direct_provider_gap_source_sequence_overflow_is_typed_and_atomic() {
6925 let (mut engine, spawn) = running_engine();
6926 let source = provider_source();
6927 engine
6928 .session_mut(instance())
6929 .provider
6930 .sources
6931 .push(ProviderSourceCursor {
6932 source: source.clone(),
6933 sequence: u64::MAX,
6934 gap_count: 0,
6935 stale: false,
6936 });
6937 let before = engine.clone();
6938 let error = engine
6939 .try_apply_observation(ObservationEnvelope {
6940 operation_id: None,
6941 instance_id: instance(),
6942 generation: spawn.generation,
6943 observation: ControlObservation::ProviderGap {
6944 source: source.clone(),
6945 source_sequence: u64::MAX,
6946 missed: 1,
6947 },
6948 })
6949 .unwrap_err();
6950 assert_eq!(
6951 error,
6952 ControlError::ProviderSourceSequenceExhausted {
6953 instance_id: instance(),
6954 generation: spawn.generation,
6955 provider_source: source,
6956 },
6957 );
6958 assert_eq!(engine, before);
6959 }
6960
6961 #[test]
6962 fn direct_provider_gap_rejects_non_exact_source_sequence() {
6963 let (mut engine, spawn) = running_engine();
6964 engine.apply_observation(ObservationEnvelope {
6965 operation_id: None,
6966 instance_id: instance(),
6967 generation: spawn.generation,
6968 observation: ControlObservation::ProviderGap {
6969 source: provider_source(),
6970 source_sequence: 2,
6971 missed: 1,
6972 },
6973 });
6974 let provider = &engine.snapshot().sessions[0].provider;
6975 assert_eq!(provider.gap_count, 0);
6976 assert!(provider.sources.is_empty());
6977 assert!(engine.drain_events().iter().any(|event| matches!(
6978 event.event,
6979 ControlEventKind::ObservationIgnored {
6980 reason: ObservationIgnoredReason::StaleProviderEvent,
6981 }
6982 )));
6983 }
6984
6985 #[test]
6986 fn removed_lifecycle_observations_cannot_mutate_a_reregistered_instance() {
6987 let (mut engine, first_spawn) = running_engine();
6988 engine.apply_observation(ObservationEnvelope {
6989 operation_id: None,
6990 instance_id: instance(),
6991 generation: first_spawn.generation,
6992 observation: ControlObservation::ProcessExited {
6993 exit_code: Some(0),
6994 final_terminal: None,
6995 },
6996 });
6997 engine
6998 .apply_command(CommandEnvelope {
6999 id: CommandId(3),
7000 command: ControlCommand::Remove {
7001 instance_id: instance(),
7002 },
7003 })
7004 .unwrap();
7005 engine.apply_command(register(4)).unwrap();
7006 engine.apply_command(start(5)).unwrap();
7007 let second_spawn = engine.drain_effects().pop().unwrap();
7008 engine.apply_observation(ObservationEnvelope {
7009 operation_id: Some(second_spawn.operation_id),
7010 instance_id: instance(),
7011 generation: second_spawn.generation,
7012 observation: ControlObservation::Spawned {
7013 process_id: Some(84),
7014 },
7015 });
7016 let before = engine.snapshot().sessions[0].clone();
7017 engine.drain_events();
7018
7019 engine.apply_observation(ObservationEnvelope {
7020 operation_id: None,
7021 instance_id: instance(),
7022 generation: first_spawn.generation,
7023 observation: ControlObservation::ProcessExited {
7024 exit_code: Some(9),
7025 final_terminal: None,
7026 },
7027 });
7028 engine.apply_observation(ObservationEnvelope {
7029 operation_id: None,
7030 instance_id: instance(),
7031 generation: first_spawn.generation,
7032 observation: ControlObservation::ProviderEvent {
7033 source: provider_source(),
7034 sequence: 1,
7035 event: ProviderEvent::SessionIdentityObserved {
7036 identity: ProviderSessionIdentity {
7037 key: ProviderSessionKey::SessionId,
7038 id: "stale-provider-session".to_owned(),
7039 transcript_path: None,
7040 },
7041 },
7042 },
7043 });
7044
7045 assert_eq!(engine.snapshot().sessions[0], before);
7046 let ignored = engine
7047 .drain_events()
7048 .into_iter()
7049 .filter(|event| {
7050 matches!(
7051 event.event,
7052 ControlEventKind::ObservationIgnored {
7053 reason: ObservationIgnoredReason::StaleGeneration
7054 }
7055 )
7056 })
7057 .count();
7058 assert_eq!(ignored, 2);
7059 }
7060
7061 #[test]
7062 fn start_generation_exhaustion_is_atomic() {
7063 let mut engine = Gate4AgentEngine::new();
7064 engine.apply_command(register(1)).unwrap();
7065 let exhausted = SessionGeneration(u64::MAX);
7066 engine.session_mut(instance()).generation = exhausted;
7067 engine.generation_watermarks.insert(instance(), exhausted);
7068 let before = engine.clone();
7069
7070 let error = engine.apply_command(start(2)).unwrap_err();
7071
7072 assert_eq!(
7073 error,
7074 ControlError::GenerationExhausted {
7075 instance_id: instance(),
7076 generation: exhausted,
7077 }
7078 );
7079 assert_eq!(engine, before);
7080 }
7081
7082 #[test]
7083 fn start_generation_exhaustion_preserves_queued_history_effect() {
7084 let mut engine = Gate4AgentEngine::new();
7085 engine.apply_command(register(1)).unwrap();
7086 engine
7087 .apply_command(CommandEnvelope {
7088 id: CommandId(2),
7089 command: ControlCommand::DiscoverHistory {
7090 instance_id: instance(),
7091 query: HistoryQuery {
7092 working_directory: None,
7093 limit: 1,
7094 },
7095 },
7096 })
7097 .unwrap();
7098 let exhausted = SessionGeneration(u64::MAX);
7099 engine.session_mut(instance()).generation = exhausted;
7100 engine.generation_watermarks.insert(instance(), exhausted);
7101 let before = engine.clone();
7102
7103 let error = engine.apply_command(start(3)).unwrap_err();
7104
7105 assert_eq!(
7106 error,
7107 ControlError::GenerationExhausted {
7108 instance_id: instance(),
7109 generation: exhausted,
7110 }
7111 );
7112 assert_eq!(engine, before);
7113 assert!(matches!(
7114 engine.drain_effects().as_slice(),
7115 [EffectEnvelope {
7116 effect: ControlEffect::DiscoverHistory { .. },
7117 ..
7118 }]
7119 ));
7120 }
7121
7122 #[test]
7123 fn multi_event_rollback_retires_the_event_sequence_terminally() {
7124 let (mut engine, spawn) = running_engine();
7125 engine.apply_observation(ObservationEnvelope {
7126 operation_id: None,
7127 instance_id: instance(),
7128 generation: spawn.generation,
7129 observation: ControlObservation::ProviderEvent {
7130 source: provider_source(),
7131 sequence: 1,
7132 event: ProviderEvent::InteractionRequested {
7133 request_id: Some("approval-counter-boundary".to_owned()),
7134 interaction_kind: ProviderInteractionKind::Approval,
7135 tool_name: "shell".to_owned(),
7136 title: None,
7137 prompt: "approve".to_owned(),
7138 options: Vec::new(),
7139 agent_id: None,
7140 },
7141 },
7142 });
7143 engine
7144 .apply_command(CommandEnvelope {
7145 id: CommandId(3),
7146 command: ControlCommand::SendInput {
7147 instance_id: instance(),
7148 action: InputAction::TerminalControl(TerminalControl::Interrupt),
7149 },
7150 })
7151 .unwrap();
7152 let interrupt = engine.drain_effects().pop().unwrap();
7153 engine.drain_events();
7154 engine.next_event_sequence = Some(u64::MAX);
7155 let before = engine.snapshot();
7156
7157 let error = engine
7158 .try_apply_observation(ObservationEnvelope {
7159 operation_id: Some(interrupt.operation_id),
7160 instance_id: instance(),
7161 generation: spawn.generation,
7162 observation: ControlObservation::InputCompleted,
7163 })
7164 .unwrap_err();
7165
7166 assert_eq!(error, ControlError::EventSequenceExhausted);
7167 let after = engine.snapshot();
7168 assert_eq!(after.sessions, before.sessions);
7169 assert_eq!(after.revision, before.revision);
7170 assert!(!before.health.event_sequence_exhausted);
7171 assert!(after.health.event_sequence_exhausted);
7172 assert!(engine.drain_events().is_empty());
7173 }
7174
7175 #[test]
7176 fn stale_generation_cannot_mutate_restarted_session() {
7177 let (mut engine, first_spawn) = running_engine();
7178 engine.apply_observation(ObservationEnvelope {
7179 operation_id: None,
7180 instance_id: instance(),
7181 generation: first_spawn.generation,
7182 observation: ControlObservation::ProcessExited {
7183 exit_code: Some(0),
7184 final_terminal: None,
7185 },
7186 });
7187 engine.apply_command(start(4)).unwrap();
7188 let second_spawn = engine.drain_effects().pop().unwrap();
7189 assert!(second_spawn.generation.0 > first_spawn.generation.0);
7190
7191 engine.apply_observation(ObservationEnvelope {
7192 operation_id: Some(first_spawn.operation_id),
7193 instance_id: instance(),
7194 generation: first_spawn.generation,
7195 observation: ControlObservation::Spawned {
7196 process_id: Some(99),
7197 },
7198 });
7199 let session = &engine.snapshot().sessions[0];
7200 assert_eq!(session.generation, second_spawn.generation);
7201 assert_eq!(session.status, SessionStatus::Starting);
7202 assert_eq!(session.process_id, None);
7203 assert!(engine.drain_events().iter().any(|event| matches!(
7204 event.event,
7205 ControlEventKind::ObservationIgnored {
7206 reason: ObservationIgnoredReason::StaleGeneration
7207 }
7208 )));
7209 }
7210
7211 #[test]
7212 fn active_instance_cannot_be_removed_through_a_second_door() {
7213 let (mut engine, _) = running_engine();
7214 let error = engine
7215 .apply_command(CommandEnvelope {
7216 id: CommandId(5),
7217 command: ControlCommand::Remove {
7218 instance_id: instance(),
7219 },
7220 })
7221 .unwrap_err();
7222 assert!(matches!(error, ControlError::InvalidTransition { .. }));
7223 assert_eq!(engine.snapshot().sessions.len(), 1);
7224 }
7225
7226 #[test]
7227 fn history_has_independent_correlation_and_full_snapshot_results() {
7228 let mut engine = Gate4AgentEngine::new();
7229 engine.apply_command(register(1)).unwrap();
7230 engine.drain_events();
7231 engine
7232 .apply_command(CommandEnvelope {
7233 id: CommandId(2),
7234 command: ControlCommand::DiscoverHistory {
7235 instance_id: instance(),
7236 query: HistoryQuery {
7237 working_directory: Some("/repo".to_owned()),
7238 limit: 4,
7239 },
7240 },
7241 })
7242 .unwrap();
7243 let discovery = engine.drain_effects().pop().unwrap();
7244 let session = &engine.snapshot().sessions[0];
7245 assert_eq!(session.pending_operation, None);
7246 assert_eq!(
7247 session
7248 .history
7249 .pending
7250 .as_ref()
7251 .map(|pending| pending.operation_id),
7252 Some(discovery.operation_id)
7253 );
7254
7255 let candidate = HistoryCandidateSummary {
7256 id: "hist_fixture_1".to_owned(),
7257 session_id_hint: "session-1".to_owned(),
7258 modified_at_unix_ms: Some(42),
7259 };
7260 engine.apply_observation(ObservationEnvelope {
7261 operation_id: Some(discovery.operation_id),
7262 instance_id: instance(),
7263 generation: discovery.generation,
7264 observation: ControlObservation::HistoryDiscovered {
7265 candidates: vec![candidate.clone()],
7266 },
7267 });
7268 assert_eq!(
7269 engine.snapshot().sessions[0].history.candidates,
7270 vec![candidate]
7271 );
7272
7273 engine
7274 .apply_command(CommandEnvelope {
7275 id: CommandId(3),
7276 command: ControlCommand::LoadHistory {
7277 instance_id: instance(),
7278 candidate_id: "hist_fixture_1".to_owned(),
7279 },
7280 })
7281 .unwrap();
7282 let load = engine.drain_effects().pop().unwrap();
7283 engine.apply_observation(ObservationEnvelope {
7284 operation_id: Some(load.operation_id),
7285 instance_id: instance(),
7286 generation: load.generation,
7287 observation: ControlObservation::HistoryLoaded {
7288 session: HistorySessionRecord {
7289 session_id: "session-1".to_owned(),
7290 title: Some("title".to_owned()),
7291 cwd: Some("/repo".to_owned()),
7292 model: Some("model".to_owned()),
7293 message_count: 1,
7294 completed_turn_count: None,
7295 total_tokens: 7,
7296 messages: vec![HistoryMessageRecord {
7297 role: HistoryMessageRole::User,
7298 text: "hello".to_owned(),
7299 }],
7300 },
7301 },
7302 });
7303 let history = &engine.snapshot().sessions[0].history;
7304 assert!(history.pending.is_none());
7305 assert_eq!(history.loaded.as_ref().unwrap().session_id, "session-1");
7306 assert!(engine.drain_events().iter().any(|event| matches!(
7307 &event.event,
7308 ControlEventKind::HistoryLoaded { session_id } if session_id == "session-1"
7309 )));
7310 }
7311
7312 #[test]
7313 fn session_start_purges_queued_generation_bound_history_work() {
7314 let mut engine = Gate4AgentEngine::new();
7315 engine.apply_command(register(1)).unwrap();
7316 engine
7317 .apply_command(CommandEnvelope {
7318 id: CommandId(2),
7319 command: ControlCommand::DiscoverHistory {
7320 instance_id: instance(),
7321 query: HistoryQuery {
7322 working_directory: None,
7323 limit: 4,
7324 },
7325 },
7326 })
7327 .unwrap();
7328 let before_start = engine.snapshot().sessions[0].clone();
7329 let history_operation_id = before_start
7330 .history
7331 .pending
7332 .as_ref()
7333 .expect("history request must be pending")
7334 .operation_id;
7335 engine.apply_command(start(3)).unwrap();
7336 let effects = engine.drain_effects();
7337 assert_eq!(effects.len(), 1);
7338 assert!(matches!(effects[0].effect, ControlEffect::Spawn { .. }));
7339 let snapshot = engine.snapshot();
7340 assert!(snapshot.sessions[0].history.pending.is_none());
7341 assert!(snapshot.sessions[0].history.candidates.is_empty());
7342
7343 engine.apply_observation(ObservationEnvelope {
7344 operation_id: Some(history_operation_id),
7345 instance_id: instance(),
7346 generation: before_start.generation,
7347 observation: ControlObservation::HistoryDiscovered {
7348 candidates: Vec::new(),
7349 },
7350 });
7351 assert!(engine.drain_events().iter().any(|event| matches!(
7352 event.event,
7353 ControlEventKind::ObservationIgnored {
7354 reason: ObservationIgnoredReason::StaleGeneration
7355 }
7356 )));
7357 }
7358
7359 #[test]
7360 fn remove_rejects_pending_capability_probe_without_dropping_its_effect() {
7361 let mut engine = Gate4AgentEngine::new();
7362 engine.apply_command(register(1)).unwrap();
7363 engine
7364 .apply_command(CommandEnvelope {
7365 id: CommandId(2),
7366 command: ControlCommand::ProbeCapabilities {
7367 instance_id: instance(),
7368 request: CapabilityProbeRequest {
7369 working_directory: ".".to_owned(),
7370 },
7371 },
7372 })
7373 .unwrap();
7374 let pending = engine.snapshot().sessions[0]
7375 .capabilities
7376 .pending
7377 .as_ref()
7378 .expect("capability probe must be pending")
7379 .operation_id;
7380 let before = engine.clone();
7381
7382 let error = engine
7383 .apply_command(CommandEnvelope {
7384 id: CommandId(3),
7385 command: ControlCommand::Remove {
7386 instance_id: instance(),
7387 },
7388 })
7389 .unwrap_err();
7390
7391 assert_eq!(
7392 error,
7393 ControlError::CapabilityProbeOperationPending {
7394 operation_id: pending,
7395 }
7396 );
7397 assert_eq!(engine, before);
7398 assert!(matches!(
7399 engine.drain_effects().as_slice(),
7400 [EffectEnvelope {
7401 effect: ControlEffect::ProbeCapabilities { .. },
7402 ..
7403 }]
7404 ));
7405 }
7406
7407 #[test]
7408 fn remove_rejects_pending_history_without_dropping_its_effect() {
7409 let mut engine = Gate4AgentEngine::new();
7410 engine.apply_command(register(1)).unwrap();
7411 engine
7412 .apply_command(CommandEnvelope {
7413 id: CommandId(2),
7414 command: ControlCommand::DiscoverHistory {
7415 instance_id: instance(),
7416 query: HistoryQuery {
7417 working_directory: None,
7418 limit: 1,
7419 },
7420 },
7421 })
7422 .unwrap();
7423 let pending = engine.snapshot().sessions[0]
7424 .history
7425 .pending
7426 .as_ref()
7427 .expect("history request must be pending")
7428 .operation_id;
7429 let before = engine.clone();
7430
7431 let error = engine
7432 .apply_command(CommandEnvelope {
7433 id: CommandId(3),
7434 command: ControlCommand::Remove {
7435 instance_id: instance(),
7436 },
7437 })
7438 .unwrap_err();
7439
7440 assert_eq!(
7441 error,
7442 ControlError::HistoryOperationPending {
7443 operation_id: pending,
7444 }
7445 );
7446 assert_eq!(engine, before);
7447 assert!(matches!(
7448 engine.drain_effects().as_slice(),
7449 [EffectEnvelope {
7450 effect: ControlEffect::DiscoverHistory { .. },
7451 ..
7452 }]
7453 ));
7454 }
7455
7456 #[test]
7457 fn invalid_history_result_cannot_clear_the_correlated_operation() {
7458 let mut engine = Gate4AgentEngine::new();
7459 engine.apply_command(register(1)).unwrap();
7460 engine
7461 .apply_command(CommandEnvelope {
7462 id: CommandId(2),
7463 command: ControlCommand::DiscoverHistory {
7464 instance_id: instance(),
7465 query: HistoryQuery {
7466 working_directory: None,
7467 limit: 1,
7468 },
7469 },
7470 })
7471 .unwrap();
7472 let discovery = engine.drain_effects().pop().unwrap();
7473
7474 engine.apply_observation(ObservationEnvelope {
7475 operation_id: Some(discovery.operation_id),
7476 instance_id: instance(),
7477 generation: discovery.generation,
7478 observation: ControlObservation::HistoryDiscovered {
7479 candidates: vec![
7480 HistoryCandidateSummary {
7481 id: "hist_duplicate".to_owned(),
7482 session_id_hint: "session-1".to_owned(),
7483 modified_at_unix_ms: None,
7484 };
7485 2
7486 ],
7487 },
7488 });
7489
7490 let snapshot = engine.snapshot();
7491 assert_eq!(
7492 snapshot.sessions[0]
7493 .history
7494 .pending
7495 .as_ref()
7496 .map(|pending| pending.operation_id),
7497 Some(discovery.operation_id)
7498 );
7499 assert!(snapshot.sessions[0].history.candidates.is_empty());
7500 assert!(engine.drain_events().iter().any(|event| matches!(
7501 event.event,
7502 ControlEventKind::ObservationIgnored {
7503 reason: ObservationIgnoredReason::InvalidHistoryObservation
7504 }
7505 )));
7506 }
7507
7508 #[test]
7509 fn resume_is_authorized_before_a_new_generation_can_spawn() {
7510 let (mut engine, identity) = inactive_engine_with_provider_session();
7511 let previous_generation = engine.snapshot().sessions[0].generation;
7512 engine
7513 .apply_command(CommandEnvelope {
7514 id: CommandId(10),
7515 command: ControlCommand::Resume {
7516 instance_id: instance(),
7517 runtime_policy: verified_runtime_policy(),
7518 target: ResumeTarget::CurrentProvider,
7519 request: ResumeLaunchRequest {
7520 working_directory: ".".to_owned(),
7521 terminal_size: TerminalSize {
7522 rows: 24,
7523 columns: 80,
7524 },
7525 initial_prompt: None,
7526 },
7527 },
7528 })
7529 .unwrap();
7530 let authorize = engine.drain_effects().pop().unwrap();
7531 assert_eq!(authorize.generation, previous_generation);
7532 assert!(matches!(
7533 &authorize.effect,
7534 ControlEffect::AuthorizeResume {
7535 target: ResumeAuthorityTarget::ProviderSession { identity: authorized },
7536 ..
7537 } if authorized == &identity
7538 ));
7539 assert_eq!(
7540 engine.snapshot().sessions[0]
7541 .resume
7542 .pending
7543 .as_ref()
7544 .map(|pending| pending.phase),
7545 Some(ResumePhase::Authorizing)
7546 );
7547
7548 engine.apply_observation(ObservationEnvelope {
7549 operation_id: Some(authorize.operation_id),
7550 instance_id: instance(),
7551 generation: authorize.generation,
7552 observation: ControlObservation::ResumeAuthorized {
7553 provider_session: identity.clone(),
7554 },
7555 });
7556 let spawn = engine.drain_effects().pop().unwrap();
7557 assert_eq!(spawn.operation_id, authorize.operation_id);
7558 assert_eq!(spawn.generation.0, previous_generation.0 + 1);
7559 assert!(matches!(
7560 &spawn.effect,
7561 ControlEffect::SpawnResume {
7562 transport: TransportKind::Pty,
7563 provider_session,
7564 runtime_policy,
7565 ..
7566 } if provider_session == &identity
7567 && *runtime_policy == verified_runtime_policy()
7568 ));
7569 assert_eq!(
7570 engine.snapshot().sessions[0].status,
7571 SessionStatus::Starting
7572 );
7573
7574 engine.apply_observation(ObservationEnvelope {
7575 operation_id: Some(spawn.operation_id),
7576 instance_id: instance(),
7577 generation: spawn.generation,
7578 observation: ControlObservation::Spawned {
7579 process_id: Some(77),
7580 },
7581 });
7582 let session = &engine.snapshot().sessions[0];
7583 assert_eq!(session.status, SessionStatus::Running);
7584 assert!(session.resume.pending.is_none());
7585 assert_eq!(session.provider.session.as_ref(), Some(&identity));
7586 assert_eq!(
7587 session.resume.last_session,
7588 Some(ResumeSessionSummary::from(&identity))
7589 );
7590 assert!(engine.drain_events().iter().any(|event| matches!(
7591 &event.event,
7592 ControlEventKind::Resumed { session, process_id: Some(77) }
7593 if session.id == "provider-session-1"
7594 )));
7595 }
7596
7597 #[test]
7598 fn explicit_provider_session_resume_still_requires_matching_authority() {
7599 let mut engine = Gate4AgentEngine::new();
7600 engine.apply_command(register(1)).unwrap();
7601 engine.drain_events();
7602 let identity = ProviderSessionIdentity {
7603 key: ProviderSessionKey::SessionId,
7604 id: "durable-provider-session".to_owned(),
7605 transcript_path: None,
7606 };
7607 engine
7608 .apply_command(CommandEnvelope {
7609 id: CommandId(2),
7610 command: ControlCommand::Resume {
7611 instance_id: instance(),
7612 runtime_policy: verified_runtime_policy(),
7613 target: ResumeTarget::ProviderSession {
7614 identity: identity.clone(),
7615 },
7616 request: ResumeLaunchRequest {
7617 working_directory: ".".to_owned(),
7618 terminal_size: TerminalSize {
7619 rows: 24,
7620 columns: 80,
7621 },
7622 initial_prompt: None,
7623 },
7624 },
7625 })
7626 .unwrap();
7627 let authorize = engine.drain_effects().pop().unwrap();
7628 assert!(matches!(
7629 &authorize.effect,
7630 ControlEffect::AuthorizeResume {
7631 target: ResumeAuthorityTarget::ProviderSession { identity: requested },
7632 ..
7633 } if requested == &identity
7634 ));
7635 assert_eq!(engine.snapshot().sessions[0].provider.session, None);
7636 assert_eq!(
7637 engine.snapshot().sessions[0].status,
7638 SessionStatus::Registered
7639 );
7640
7641 engine.apply_observation(ObservationEnvelope {
7642 operation_id: Some(authorize.operation_id),
7643 instance_id: instance(),
7644 generation: authorize.generation,
7645 observation: ControlObservation::ResumeAuthorized {
7646 provider_session: ProviderSessionIdentity {
7647 key: ProviderSessionKey::SessionId,
7648 id: "wrong-provider-session".to_owned(),
7649 transcript_path: None,
7650 },
7651 },
7652 });
7653 assert!(engine.drain_effects().is_empty());
7654 assert_eq!(
7655 engine.snapshot().sessions[0].status,
7656 SessionStatus::Registered
7657 );
7658
7659 engine.apply_observation(ObservationEnvelope {
7660 operation_id: Some(authorize.operation_id),
7661 instance_id: instance(),
7662 generation: authorize.generation,
7663 observation: ControlObservation::ResumeAuthorized {
7664 provider_session: identity.clone(),
7665 },
7666 });
7667 let spawn = engine.drain_effects().pop().unwrap();
7668 assert_eq!(spawn.generation, SessionGeneration(1));
7669 assert!(matches!(
7670 spawn.effect,
7671 ControlEffect::SpawnResume {
7672 provider_session,
7673 ..
7674 } if provider_session == identity
7675 ));
7676 }
7677
7678 #[test]
7679 fn pipe_resume_requires_and_preserves_a_normalized_initial_prompt() {
7680 let (mut engine, identity) = inactive_engine_with_provider_session();
7681 engine.session_mut(instance()).transport = TransportKind::Pipe;
7682
7683 let missing_prompt = engine.apply_command(CommandEnvelope {
7684 id: CommandId(20),
7685 command: ControlCommand::Resume {
7686 instance_id: instance(),
7687 runtime_policy: verified_runtime_policy(),
7688 target: ResumeTarget::CurrentProvider,
7689 request: ResumeLaunchRequest {
7690 working_directory: ".".to_owned(),
7691 terminal_size: TerminalSize {
7692 rows: 24,
7693 columns: 80,
7694 },
7695 initial_prompt: None,
7696 },
7697 },
7698 });
7699 assert_eq!(missing_prompt, Err(ControlError::MissingInitialPrompt));
7700 assert!(engine.drain_effects().is_empty());
7701
7702 engine.session_mut(instance()).transport = TransportKind::Acp;
7703 assert!(matches!(
7704 engine.apply_command(CommandEnvelope {
7705 id: CommandId(21),
7706 command: ControlCommand::Resume {
7707 instance_id: instance(),
7708 runtime_policy: verified_runtime_policy(),
7709 target: ResumeTarget::CurrentProvider,
7710 request: ResumeLaunchRequest {
7711 working_directory: ".".to_owned(),
7712 terminal_size: TerminalSize {
7713 rows: 24,
7714 columns: 80,
7715 },
7716 initial_prompt: Some("continue".to_owned()),
7717 },
7718 },
7719 }),
7720 Err(ControlError::UnsupportedTransportOperation {
7721 transport: TransportKind::Acp,
7722 ..
7723 })
7724 ));
7725 engine.session_mut(instance()).transport = TransportKind::Pipe;
7726
7727 engine
7728 .apply_command(CommandEnvelope {
7729 id: CommandId(22),
7730 command: ControlCommand::Resume {
7731 instance_id: instance(),
7732 runtime_policy: verified_runtime_policy(),
7733 target: ResumeTarget::CurrentProvider,
7734 request: ResumeLaunchRequest {
7735 working_directory: ".".to_owned(),
7736 terminal_size: TerminalSize {
7737 rows: 24,
7738 columns: 80,
7739 },
7740 initial_prompt: Some("continue\u{0000}now".to_owned()),
7741 },
7742 },
7743 })
7744 .unwrap();
7745 let authorize = engine.drain_effects().pop().unwrap();
7746 assert!(matches!(
7747 &authorize.effect,
7748 ControlEffect::AuthorizeResume { request, .. }
7749 if request.initial_prompt.as_deref() == Some("continue<U+0000>now")
7750 ));
7751
7752 engine.apply_observation(ObservationEnvelope {
7753 operation_id: Some(authorize.operation_id),
7754 instance_id: instance(),
7755 generation: authorize.generation,
7756 observation: ControlObservation::ResumeAuthorized {
7757 provider_session: identity,
7758 },
7759 });
7760 assert!(matches!(
7761 engine.drain_effects().pop().unwrap().effect,
7762 ControlEffect::SpawnResume {
7763 transport: TransportKind::Pipe,
7764 request,
7765 ..
7766 } if request.initial_prompt.as_deref() == Some("continue<U+0000>now")
7767 ));
7768 }
7769
7770 #[test]
7771 fn resume_authorization_generation_exhaustion_is_atomic_and_ignored() {
7772 let (mut engine, identity) = inactive_engine_with_provider_session();
7773 let exhausted = SessionGeneration(u64::MAX);
7774 engine.session_mut(instance()).generation = exhausted;
7775 engine.generation_watermarks.insert(instance(), exhausted);
7776 engine
7777 .apply_command(CommandEnvelope {
7778 id: CommandId(10),
7779 command: ControlCommand::Resume {
7780 instance_id: instance(),
7781 runtime_policy: verified_runtime_policy(),
7782 target: ResumeTarget::CurrentProvider,
7783 request: ResumeLaunchRequest {
7784 working_directory: ".".to_owned(),
7785 terminal_size: TerminalSize {
7786 rows: 24,
7787 columns: 80,
7788 },
7789 initial_prompt: None,
7790 },
7791 },
7792 })
7793 .unwrap();
7794 let authorize = engine.drain_effects().pop().unwrap();
7795 engine.drain_events();
7796 let before = engine.snapshot().sessions[0].clone();
7797
7798 engine.apply_observation(ObservationEnvelope {
7799 operation_id: Some(authorize.operation_id),
7800 instance_id: instance(),
7801 generation: authorize.generation,
7802 observation: ControlObservation::ResumeAuthorized {
7803 provider_session: identity,
7804 },
7805 });
7806
7807 assert_eq!(engine.snapshot().sessions[0], before);
7808 assert_eq!(
7809 engine.generation_watermarks.get(&instance()),
7810 Some(&exhausted)
7811 );
7812 assert!(engine.drain_effects().is_empty());
7813 assert!(engine.drain_events().iter().any(|event| matches!(
7814 event.event,
7815 ControlEventKind::ObservationIgnored {
7816 reason: ObservationIgnoredReason::GenerationExhausted
7817 }
7818 )));
7819 }
7820
7821 #[test]
7822 fn resume_spawn_failure_retains_the_authorized_identity_for_retry() {
7823 let (mut engine, identity) = inactive_engine_with_provider_session();
7824 engine
7825 .apply_command(CommandEnvelope {
7826 id: CommandId(12),
7827 command: ControlCommand::Resume {
7828 instance_id: instance(),
7829 runtime_policy: verified_runtime_policy(),
7830 target: ResumeTarget::CurrentProvider,
7831 request: ResumeLaunchRequest {
7832 working_directory: ".".to_owned(),
7833 terminal_size: TerminalSize {
7834 rows: 24,
7835 columns: 80,
7836 },
7837 initial_prompt: None,
7838 },
7839 },
7840 })
7841 .unwrap();
7842 let authorize = engine.drain_effects().pop().unwrap();
7843 engine.apply_observation(ObservationEnvelope {
7844 operation_id: Some(authorize.operation_id),
7845 instance_id: instance(),
7846 generation: authorize.generation,
7847 observation: ControlObservation::ResumeAuthorized {
7848 provider_session: identity.clone(),
7849 },
7850 });
7851 let spawn = engine.drain_effects().pop().unwrap();
7852 engine.apply_observation(ObservationEnvelope {
7853 operation_id: Some(spawn.operation_id),
7854 instance_id: instance(),
7855 generation: spawn.generation,
7856 observation: ControlObservation::SpawnFailed {
7857 message: "controlled spawn failure".to_owned(),
7858 },
7859 });
7860
7861 let session = &engine.snapshot().sessions[0];
7862 assert_eq!(session.provider.session.as_ref(), Some(&identity));
7863 assert_eq!(
7864 session.resume.last_error.as_deref(),
7865 Some("controlled spawn failure")
7866 );
7867 assert!(session.resume.last_session.is_none());
7868 engine
7869 .apply_command(CommandEnvelope {
7870 id: CommandId(13),
7871 command: ControlCommand::Resume {
7872 instance_id: instance(),
7873 runtime_policy: verified_runtime_policy(),
7874 target: ResumeTarget::CurrentProvider,
7875 request: ResumeLaunchRequest {
7876 working_directory: ".".to_owned(),
7877 terminal_size: TerminalSize {
7878 rows: 24,
7879 columns: 80,
7880 },
7881 initial_prompt: None,
7882 },
7883 },
7884 })
7885 .unwrap();
7886 assert!(matches!(
7887 engine.drain_effects().pop().unwrap().effect,
7888 ControlEffect::AuthorizeResume { .. }
7889 ));
7890 }
7891
7892 #[test]
7893 fn resume_denial_preserves_the_inactive_session_generation_and_status() {
7894 let (mut engine, _) = inactive_engine_with_provider_session();
7895 let before = engine.snapshot().sessions[0].clone();
7896 engine
7897 .apply_command(CommandEnvelope {
7898 id: CommandId(11),
7899 command: ControlCommand::Resume {
7900 instance_id: instance(),
7901 runtime_policy: verified_runtime_policy(),
7902 target: ResumeTarget::CurrentProvider,
7903 request: ResumeLaunchRequest {
7904 working_directory: ".".to_owned(),
7905 terminal_size: TerminalSize {
7906 rows: 24,
7907 columns: 80,
7908 },
7909 initial_prompt: None,
7910 },
7911 },
7912 })
7913 .unwrap();
7914 let authorize = engine.drain_effects().pop().unwrap();
7915 engine.apply_observation(ObservationEnvelope {
7916 operation_id: Some(authorize.operation_id),
7917 instance_id: instance(),
7918 generation: authorize.generation,
7919 observation: ControlObservation::ResumeDenied {
7920 reason: "vendor login is required".to_owned(),
7921 },
7922 });
7923
7924 let session = &engine.snapshot().sessions[0];
7925 assert_eq!(session.generation, before.generation);
7926 assert_eq!(session.status, before.status);
7927 assert_eq!(
7928 session.resume.last_error.as_deref(),
7929 Some("vendor login is required")
7930 );
7931 assert!(session.resume.pending.is_none());
7932 assert!(engine.drain_effects().is_empty());
7933 }
7934
7935 #[test]
7936 fn history_resume_requires_the_loaded_candidate_and_exact_parsed_session() {
7937 let mut engine = Gate4AgentEngine::new();
7938 engine.apply_command(register(1)).unwrap();
7939 engine
7940 .apply_command(CommandEnvelope {
7941 id: CommandId(2),
7942 command: ControlCommand::DiscoverHistory {
7943 instance_id: instance(),
7944 query: HistoryQuery {
7945 working_directory: None,
7946 limit: 4,
7947 },
7948 },
7949 })
7950 .unwrap();
7951 let discover = engine.drain_effects().pop().unwrap();
7952 engine.apply_observation(ObservationEnvelope {
7953 operation_id: Some(discover.operation_id),
7954 instance_id: instance(),
7955 generation: discover.generation,
7956 observation: ControlObservation::HistoryDiscovered {
7957 candidates: vec![HistoryCandidateSummary {
7958 id: "hist_resume_1".to_owned(),
7959 session_id_hint: "hint-only".to_owned(),
7960 modified_at_unix_ms: None,
7961 }],
7962 },
7963 });
7964 let request = ResumeLaunchRequest {
7965 working_directory: ".".to_owned(),
7966 terminal_size: TerminalSize {
7967 rows: 24,
7968 columns: 80,
7969 },
7970 initial_prompt: None,
7971 };
7972 assert_eq!(
7973 engine.apply_command(CommandEnvelope {
7974 id: CommandId(3),
7975 command: ControlCommand::Resume {
7976 instance_id: instance(),
7977 runtime_policy: verified_runtime_policy(),
7978 target: ResumeTarget::HistoryCandidate {
7979 candidate_id: "hist_resume_1".to_owned(),
7980 },
7981 request: request.clone(),
7982 },
7983 }),
7984 Err(ControlError::HistoryCandidateNotLoaded)
7985 );
7986
7987 engine
7988 .apply_command(CommandEnvelope {
7989 id: CommandId(4),
7990 command: ControlCommand::LoadHistory {
7991 instance_id: instance(),
7992 candidate_id: "hist_resume_1".to_owned(),
7993 },
7994 })
7995 .unwrap();
7996 let load = engine.drain_effects().pop().unwrap();
7997 engine.apply_observation(ObservationEnvelope {
7998 operation_id: Some(load.operation_id),
7999 instance_id: instance(),
8000 generation: load.generation,
8001 observation: ControlObservation::HistoryLoaded {
8002 session: HistorySessionRecord {
8003 session_id: "parsed-session-1".to_owned(),
8004 title: None,
8005 cwd: None,
8006 model: None,
8007 message_count: 0,
8008 completed_turn_count: None,
8009 total_tokens: 0,
8010 messages: Vec::new(),
8011 },
8012 },
8013 });
8014 engine
8015 .apply_command(CommandEnvelope {
8016 id: CommandId(5),
8017 command: ControlCommand::Resume {
8018 instance_id: instance(),
8019 runtime_policy: verified_runtime_policy(),
8020 target: ResumeTarget::HistoryCandidate {
8021 candidate_id: "hist_resume_1".to_owned(),
8022 },
8023 request,
8024 },
8025 })
8026 .unwrap();
8027 let authorize = engine.drain_effects().pop().unwrap();
8028 engine.apply_observation(ObservationEnvelope {
8029 operation_id: Some(authorize.operation_id),
8030 instance_id: instance(),
8031 generation: authorize.generation,
8032 observation: ControlObservation::ResumeAuthorized {
8033 provider_session: ProviderSessionIdentity {
8034 key: ProviderSessionKey::SessionId,
8035 id: "wrong-session".to_owned(),
8036 transcript_path: None,
8037 },
8038 },
8039 });
8040 assert!(engine.drain_effects().is_empty());
8041 assert_eq!(
8042 engine.snapshot().sessions[0]
8043 .resume
8044 .pending
8045 .as_ref()
8046 .map(|pending| pending.phase),
8047 Some(ResumePhase::Authorizing)
8048 );
8049 assert!(engine.drain_events().iter().any(|event| matches!(
8050 event.event,
8051 ControlEventKind::ObservationIgnored {
8052 reason: ObservationIgnoredReason::InvalidResumeObservation
8053 }
8054 )));
8055 }
8056
8057 #[test]
8058 fn replay_is_deterministic() {
8059 fn replay() -> (ControlSnapshot, Vec<EffectEnvelope>, Vec<ControlEvent>) {
8060 let mut engine = Gate4AgentEngine::new();
8061 engine.apply_command(register(1)).unwrap();
8062 engine.apply_command(start(2)).unwrap();
8063 let effect = engine.drain_effects().pop().unwrap();
8064 engine.apply_observation(ObservationEnvelope {
8065 operation_id: Some(effect.operation_id),
8066 instance_id: effect.instance_id,
8067 generation: effect.generation,
8068 observation: ControlObservation::Spawned {
8069 process_id: Some(42),
8070 },
8071 });
8072 (engine.snapshot(), vec![effect], engine.drain_events())
8073 }
8074
8075 assert_eq!(replay(), replay());
8076 }
8077}