Skip to main content

gate4agent_engine/
lib.rs

1//! Deterministic single-writer lifecycle engine for gate4agent sessions.
2
3use 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/// Owns logical session lifecycle state. External work is emitted as effects
45/// and can change observed state only after a matching observation.
46#[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                            // No observation exists yet for a freshly registered session.
147                            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 = &current_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(&current.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            (&current.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                        // A new generation must not inherit the old one's
642                        // screen classification -- nothing has observed
643                        // this generation's PTY yet.
644                        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            // The observation is the authority for `SessionSnapshot::screen_state`;
880            // `TerminalFrame::screen_state` is only a per-frame stamp for terminal
881            // subscribers reading a single frame in isolation. Both are read off
882            // the same value computed at the node, but only this observation sets
883            // the snapshot field a caller consults without subscribing to frames.
884            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(&current.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                // One refused observation must never wedge every later
915                // observation on this source for the rest of the session:
916                // advance the per-source cursor to `*sequence` exactly as an
917                // admitted observation would have left it -- the sender
918                // (`gate4agent-shell-native`'s `drain_provider_stream`)
919                // advances its own counter the moment it hands an event off,
920                // independent of whether the engine goes on to admit it, so
921                // this per-source cursor stays the SOLE authority for "how
922                // far this source has been consumed" instead of growing a
923                // second, sender-side notion of it. The refusal itself must
924                // not be silent either: it mints a `ProviderEvent::Error`
925                // naming the withheld capability, the same shape `ingest_
926                // provider`'s batched refusal handling uses, so an
927                // operator watching a live subscription sees a refusal
928                // instead of nothing.
929                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(&current.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            // A new generation must not inherit the old one's screen
1538            // classification -- nothing has observed this generation's
1539            // PTY yet.
1540            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            // Lifecycle hooks are retired (owner ruling 2026-09-25): a
2285            // session is observed through ACP where the provider has it, or
2286            // through the PTY-parsing verification chain otherwise -- every
2287            // source family proves its semantics through `SemanticReadiness`.
2288            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            // The command is still refused -- this call still returns `Err`
2305            // below, exactly as it always has, so a caller inspecting THIS
2306            // command's own outcome (and `gate4agent-kernel`'s existing
2307            // `record_command_rejection`, which fires off that same `Err`)
2308            // keeps seeing precisely what it always saw. What changes is
2309            // what a refusal leaves behind for every LATER batch on this
2310            // source: one refused batch must never wedge every later one for
2311            // the rest of the session. The sender (`gate4agent-shell-
2312            // native`'s `drain_provider_stream`) advances its own sequence
2313            // counter the moment it hands an event off, independent of
2314            // whether the engine goes on to admit it, so this per-source
2315            // cursor is kept the SOLE authority for "how far this source has
2316            // been consumed" rather than growing a second, sender-side
2317            // notion of it: advancing the cursor here, exactly as an
2318            // admitted batch would have left it, is what lets the next
2319            // (capability-clean) batch land instead of reading as stale
2320            // forever. The refusal itself must not be silent either: it
2321            // mints a `ProviderEvent::Error` naming the withheld capability
2322            // and how many events it took down with it -- the same
2323            // observation shape a real provider-reported error uses -- so an
2324            // operator watching a live subscription sees a refusal instead
2325            // of nothing.
2326            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    /// Request a switch of the session's ACP `session/set_mode` mode.
2544    /// ACP-transport only, same gating shape as [`resize`](Self::resize):
2545    /// the session must be `Running`, ACP-transport, and have no other
2546    /// operation already pending. Whether `mode_id` is one the agent
2547    /// actually offered is the shell executor's job (`AcpSession::
2548    /// available_modes`), not this engine's -- it has no live catalogue of
2549    /// its own to check against.
2550    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    /// Request an ACP `session/set_config_option` value change. Same
2616    /// gating shape as [`set_session_mode`](Self::set_session_mode);
2617    /// `value_json` is carried as pre-serialized JSON text (see
2618    /// `ControlCommand::SetSessionConfigOption`'s doc comment for why) and
2619    /// is only bounds-checked here, never parsed -- this crate does not
2620    /// depend on `serde_json`.
2621    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    /// Request a switch of the session's active model via a provider vendor
2694    /// extension. Same gating shape as [`set_session_mode`](Self::set_session_mode);
2695    /// see `ControlCommand::SetSessionModel`'s doc comment for why this is
2696    /// wired end to end despite no current shell executor being able to
2697    /// honor it for any known provider.
2698    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    /// Starting a new generation invalidates history work, while host-scoped
2863    /// capability discovery intentionally remains valid across that boundary.
2864    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::ProviderEvent` is deliberately NOT classified
3030        // here -- unlike every other kind this function gates, one denied
3031        // provider event must not be dropped silently: the dedicated
3032        // `ControlObservation::ProviderEvent` arm in `apply_observation_in_
3033        // place` runs the same capability check itself, and on denial keeps
3034        // the per-source cursor advancing (so the next observation on this
3035        // source is never read as stale) and mints a visible
3036        // `ProviderEvent::Error` naming the withheld capability instead of
3037        // vanishing into `ObservationIgnoredReason::ProviderRuntimePolicyDenied`.
3038        ControlObservation::ProviderGap { .. } => ProviderRuntimeCapability::SemanticReadiness,
3039        // Resume authority observations are correlated to an exact pending
3040        // operation and are not provider telemetry. Command admission already
3041        // requires the semantic resume contract when a prompt is injected;
3042        // a prompt-free provider-native PTY resume is intentionally raw.
3043        ControlObservation::ResumeAuthorized { .. }
3044        | ControlObservation::ResumeDenied { .. }
3045        | ControlObservation::ResumeFailed { .. } => return None,
3046        _ => return None,
3047    };
3048    (!runtime_policy.admits(capability)).then_some(capability)
3049}
3050
3051/// The capability `event` needs but `runtime_policy` withholds, or `None`
3052/// if `event` is already admitted. Shared by both provider-event ingress
3053/// paths (`Gate4AgentEngine::ingest_provider`'s batched command path and
3054/// `apply_observation_in_place`'s single-event observation path) so a
3055/// session-identity-carrying event is gated on `ProviderSessionIdentity`
3056/// alone, and every other event needs `SemanticReadiness` first and (only if
3057/// it independently carries identity) `ProviderSessionIdentity` too, the
3058/// same rule on both paths.
3059fn 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
3077/// The visible refusal text a denied provider-event batch or observation
3078/// mints as a `ProviderEvent::Error` -- named after the withheld capability
3079/// and how many events it took down with it, so an operator watching a live
3080/// subscription sees a refusal instead of nothing.
3081fn 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        // ACP session/update coverage beyond text/tool/turn streaming:
3492        // these carry real data (still recorded verbatim on `snapshot.
3493        // last_event` below), but none of them has an established snapshot
3494        // field or `lead_activity` semantics yet -- adding one is a
3495        // deliberate engine-schema decision, not a side effect of parsing
3496        // more of the wire protocol, so it is left to a follow-up that owns
3497        // that decision.
3498        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    /// Registers `instance()` on ACP transport instead of `register`'s
3853    /// hardcoded PTY -- `start()` requires `RawPtyLifecycle` only when
3854    /// `session.snapshot.transport == TransportKind::Pty` (see `Gate4Agent
3855    /// Engine::start`), so an ACP-shaped runtime policy (`raw_pty_lifecycle:
3856    /// false`) can only be started against a session actually registered as
3857    /// ACP, never against `register`'s PTY session with the transport
3858    /// overridden afterward -- that ordering hits the PTY check before the
3859    /// override ever runs.
3860    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    /// The ACP-transport counterpart to `running_engine_with_policy`:
3938    /// registers `instance()` on ACP transport up front (see `register_acp`)
3939    /// so an ACP-shaped `runtime_policy` (`raw_pty_lifecycle: false`) can
3940    /// actually reach `Spawned` instead of failing `start()`'s PTY-only
3941    /// `RawPtyLifecycle` requirement.
3942    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    /// Both provider-event ingress paths fail closed on a policy that
4006    /// withholds the capability an event needs -- but, unlike before this
4007    /// fix, a refusal is no longer silence AND no longer wedges the next,
4008    /// capability-clean call on the same source: each denial below still
4009    /// advances the per-source cursor (`sequence` climbs by one per call,
4010    /// never resets), and the observation-path refusals mint a visible
4011    /// `ProviderEvent::Error` in place of the old, invisible
4012    /// `ObservationIgnoredReason::ProviderRuntimePolicyDenied`. The command
4013    /// path keeps returning `Err(ProviderRuntimePolicyDenied)` from
4014    /// `apply_command` itself, exactly as it always has --
4015    /// `acp_runtime_policy_survives_a_capability_refusal_and_surfaces_it`
4016    /// is the focused regression test for the visibility/survival change
4017    /// itself.
4018    #[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    /// The ACP-shaped counterpart to `runtime_policy_keeps_raw_input_and_
4101    /// rejects_structured_prompt` above -- the two together pin the boundary
4102    /// from both sides. An ACP session's `runtime_policy` grants
4103    /// `SemanticReadiness`/`StructuredPrompt` as facts of the ACP protocol
4104    /// itself (no PTY to verify: `raw_pty_lifecycle` is false), so
4105    /// `SubmitPrompt` is admitted here where it used to come back
4106    /// `ProviderRuntimePolicyDenied` for every ACP session (the live-measured
4107    /// defect this fix closes). A PTY-only action landing on the same
4108    /// session is still refused -- by the transport `match` itself
4109    /// (`UnsupportedTransportOperation`), never by a policy flag that
4110    /// happens to be false.
4111    #[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    /// The observation-ingress counterpart: an ACP session's
4155    /// `SemanticReadiness` grant lets a non-identity `ProviderEvent` reach
4156    /// the session's provider sequence instead of being dropped as
4157    /// `ObservationIgnoredReason::ProviderRuntimePolicyDenied` -- the exact
4158    /// defect that left a live ACP agent-stream subscription observing zero
4159    /// frames before this fix.
4160    #[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    /// The observation-ingress counterpart to `ingest_provider`'s batched
4186    /// refusal handling, and the direct regression test for the
4187    /// live-measured defect this fix closes: an ACP-shaped policy that still
4188    /// withholds `ProviderSessionIdentity` refuses a `SessionIdentityObserved`
4189    /// observation, but that ONE refusal must (a) not wedge the session's
4190    /// provider stream for the rest of its life -- a later, capability-clean
4191    /// observation on the SAME source must still land -- and (b) surface
4192    /// itself as a visible `ProviderEvent::Error` naming the withheld
4193    /// capability, never vanish into
4194    /// `ObservationIgnoredReason::ProviderRuntimePolicyDenied` silence. Live,
4195    /// an ACP session's `session/new` response was denied exactly this way
4196    /// and the subscription observed nothing at all afterward.
4197    #[test]
4198    fn acp_runtime_policy_survives_a_capability_refusal_and_surfaces_it() {
4199        // `SemanticReadiness`/`StructuredPrompt` granted (the ACP protocol
4200        // shape); `ProviderSessionIdentity` deliberately withheld so this
4201        // fixture exercises the exact refusal this test is about.
4202        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        // The NEXT, capability-clean observation on the same source must
4246        // land -- the refusal above must not have wedged the per-source
4247        // cursor for the rest of the session.
4248        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        // `raw_pty_lifecycle: false` with `semantic_readiness`/
4277        // `structured_prompt`/`provider_session_identity: true` and nothing
4278        // else is no longer a construction defect on its own -- it is
4279        // exactly the shape `provider_runtime::policy_for_transport` grants
4280        // an ACP session (`session/prompt`/`session/update` are mandatory
4281        // ACP protocol surface, and `session/new` returns a `sessionId`
4282        // under that same specification, neither a PTY-terminal-text
4283        // inference). `hook_semantics: true` here keeps this fixture on the
4284        // side of the invariant that IS still a defect: nothing derives
4285        // hook-sourced semantics for a transport other than a raw PTY
4286        // lifecycle, so granting it without `raw_pty_lifecycle` stays
4287        // rejected by `ProviderRuntimePolicy::validate` itself, before
4288        // `start()` even inspects this (PTY) session's transport.
4289        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(&current));
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}