Skip to main content

gate4agent_types/
control.rs

1use crate::{
2    AdapterBinding, AdapterFamily, AgentId, CapabilityModelSummary, CapabilityProbeFailure,
3    CapabilityProbeRequest, CapabilitySnapshot, HistoryCandidateSummary, HistoryOperation,
4    HistoryQuery, HistorySessionRecord, HistorySnapshot, InputAction, InputPrepareError,
5    PreparedInput, PreparedInputKind, ResumeAuthorityTarget, ResumeLaunchRequest,
6    ResumeSessionSummary, ResumeSnapshot, ResumeTarget, SessionOptionSelection,
7};
8use serde::{Deserialize, Deserializer, Serialize};
9use thiserror::Error;
10
11pub const CONTROL_PROTOCOL_VERSION: u16 = 28;
12pub const CONTROL_SESSIONS_MAX: usize = 512;
13pub const CONTROL_INSTANCE_IDENTITIES_CAPACITY: u32 = 4_096;
14pub const CONTROL_INSTANCE_IDENTITIES_MAX: usize = CONTROL_INSTANCE_IDENTITIES_CAPACITY as usize;
15pub const TERMINAL_ROWS_MAX: u16 = 1_000;
16pub const TERMINAL_COLUMNS_MAX: u16 = 1_000;
17pub const WORKING_DIRECTORY_MAX_BYTES: usize = 32_768;
18pub const PROVIDER_INGRESS_EVENTS_MAX: usize = 32;
19pub const PROVIDER_EVENT_TEXT_MAX_BYTES: usize = 262_144;
20pub const PROVIDER_EVENT_ID_MAX_BYTES: usize = 512;
21pub const PROVIDER_EVENT_TOOLS_MAX: usize = 256;
22pub const PROVIDER_INTERACTIONS_MAX: usize = 64;
23pub const PROVIDER_INTERACTION_RESPONSE_MAX_BYTES: usize = 32_768;
24pub const PROVIDER_INTERACTION_FAILURE_MAX_BYTES: usize = 4_096;
25pub const PROVIDER_SUBAGENTS_MAX: usize = 64;
26pub const PROVIDER_SESSION_LOCATOR_MAX_BYTES: usize = 32_768;
27
28#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
29#[serde(transparent)]
30pub struct AgentInstanceId(pub u64);
31
32#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
33#[serde(transparent)]
34pub struct CommandId(pub u64);
35
36#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
37#[serde(transparent)]
38pub struct OperationId(pub u64);
39
40#[derive(
41    Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize,
42)]
43#[serde(transparent)]
44pub struct SessionGeneration(pub u64);
45
46#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
47pub struct StartRequest {
48    pub working_directory: String,
49    pub terminal_size: TerminalSize,
50    #[serde(default)]
51    pub initial_prompt: Option<String>,
52    #[serde(default)]
53    pub session_options: Option<SessionOptionSelection>,
54}
55
56#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
57pub struct ProviderRuntimePolicy {
58    pub raw_pty_lifecycle: bool,
59    pub semantic_readiness: bool,
60    pub structured_prompt: bool,
61    pub provider_session_identity: bool,
62    pub semantic_resume: bool,
63}
64
65impl ProviderRuntimePolicy {
66    pub fn new(
67        raw_pty_lifecycle: bool,
68        semantic_readiness: bool,
69        structured_prompt: bool,
70        provider_session_identity: bool,
71        semantic_resume: bool,
72    ) -> Result<Self, ProviderRuntimePolicyError> {
73        let policy = Self {
74            raw_pty_lifecycle,
75            semantic_readiness,
76            structured_prompt,
77            provider_session_identity,
78            semantic_resume,
79        };
80        policy.validate()?;
81        Ok(policy)
82    }
83
84    pub const fn raw_pty() -> Self {
85        Self {
86            raw_pty_lifecycle: true,
87            semantic_readiness: false,
88            structured_prompt: false,
89            provider_session_identity: false,
90            semantic_resume: false,
91        }
92    }
93
94    pub fn validate(self) -> Result<(), ProviderRuntimePolicyError> {
95        if (self.semantic_readiness
96            || self.structured_prompt
97            || self.provider_session_identity
98            || self.semantic_resume)
99            && !self.raw_pty_lifecycle
100        {
101            return Err(ProviderRuntimePolicyError::SemanticCapabilityRequiresRawPty);
102        }
103        if self.structured_prompt && !self.semantic_readiness {
104            return Err(ProviderRuntimePolicyError::StructuredPromptRequiresReadiness);
105        }
106        if self.semantic_resume && !self.provider_session_identity {
107            return Err(ProviderRuntimePolicyError::ResumeRequiresSessionIdentity);
108        }
109        Ok(())
110    }
111
112    pub const fn admits(self, capability: ProviderRuntimeCapability) -> bool {
113        match capability {
114            ProviderRuntimeCapability::RawPtyLifecycle => self.raw_pty_lifecycle,
115            ProviderRuntimeCapability::SemanticReadiness => self.semantic_readiness,
116            ProviderRuntimeCapability::StructuredPrompt => self.structured_prompt,
117            ProviderRuntimeCapability::ProviderSessionIdentity => {
118                self.provider_session_identity
119            }
120            ProviderRuntimeCapability::SemanticResume => self.semantic_resume,
121        }
122    }
123}
124
125impl<'de> Deserialize<'de> for ProviderRuntimePolicy {
126    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
127    where
128        D: Deserializer<'de>,
129    {
130        #[derive(Deserialize)]
131        struct WirePolicy {
132            raw_pty_lifecycle: bool,
133            semantic_readiness: bool,
134            structured_prompt: bool,
135            provider_session_identity: bool,
136            semantic_resume: bool,
137        }
138
139        let wire = WirePolicy::deserialize(deserializer)?;
140        Self::new(
141            wire.raw_pty_lifecycle,
142            wire.semantic_readiness,
143            wire.structured_prompt,
144            wire.provider_session_identity,
145            wire.semantic_resume,
146        )
147        .map_err(serde::de::Error::custom)
148    }
149}
150
151#[derive(Clone, Copy, Debug, Eq, Error, PartialEq, Serialize, Deserialize)]
152#[serde(rename_all = "kebab-case")]
153pub enum ProviderRuntimePolicyError {
154    #[error("semantic provider capabilities require the raw PTY lifecycle")]
155    SemanticCapabilityRequiresRawPty,
156    #[error("structured prompts require semantic readiness")]
157    StructuredPromptRequiresReadiness,
158    #[error("semantic resume requires provider session identity")]
159    ResumeRequiresSessionIdentity,
160}
161
162#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
163#[serde(rename_all = "kebab-case")]
164pub enum ProviderRuntimeCapability {
165    RawPtyLifecycle,
166    SemanticReadiness,
167    StructuredPrompt,
168    ProviderSessionIdentity,
169    SemanticResume,
170}
171
172#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
173pub struct TerminalSize {
174    pub rows: u16,
175    pub columns: u16,
176}
177
178#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
179#[serde(rename_all = "kebab-case")]
180pub enum TerminalMouseProtocolEncoding {
181    #[default]
182    Default,
183    Utf8,
184    Sgr,
185}
186
187#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
188pub struct TerminalFrame {
189    pub sequence: u64,
190    pub size: TerminalSize,
191    pub cursor_row: u16,
192    pub cursor_column: u16,
193    pub contents: String,
194    pub formatted: Vec<u8>,
195    #[serde(default)]
196    pub scrollback_formatted: Vec<Vec<u8>>,
197    #[serde(default)]
198    pub alternate_screen: bool,
199    #[serde(default)]
200    pub mouse_protocol_enabled: bool,
201    #[serde(default)]
202    pub mouse_protocol_encoding: TerminalMouseProtocolEncoding,
203}
204
205pub const FOREGROUND_PROCESS_NAME_MAX_BYTES: usize = 512;
206
207#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
208#[serde(tag = "kind", rename_all = "kebab-case")]
209pub enum ForegroundProcessKind {
210    Agent { agent_id: AgentId },
211    Shell,
212    Other,
213}
214
215#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
216pub struct ForegroundProcess {
217    pub root_process_id: u32,
218    pub process_id: u32,
219    pub process_name: String,
220    pub kind: ForegroundProcessKind,
221}
222
223impl ForegroundProcess {
224    pub fn is_valid_for(&self, session_agent_id: &AgentId) -> bool {
225        self.root_process_id > 0
226            && self.process_id > 0
227            && !self.process_name.trim().is_empty()
228            && self.process_name.len() <= FOREGROUND_PROCESS_NAME_MAX_BYTES
229            && !self.process_name.chars().any(char::is_control)
230            && match &self.kind {
231                ForegroundProcessKind::Agent { agent_id } => agent_id == session_agent_id,
232                ForegroundProcessKind::Shell | ForegroundProcessKind::Other => true,
233            }
234    }
235}
236
237#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
238#[serde(rename_all = "kebab-case")]
239pub enum ForegroundAuthority {
240    #[default]
241    Unknown,
242    Confirmed,
243    Stale,
244}
245
246#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
247pub struct ForegroundSnapshot {
248    pub authority: ForegroundAuthority,
249    pub process: Option<ForegroundProcess>,
250    pub stale_reason: Option<String>,
251}
252
253impl TerminalSize {
254    pub fn is_valid(self) -> bool {
255        (1..=TERMINAL_ROWS_MAX).contains(&self.rows)
256            && (1..=TERMINAL_COLUMNS_MAX).contains(&self.columns)
257    }
258}
259
260#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
261#[serde(rename_all = "kebab-case")]
262pub enum TransportKind {
263    Pty,
264    Pipe,
265    Acp,
266}
267
268#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
269pub struct CommandEnvelope {
270    pub protocol_version: u16,
271    pub id: CommandId,
272    pub command: ControlCommand,
273}
274
275#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
276#[serde(tag = "kind", rename_all = "kebab-case")]
277pub enum ControlCommand {
278    Register {
279        instance_id: AgentInstanceId,
280        agent_id: AgentId,
281        transport: TransportKind,
282    },
283    Start {
284        instance_id: AgentInstanceId,
285        runtime_policy: ProviderRuntimePolicy,
286        request: StartRequest,
287    },
288    Stop {
289        instance_id: AgentInstanceId,
290        force: bool,
291    },
292    SendInput {
293        instance_id: AgentInstanceId,
294        action: InputAction,
295    },
296    Resize {
297        instance_id: AgentInstanceId,
298        size: TerminalSize,
299    },
300    RefreshForeground {
301        instance_id: AgentInstanceId,
302    },
303    ProbeCapabilities {
304        instance_id: AgentInstanceId,
305        request: CapabilityProbeRequest,
306    },
307    DiscoverHistory {
308        instance_id: AgentInstanceId,
309        query: HistoryQuery,
310    },
311    LoadHistory {
312        instance_id: AgentInstanceId,
313        candidate_id: String,
314    },
315    Resume {
316        instance_id: AgentInstanceId,
317        target: ResumeTarget,
318        runtime_policy: ProviderRuntimePolicy,
319        request: ResumeLaunchRequest,
320    },
321    ResolveInteraction {
322        instance_id: AgentInstanceId,
323        generation: SessionGeneration,
324        interaction_id: ProviderInteractionId,
325        response: ProviderInteractionResponse,
326    },
327    IngestProvider {
328        instance_id: AgentInstanceId,
329        generation: SessionGeneration,
330        source: ProviderSource,
331        source_sequence: u64,
332        events: Vec<ProviderEvent>,
333    },
334    Remove {
335        instance_id: AgentInstanceId,
336    },
337}
338
339impl ControlCommand {
340    pub fn instance_id(&self) -> AgentInstanceId {
341        match self {
342            Self::Register { instance_id, .. }
343            | Self::Start { instance_id, .. }
344            | Self::Stop { instance_id, .. }
345            | Self::SendInput { instance_id, .. }
346            | Self::Resize { instance_id, .. }
347            | Self::RefreshForeground { instance_id }
348            | Self::ProbeCapabilities { instance_id, .. }
349            | Self::DiscoverHistory { instance_id, .. }
350            | Self::LoadHistory { instance_id, .. }
351            | Self::Resume { instance_id, .. }
352            | Self::ResolveInteraction { instance_id, .. }
353            | Self::IngestProvider { instance_id, .. }
354            | Self::Remove { instance_id } => *instance_id,
355        }
356    }
357}
358
359#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
360pub struct EffectEnvelope {
361    pub protocol_version: u16,
362    pub operation_id: OperationId,
363    pub instance_id: AgentInstanceId,
364    pub generation: SessionGeneration,
365    pub effect: ControlEffect,
366}
367
368/// Route proof an effect executor must obtain immediately before a PTY write.
369///
370/// This is carried by the effect rather than inferred by a product shell so
371/// local, hosted, and future browser-facing executors enforce the same rule.
372#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
373#[serde(tag = "kind", rename_all = "kebab-case")]
374pub enum ForegroundRequirement {
375    /// Explicit terminal text and controls are direct user terminal input.
376    Any,
377    /// Semantic input must target the session's configured agent.
378    Agent { agent_id: AgentId },
379    /// Intentional shell syntax may be written only while a shell owns the PTY.
380    Shell,
381}
382
383#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
384#[serde(tag = "kind", rename_all = "kebab-case")]
385pub enum ControlEffect {
386    Spawn {
387        agent_id: AgentId,
388        transport: TransportKind,
389        runtime_policy: ProviderRuntimePolicy,
390        request: StartRequest,
391    },
392    Stop {
393        force: bool,
394    },
395    WriteInput {
396        input: PreparedInput,
397        required_foreground: ForegroundRequirement,
398    },
399    SubmitPrompt {
400        prompt: String,
401    },
402    Interrupt,
403    Resize {
404        size: TerminalSize,
405    },
406    ObserveForeground,
407    ProbeCapabilities {
408        agent_id: AgentId,
409        request: CapabilityProbeRequest,
410    },
411    DiscoverHistory {
412        agent_id: AgentId,
413        query: HistoryQuery,
414    },
415    LoadHistory {
416        agent_id: AgentId,
417        candidate_id: String,
418    },
419    AuthorizeResume {
420        agent_id: AgentId,
421        target: ResumeAuthorityTarget,
422        request: ResumeLaunchRequest,
423    },
424    SpawnResume {
425        agent_id: AgentId,
426        transport: TransportKind,
427        provider_session: ProviderSessionIdentity,
428        runtime_policy: ProviderRuntimePolicy,
429        request: ResumeLaunchRequest,
430    },
431    ResolveInteraction {
432        target: ProviderInteractionTarget,
433        response: ProviderInteractionResponse,
434    },
435}
436
437#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
438pub struct ObservationEnvelope {
439    pub protocol_version: u16,
440    pub operation_id: Option<OperationId>,
441    pub instance_id: AgentInstanceId,
442    pub generation: SessionGeneration,
443    pub observation: ControlObservation,
444}
445
446#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
447#[serde(tag = "kind", rename_all = "kebab-case")]
448pub enum ControlObservation {
449    Spawned {
450        process_id: Option<u32>,
451    },
452    SpawnFailed {
453        message: String,
454    },
455    ProcessExited {
456        exit_code: Option<i32>,
457        final_terminal: Option<TerminalFrame>,
458    },
459    StopCompleted {
460        forced: bool,
461        exit_code: Option<i32>,
462        final_terminal: Option<TerminalFrame>,
463    },
464    StopFailed {
465        message: String,
466    },
467    InputCompleted,
468    InputFailed {
469        message: String,
470    },
471    ResizeCompleted {
472        size: TerminalSize,
473    },
474    ResizeFailed {
475        message: String,
476    },
477    ForegroundObserved {
478        process: ForegroundProcess,
479    },
480    ForegroundFailed {
481        message: String,
482    },
483    CapabilitiesProbed {
484        session_option_models: Vec<CapabilityModelSummary>,
485    },
486    CapabilityProbeFailed {
487        failure: CapabilityProbeFailure,
488    },
489    HistoryDiscovered {
490        candidates: Vec<HistoryCandidateSummary>,
491    },
492    HistoryLoaded {
493        session: HistorySessionRecord,
494    },
495    HistoryFailed {
496        message: String,
497    },
498    ResumeAuthorized {
499        provider_session: ProviderSessionIdentity,
500    },
501    ResumeDenied {
502        reason: String,
503    },
504    ResumeFailed {
505        message: String,
506    },
507    InteractionResolutionCompleted {
508        interaction_id: ProviderInteractionId,
509    },
510    InteractionResolutionFailed {
511        interaction_id: ProviderInteractionId,
512        message: String,
513    },
514    TerminalFrame {
515        frame: TerminalFrame,
516    },
517    TerminalStale {
518        message: String,
519    },
520    ProviderEvent {
521        source: ProviderSource,
522        sequence: u64,
523        event: ProviderEvent,
524    },
525    ProviderGap {
526        source: ProviderSource,
527        source_sequence: u64,
528        missed: u64,
529    },
530}
531
532impl ControlObservation {
533    pub fn requires_operation_id(&self) -> bool {
534        !matches!(
535            self,
536            Self::ProcessExited { .. }
537                | Self::TerminalFrame { .. }
538                | Self::TerminalStale { .. }
539                | Self::ProviderEvent { .. }
540                | Self::ProviderGap { .. }
541        )
542    }
543}
544
545#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
546#[serde(tag = "kind", rename_all = "kebab-case")]
547pub enum SessionStatus {
548    Registered,
549    Starting,
550    Running,
551    Stopping,
552    Exited { exit_code: Option<i32> },
553    Failed { message: String },
554}
555
556#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
557pub struct SessionSnapshot {
558    pub instance_id: AgentInstanceId,
559    pub agent_id: AgentId,
560    pub transport: TransportKind,
561    pub generation: SessionGeneration,
562    pub status: SessionStatus,
563    pub pending_operation: Option<OperationId>,
564    pub pending_input: Option<PreparedInputKind>,
565    pub process_id: Option<u32>,
566    pub terminal_size: Option<TerminalSize>,
567    pub terminal_frame: Option<TerminalFrame>,
568    pub terminal_stale: Option<String>,
569    pub session_options: Option<SessionOptionSelection>,
570    pub capabilities: CapabilitySnapshot,
571    pub history: HistorySnapshot,
572    pub resume: ResumeSnapshot,
573    pub foreground: ForegroundSnapshot,
574    pub provider: ProviderSnapshot,
575}
576
577#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
578pub struct TokenUsage {
579    pub input_tokens: u64,
580    pub output_tokens: u64,
581    pub cache_read_tokens: u64,
582    pub cache_write_tokens: u64,
583    pub reasoning_tokens: u64,
584    pub context_window: Option<u64>,
585}
586
587#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
588pub struct ContextWindowUsage {
589    pub uncached_input_tokens: u64,
590    pub cache_read_tokens: u64,
591    pub cache_write_tokens: u64,
592    pub output_tokens: u64,
593    pub unattributed_tokens: u64,
594    pub used_tokens: u64,
595    pub capacity_tokens: u64,
596}
597
598impl ContextWindowUsage {
599    pub fn validate(&self) -> Result<(), ProviderEventValidationError> {
600        if self.capacity_tokens == 0 {
601            return Err(ProviderEventValidationError::ZeroContextWindowCapacity);
602        }
603        let segment_sum = self
604            .uncached_input_tokens
605            .checked_add(self.cache_read_tokens)
606            .and_then(|sum| sum.checked_add(self.cache_write_tokens))
607            .and_then(|sum| sum.checked_add(self.output_tokens))
608            .and_then(|sum| sum.checked_add(self.unattributed_tokens))
609            .ok_or(ProviderEventValidationError::ContextWindowSegmentsOverflow)?;
610        if segment_sum != self.used_tokens {
611            return Err(ProviderEventValidationError::ContextWindowSegmentsMismatch {
612                segment_sum,
613                used_tokens: self.used_tokens,
614            });
615        }
616        Ok(())
617    }
618}
619
620#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
621#[serde(rename_all = "kebab-case")]
622pub enum ProviderInteractionKind {
623    Approval,
624    Question,
625}
626
627#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
628#[serde(rename_all = "kebab-case")]
629pub enum ProviderInteractionOutcome {
630    Approved,
631    Answered,
632    Denied,
633    Interrupted,
634    TurnEnded,
635    Superseded,
636}
637
638#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
639#[serde(rename_all = "kebab-case")]
640pub enum ProviderInteractionResponseKind {
641    ApproveOnce,
642    Deny,
643    Answer,
644}
645
646#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
647#[serde(tag = "kind", rename_all = "kebab-case")]
648pub enum ProviderInteractionResponse {
649    ApproveOnce,
650    Deny,
651    Answer { text: String },
652}
653
654impl ProviderInteractionResponse {
655    pub fn kind(&self) -> ProviderInteractionResponseKind {
656        match self {
657            Self::ApproveOnce => ProviderInteractionResponseKind::ApproveOnce,
658            Self::Deny => ProviderInteractionResponseKind::Deny,
659            Self::Answer { .. } => ProviderInteractionResponseKind::Answer,
660        }
661    }
662
663    pub fn outcome(&self) -> ProviderInteractionOutcome {
664        match self {
665            Self::ApproveOnce => ProviderInteractionOutcome::Approved,
666            Self::Deny => ProviderInteractionOutcome::Denied,
667            Self::Answer { .. } => ProviderInteractionOutcome::Answered,
668        }
669    }
670
671    pub fn validate_for(
672        &self,
673        interaction_kind: ProviderInteractionKind,
674    ) -> Result<(), ProviderInteractionResponseError> {
675        match (interaction_kind, self) {
676            (ProviderInteractionKind::Approval, Self::ApproveOnce)
677            | (ProviderInteractionKind::Approval, Self::Deny)
678            | (ProviderInteractionKind::Question, Self::Deny) => Ok(()),
679            (ProviderInteractionKind::Question, Self::Answer { text }) => {
680                if text.trim().is_empty() {
681                    return Err(ProviderInteractionResponseError::EmptyAnswer);
682                }
683                let has_unsafe_control = text.chars().any(|character| {
684                    character.is_control() && !matches!(character, '\n' | '\r' | '\t')
685                });
686                if text.len() > PROVIDER_INTERACTION_RESPONSE_MAX_BYTES || has_unsafe_control {
687                    return Err(ProviderInteractionResponseError::InvalidAnswer {
688                        max: PROVIDER_INTERACTION_RESPONSE_MAX_BYTES,
689                    });
690                }
691                Ok(())
692            }
693            (ProviderInteractionKind::Approval, Self::Answer { .. }) => {
694                Err(ProviderInteractionResponseError::AnswerRequiresQuestion)
695            }
696            (ProviderInteractionKind::Question, Self::ApproveOnce) => {
697                Err(ProviderInteractionResponseError::ApprovalRequiresApproval)
698            }
699        }
700    }
701}
702
703#[derive(Clone, Debug, Eq, Error, PartialEq)]
704pub enum ProviderInteractionResponseError {
705    #[error("interaction answer is required")]
706    EmptyAnswer,
707    #[error("interaction answer contains controls or exceeds {max} bytes")]
708    InvalidAnswer { max: usize },
709    #[error("an answer response requires a question interaction")]
710    AnswerRequiresQuestion,
711    #[error("an approve-once response requires an approval interaction")]
712    ApprovalRequiresApproval,
713}
714
715#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
716#[serde(rename_all = "kebab-case")]
717pub enum ProviderSessionKey {
718    SessionId,
719    ConversationId,
720}
721
722#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
723pub struct ProviderSessionIdentity {
724    pub key: ProviderSessionKey,
725    pub id: String,
726    pub transcript_path: Option<String>,
727}
728
729impl ProviderSessionIdentity {
730    pub fn validate(&self) -> Result<(), ProviderEventValidationError> {
731        validate_required("provider session id", &self.id, PROVIDER_EVENT_ID_MAX_BYTES)?;
732        if self.id.starts_with('-') {
733            return Err(ProviderEventValidationError::InvalidField {
734                field: "provider session id",
735                max: PROVIDER_EVENT_ID_MAX_BYTES,
736            });
737        }
738        if let Some(path) = &self.transcript_path {
739            validate_required(
740                "provider transcript path",
741                path,
742                PROVIDER_SESSION_LOCATOR_MAX_BYTES,
743            )?;
744        }
745        Ok(())
746    }
747}
748
749#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
750pub struct ProviderSubagent {
751    pub source: ProviderSource,
752    pub provider_agent_id: String,
753    pub agent_type: Option<String>,
754    pub description: Option<String>,
755}
756
757#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
758#[serde(tag = "kind", rename_all = "kebab-case")]
759pub enum ProviderEvent {
760    SessionStarted {
761        session_id: String,
762        model: String,
763        tools: Vec<String>,
764    },
765    SessionIdentityObserved {
766        identity: ProviderSessionIdentity,
767    },
768    TurnStarted {
769        prompt: Option<String>,
770    },
771    WorkingObserved,
772    Text {
773        text: String,
774        is_delta: bool,
775    },
776    Thinking {
777        text: String,
778    },
779    ToolStarted {
780        id: String,
781        name: String,
782        input_json: String,
783        agent_id: Option<String>,
784    },
785    ToolCompleted {
786        id: String,
787        output: String,
788        is_error: bool,
789        duration_ms: Option<u64>,
790        agent_id: Option<String>,
791    },
792    TurnCompleted {
793        usage: TokenUsage,
794        is_cumulative: bool,
795    },
796    ContextWindowUsage {
797        usage: ContextWindowUsage,
798    },
799    TurnInterrupted,
800    SessionEnded {
801        result: String,
802        cost_usd: Option<String>,
803        is_error: bool,
804    },
805    Error {
806        message: String,
807    },
808    Ready,
809    InteractionRequested {
810        request_id: Option<String>,
811        interaction_kind: ProviderInteractionKind,
812        tool_name: String,
813        prompt: String,
814        agent_id: Option<String>,
815    },
816    InteractionResolved {
817        request_id: String,
818        outcome: ProviderInteractionOutcome,
819    },
820    SubagentStarted {
821        agent_id: String,
822        agent_type: Option<String>,
823        description: Option<String>,
824    },
825    SubagentStopped {
826        agent_id: String,
827    },
828    RateLimited {
829        limit_type: String,
830        resets_at: Option<String>,
831        usage_percent: Option<String>,
832        raw_message: String,
833    },
834}
835
836impl ProviderEvent {
837    pub fn validate_ingress(&self) -> Result<(), ProviderEventValidationError> {
838        match self {
839            Self::SessionStarted {
840                session_id,
841                model,
842                tools,
843            } => {
844                validate_required("session_id", session_id, PROVIDER_EVENT_ID_MAX_BYTES)?;
845                validate_identifier("model", model, PROVIDER_EVENT_ID_MAX_BYTES)?;
846                if tools.len() > PROVIDER_EVENT_TOOLS_MAX {
847                    return Err(ProviderEventValidationError::TooManyTools {
848                        count: tools.len(),
849                        max: PROVIDER_EVENT_TOOLS_MAX,
850                    });
851                }
852                for tool in tools {
853                    validate_required("tool", tool, PROVIDER_EVENT_ID_MAX_BYTES)?;
854                }
855            }
856            Self::SessionIdentityObserved { identity } => {
857                identity.validate()?;
858            }
859            Self::TurnStarted { prompt } => {
860                if let Some(prompt) = prompt {
861                    validate_text("prompt", prompt, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
862                }
863            }
864            Self::Text { text, .. } | Self::Thinking { text } => {
865                validate_text("text", text, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
866            }
867            Self::ToolStarted {
868                id,
869                name,
870                input_json,
871                agent_id,
872            } => {
873                validate_required("tool id", id, PROVIDER_EVENT_ID_MAX_BYTES)?;
874                validate_required("tool name", name, PROVIDER_EVENT_ID_MAX_BYTES)?;
875                validate_text("tool input", input_json, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
876                validate_optional_agent_id(agent_id)?;
877            }
878            Self::ToolCompleted {
879                id,
880                output,
881                agent_id,
882                ..
883            } => {
884                validate_required("tool id", id, PROVIDER_EVENT_ID_MAX_BYTES)?;
885                validate_text("tool output", output, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
886                validate_optional_agent_id(agent_id)?;
887            }
888            Self::SessionEnded {
889                result, cost_usd, ..
890            } => {
891                validate_text("session result", result, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
892                if let Some(cost) = cost_usd {
893                    validate_identifier("cost", cost, PROVIDER_EVENT_ID_MAX_BYTES)?;
894                }
895            }
896            Self::Error { message } => {
897                validate_required_text("error", message, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
898            }
899            Self::InteractionRequested {
900                request_id,
901                interaction_kind,
902                tool_name,
903                prompt,
904                agent_id,
905            } => {
906                if let Some(request_id) = request_id {
907                    validate_required(
908                        "interaction request id",
909                        request_id,
910                        PROVIDER_EVENT_ID_MAX_BYTES,
911                    )?;
912                }
913                validate_required("interaction tool", tool_name, PROVIDER_EVENT_ID_MAX_BYTES)?;
914                if *interaction_kind == ProviderInteractionKind::Question {
915                    validate_required_text(
916                        "interaction prompt",
917                        prompt,
918                        PROVIDER_EVENT_TEXT_MAX_BYTES,
919                    )?;
920                } else {
921                    validate_text("interaction prompt", prompt, PROVIDER_EVENT_TEXT_MAX_BYTES)?;
922                }
923                validate_optional_agent_id(agent_id)?;
924            }
925            Self::InteractionResolved {
926                request_id,
927                outcome,
928            } => {
929                validate_required(
930                    "interaction request id",
931                    request_id,
932                    PROVIDER_EVENT_ID_MAX_BYTES,
933                )?;
934                if !matches!(
935                    outcome,
936                    ProviderInteractionOutcome::Approved | ProviderInteractionOutcome::Denied
937                ) {
938                    return Err(
939                        ProviderEventValidationError::InvalidInteractionResolutionOutcome {
940                            outcome: *outcome,
941                        },
942                    );
943                }
944            }
945            Self::SubagentStarted {
946                agent_id,
947                agent_type,
948                description,
949            } => {
950                validate_required("subagent id", agent_id, PROVIDER_EVENT_ID_MAX_BYTES)?;
951                if let Some(agent_type) = agent_type {
952                    validate_identifier("subagent type", agent_type, PROVIDER_EVENT_ID_MAX_BYTES)?;
953                }
954                if let Some(description) = description {
955                    validate_text(
956                        "subagent description",
957                        description,
958                        PROVIDER_EVENT_TEXT_MAX_BYTES,
959                    )?;
960                }
961            }
962            Self::SubagentStopped { agent_id } => {
963                validate_required("subagent id", agent_id, PROVIDER_EVENT_ID_MAX_BYTES)?;
964            }
965            Self::RateLimited {
966                limit_type,
967                resets_at,
968                usage_percent,
969                raw_message,
970            } => {
971                validate_required("limit type", limit_type, PROVIDER_EVENT_ID_MAX_BYTES)?;
972                for (field, value) in [
973                    ("reset time", resets_at.as_deref()),
974                    ("usage percent", usage_percent.as_deref()),
975                ] {
976                    if let Some(value) = value {
977                        validate_identifier(field, value, PROVIDER_EVENT_ID_MAX_BYTES)?;
978                    }
979                }
980                validate_text(
981                    "rate limit message",
982                    raw_message,
983                    PROVIDER_EVENT_TEXT_MAX_BYTES,
984                )?;
985            }
986            Self::WorkingObserved
987            | Self::TurnCompleted { .. }
988            | Self::TurnInterrupted
989            | Self::Ready => {}
990            Self::ContextWindowUsage { usage } => usage.validate()?,
991        }
992        Ok(())
993    }
994}
995
996fn validate_required(
997    field: &'static str,
998    value: &str,
999    max: usize,
1000) -> Result<(), ProviderEventValidationError> {
1001    if value.trim().is_empty() {
1002        return Err(ProviderEventValidationError::Empty { field });
1003    }
1004    validate_identifier(field, value, max)
1005}
1006
1007fn validate_required_text(
1008    field: &'static str,
1009    value: &str,
1010    max: usize,
1011) -> Result<(), ProviderEventValidationError> {
1012    if value.trim().is_empty() {
1013        return Err(ProviderEventValidationError::Empty { field });
1014    }
1015    validate_text(field, value, max)
1016}
1017
1018fn validate_identifier(
1019    field: &'static str,
1020    value: &str,
1021    max: usize,
1022) -> Result<(), ProviderEventValidationError> {
1023    if value.len() > max || value.chars().any(char::is_control) {
1024        return Err(ProviderEventValidationError::InvalidField { field, max });
1025    }
1026    Ok(())
1027}
1028
1029fn validate_optional_agent_id(
1030    agent_id: &Option<String>,
1031) -> Result<(), ProviderEventValidationError> {
1032    if let Some(agent_id) = agent_id {
1033        validate_required("provider agent id", agent_id, PROVIDER_EVENT_ID_MAX_BYTES)?;
1034    }
1035    Ok(())
1036}
1037
1038fn validate_text(
1039    field: &'static str,
1040    value: &str,
1041    max: usize,
1042) -> Result<(), ProviderEventValidationError> {
1043    let has_unsafe_control = value
1044        .chars()
1045        .any(|character| character.is_control() && !matches!(character, '\n' | '\r' | '\t'));
1046    if value.len() > max || has_unsafe_control {
1047        return Err(ProviderEventValidationError::InvalidField { field, max });
1048    }
1049    Ok(())
1050}
1051
1052#[derive(Clone, Debug, Error, Eq, PartialEq)]
1053pub enum ProviderEventValidationError {
1054    #[error("provider event field '{field}' is required")]
1055    Empty { field: &'static str },
1056    #[error("provider event field '{field}' contains controls or exceeds {max} bytes")]
1057    InvalidField { field: &'static str, max: usize },
1058    #[error("provider event tool count {count} exceeds {max}")]
1059    TooManyTools { count: usize, max: usize },
1060    #[error("provider interaction resolution outcome {outcome:?} is not exact")]
1061    InvalidInteractionResolutionOutcome { outcome: ProviderInteractionOutcome },
1062    #[error("context-window capacity must be non-zero")]
1063    ZeroContextWindowCapacity,
1064    #[error("context-window token segments overflow u64")]
1065    ContextWindowSegmentsOverflow,
1066    #[error("context-window token segments sum to {segment_sum}, not used_tokens {used_tokens}")]
1067    ContextWindowSegmentsMismatch { segment_sum: u64, used_tokens: u64 },
1068}
1069
1070#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
1071pub struct ProviderSource {
1072    pub family: AdapterFamily,
1073    pub binding: AdapterBinding,
1074}
1075
1076#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1077pub struct ProviderSourceCursor {
1078    pub source: ProviderSource,
1079    pub sequence: u64,
1080    pub gap_count: u64,
1081    pub stale: bool,
1082}
1083
1084#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
1085#[serde(transparent)]
1086pub struct ProviderInteractionId(pub u64);
1087
1088#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1089pub struct ProviderInteractionTarget {
1090    pub interaction_id: ProviderInteractionId,
1091    pub source: ProviderSource,
1092    pub provider_request_id: Option<String>,
1093    pub interaction_kind: ProviderInteractionKind,
1094    pub tool_name: String,
1095    pub agent_id: Option<String>,
1096}
1097
1098#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1099#[serde(tag = "kind", rename_all = "kebab-case")]
1100pub enum ProviderInteractionStatus {
1101    Pending,
1102    Resolving {
1103        operation_id: OperationId,
1104        response_kind: ProviderInteractionResponseKind,
1105    },
1106    Resolved {
1107        outcome: ProviderInteractionOutcome,
1108    },
1109}
1110
1111#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1112pub struct ProviderInteraction {
1113    pub id: ProviderInteractionId,
1114    pub source: ProviderSource,
1115    pub provider_request_id: Option<String>,
1116    pub interaction_kind: ProviderInteractionKind,
1117    pub tool_name: String,
1118    pub prompt: String,
1119    pub agent_id: Option<String>,
1120    pub resume_lead_activity: Option<ProviderActivity>,
1121    pub status: ProviderInteractionStatus,
1122}
1123
1124#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
1125#[serde(rename_all = "kebab-case")]
1126pub enum ProviderActivity {
1127    #[default]
1128    Idle,
1129    Working,
1130    WaitingForInput,
1131    Blocked,
1132}
1133
1134#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1135pub struct ActiveProviderTool {
1136    pub id: String,
1137    pub name: String,
1138    pub input_json: String,
1139}
1140
1141#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
1142pub struct ProviderSnapshot {
1143    pub sequence: u64,
1144    pub session: Option<ProviderSessionIdentity>,
1145    pub model: Option<String>,
1146    pub tools: Vec<String>,
1147    pub completed_turns: u64,
1148    pub usage: TokenUsage,
1149    pub lead_activity: ProviderActivity,
1150    pub activity: ProviderActivity,
1151    pub current_prompt: Option<String>,
1152    pub active_tools: Vec<ActiveProviderTool>,
1153    pub interactions: Vec<ProviderInteraction>,
1154    pub subagents: Vec<ProviderSubagent>,
1155    pub sources: Vec<ProviderSourceCursor>,
1156    pub last_event: Option<ProviderEvent>,
1157    pub gap_count: u64,
1158    pub stale: bool,
1159}
1160
1161#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1162pub struct ControlSnapshot {
1163    pub protocol_version: u16,
1164    pub revision: u64,
1165    pub health: ControlHealth,
1166    pub sessions: Vec<SessionSnapshot>,
1167}
1168
1169#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
1170pub struct ControlHealth {
1171    pub operation_id_exhausted: bool,
1172    pub event_sequence_exhausted: bool,
1173    pub revision_exhausted: bool,
1174    pub provider_sequence_exhausted_sessions: u32,
1175    pub retained_instance_identities: u32,
1176    pub retained_instance_identity_capacity: u32,
1177}
1178
1179impl Default for ControlHealth {
1180    fn default() -> Self {
1181        Self {
1182            operation_id_exhausted: false,
1183            event_sequence_exhausted: false,
1184            revision_exhausted: false,
1185            provider_sequence_exhausted_sessions: 0,
1186            retained_instance_identities: 0,
1187            retained_instance_identity_capacity: CONTROL_INSTANCE_IDENTITIES_CAPACITY,
1188        }
1189    }
1190}
1191
1192#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1193pub struct ControlEvent {
1194    pub protocol_version: u16,
1195    pub sequence: u64,
1196    pub command_id: Option<CommandId>,
1197    pub instance_id: AgentInstanceId,
1198    pub generation: SessionGeneration,
1199    pub event: ControlEventKind,
1200}
1201
1202#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
1203#[serde(tag = "kind", rename_all = "kebab-case")]
1204pub enum ControlEventKind {
1205    CommandRejected {
1206        message: String,
1207    },
1208    Registered,
1209    StartRequested {
1210        operation_id: OperationId,
1211    },
1212    Running {
1213        process_id: Option<u32>,
1214    },
1215    StopRequested {
1216        operation_id: OperationId,
1217        force: bool,
1218    },
1219    InputRequested {
1220        operation_id: OperationId,
1221        input_kind: PreparedInputKind,
1222    },
1223    InputCompleted {
1224        input_kind: PreparedInputKind,
1225    },
1226    InputFailed {
1227        input_kind: PreparedInputKind,
1228        message: String,
1229    },
1230    ResizeRequested {
1231        operation_id: OperationId,
1232        size: TerminalSize,
1233    },
1234    Resized {
1235        size: TerminalSize,
1236    },
1237    ResizeFailed {
1238        message: String,
1239    },
1240    ForegroundRefreshRequested {
1241        operation_id: OperationId,
1242    },
1243    ForegroundObserved {
1244        process: ForegroundProcess,
1245    },
1246    ForegroundFailed {
1247        message: String,
1248    },
1249    CapabilityProbeRequested {
1250        operation_id: OperationId,
1251    },
1252    CapabilitiesProbed {
1253        count: usize,
1254    },
1255    CapabilityProbeFailed {
1256        failure: CapabilityProbeFailure,
1257    },
1258    HistoryRequested {
1259        operation_id: OperationId,
1260        operation: HistoryOperation,
1261    },
1262    HistoryDiscovered {
1263        count: usize,
1264    },
1265    HistoryLoaded {
1266        session_id: String,
1267    },
1268    HistoryFailed {
1269        message: String,
1270    },
1271    ResumeRequested {
1272        operation_id: OperationId,
1273        target: ResumeTarget,
1274    },
1275    ResumeAuthorized {
1276        session: ResumeSessionSummary,
1277    },
1278    Resumed {
1279        session: ResumeSessionSummary,
1280        process_id: Option<u32>,
1281    },
1282    ResumeDenied {
1283        reason: String,
1284    },
1285    ResumeFailed {
1286        message: String,
1287    },
1288    TerminalStale {
1289        message: String,
1290    },
1291    ProviderEvent {
1292        sequence: u64,
1293        source: ProviderSource,
1294        source_sequence: u64,
1295        event: ProviderEvent,
1296    },
1297    ProviderGap {
1298        sequence: u64,
1299        source: ProviderSource,
1300        source_sequence: u64,
1301        missed: u64,
1302    },
1303    InteractionRequested {
1304        interaction: ProviderInteraction,
1305    },
1306    InteractionResolutionRequested {
1307        operation_id: OperationId,
1308        interaction_id: ProviderInteractionId,
1309        response_kind: ProviderInteractionResponseKind,
1310    },
1311    InteractionResolutionFailed {
1312        interaction_id: ProviderInteractionId,
1313        message: String,
1314    },
1315    InteractionResolved {
1316        interaction_id: ProviderInteractionId,
1317        outcome: ProviderInteractionOutcome,
1318    },
1319    Exited {
1320        exit_code: Option<i32>,
1321        forced: bool,
1322    },
1323    Failed {
1324        message: String,
1325    },
1326    Removed,
1327    ObservationIgnored {
1328        reason: ObservationIgnoredReason,
1329    },
1330}
1331
1332#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
1333#[serde(rename_all = "kebab-case")]
1334pub enum ObservationIgnoredReason {
1335    UnsupportedProtocolVersion,
1336    UnknownInstance,
1337    StaleGeneration,
1338    GenerationExhausted,
1339    MissingOperation,
1340    OperationMismatch,
1341    InvalidState,
1342    StaleTerminalFrame,
1343    StaleProviderEvent,
1344    InvalidForegroundObservation,
1345    InvalidCapabilityObservation,
1346    InvalidHistoryObservation,
1347    InvalidResumeObservation,
1348    InvalidInteractionObservation,
1349    ProviderRuntimePolicyDenied {
1350        capability: ProviderRuntimeCapability,
1351    },
1352}
1353
1354#[derive(Clone, Debug, Eq, Error, PartialEq, Serialize, Deserialize)]
1355#[serde(tag = "kind", rename_all = "kebab-case")]
1356pub enum ControlError {
1357    #[error("control protocol version {actual} is unsupported; expected {expected}")]
1358    UnsupportedProtocolVersion { expected: u16, actual: u16 },
1359    #[error("agent instance {instance_id:?} is already registered")]
1360    DuplicateInstance { instance_id: AgentInstanceId },
1361    #[error(
1362        "cannot register agent instance {instance_id:?}: live session capacity {max} is exhausted"
1363    )]
1364    SessionCapacityExceeded {
1365        instance_id: AgentInstanceId,
1366        max: usize,
1367    },
1368    #[error(
1369        "cannot register agent instance {instance_id:?}: retained identity capacity {max} is exhausted"
1370    )]
1371    InstanceIdentityCapacityExceeded {
1372        instance_id: AgentInstanceId,
1373        max: usize,
1374    },
1375    #[error("agent instance {instance_id:?} is not registered")]
1376    UnknownInstance { instance_id: AgentInstanceId },
1377    #[error("agent instance {instance_id:?} exhausted session generation {generation:?}")]
1378    GenerationExhausted {
1379        instance_id: AgentInstanceId,
1380        generation: SessionGeneration,
1381    },
1382    #[error("control operation identifiers are exhausted")]
1383    OperationIdExhausted,
1384    #[error("control event sequences are exhausted")]
1385    EventSequenceExhausted,
1386    #[error("control snapshot revisions are exhausted")]
1387    RevisionExhausted,
1388    #[error(
1389        "agent instance {instance_id:?} generation {generation:?} exhausted provider event sequences"
1390    )]
1391    ProviderSequenceExhausted {
1392        instance_id: AgentInstanceId,
1393        generation: SessionGeneration,
1394    },
1395    #[error(
1396        "agent instance {instance_id:?} generation {generation:?} exhausted source sequence for {provider_source:?}"
1397    )]
1398    ProviderSourceSequenceExhausted {
1399        instance_id: AgentInstanceId,
1400        generation: SessionGeneration,
1401        provider_source: ProviderSource,
1402    },
1403    #[error("agent instance {instance_id:?} already has pending operation {operation_id:?}")]
1404    OperationPending {
1405        instance_id: AgentInstanceId,
1406        operation_id: OperationId,
1407    },
1408    #[error("agent input was rejected: {error}")]
1409    InputRejected { error: InputPrepareError },
1410    #[error("provider runtime policy is invalid: {error}")]
1411    InvalidProviderRuntimePolicy { error: ProviderRuntimePolicyError },
1412    #[error("provider runtime capability {capability:?} is not admitted")]
1413    ProviderRuntimePolicyDenied {
1414        capability: ProviderRuntimeCapability,
1415    },
1416    #[error("terminal size is outside the supported bounded range")]
1417    InvalidTerminalSize,
1418    #[error("working directory is empty, too large, or contains a NUL byte")]
1419    InvalidWorkingDirectory,
1420    #[error("pipe transport requires a non-empty initial prompt")]
1421    MissingInitialPrompt,
1422    #[error("session options are invalid: {message}")]
1423    InvalidSessionOptions { message: String },
1424    #[error("capability probe request is invalid: {message}")]
1425    InvalidCapabilityProbeRequest { message: String },
1426    #[error("capability probe operation {operation_id:?} is already pending")]
1427    CapabilityProbeOperationPending { operation_id: OperationId },
1428    #[error("capability probe already settled for this agent instance")]
1429    CapabilityProbeSettled,
1430    #[error("history request is invalid: {message}")]
1431    InvalidHistoryRequest { message: String },
1432    #[error("history operation {operation_id:?} is already pending")]
1433    HistoryOperationPending { operation_id: OperationId },
1434    #[error("history candidate is not present in the current discovery snapshot")]
1435    UnknownHistoryCandidate,
1436    #[error("resume request is invalid: {message}")]
1437    InvalidResumeRequest { message: String },
1438    #[error("resume requires a canonical provider session identity")]
1439    MissingProviderSession,
1440    #[error("resume history candidate must be the currently loaded candidate")]
1441    HistoryCandidateNotLoaded,
1442    #[error("transport {transport:?} does not support {action}")]
1443    UnsupportedTransportOperation {
1444        transport: TransportKind,
1445        action: String,
1446    },
1447    #[error("agent instance {instance_id:?} cannot {action} while in state {status:?}")]
1448    InvalidTransition {
1449        instance_id: AgentInstanceId,
1450        action: String,
1451        status: SessionStatus,
1452    },
1453    #[error("provider ingress generation {actual:?} is stale; expected {expected:?}")]
1454    StaleProviderGeneration {
1455        expected: SessionGeneration,
1456        actual: SessionGeneration,
1457    },
1458    #[error("provider ingress source sequence must be greater than the current sequence")]
1459    StaleProviderSequence,
1460    #[error("provider ingress batch must contain between 1 and {max} events")]
1461    InvalidProviderBatch { max: usize },
1462    #[error("invalid provider ingress event: {message}")]
1463    InvalidProviderEvent { message: String },
1464    #[error("provider interaction generation {actual:?} is stale; expected {expected:?}")]
1465    StaleProviderInteractionGeneration {
1466        expected: SessionGeneration,
1467        actual: SessionGeneration,
1468    },
1469    #[error("provider interaction {interaction_id:?} is unknown")]
1470    UnknownProviderInteraction {
1471        interaction_id: ProviderInteractionId,
1472    },
1473    #[error("provider interaction {interaction_id:?} is not pending")]
1474    ProviderInteractionNotPending {
1475        interaction_id: ProviderInteractionId,
1476    },
1477    #[error("provider interaction response is invalid: {message}")]
1478    InvalidProviderInteractionResponse { message: String },
1479}
1480
1481impl Default for ControlSnapshot {
1482    fn default() -> Self {
1483        Self {
1484            protocol_version: CONTROL_PROTOCOL_VERSION,
1485            revision: 0,
1486            health: ControlHealth::default(),
1487            sessions: Vec::new(),
1488        }
1489    }
1490}
1491
1492#[cfg(test)]
1493mod tests {
1494    use crate::AgentId;
1495    use super::{
1496        ContextWindowUsage, ForegroundProcess, ForegroundProcessKind, ProviderEvent,
1497        ProviderEventValidationError,
1498        ProviderInteractionKind, ProviderInteractionOutcome, ProviderInteractionResponse,
1499        ProviderInteractionResponseError, ProviderRuntimeCapability, ProviderRuntimePolicy,
1500        ProviderRuntimePolicyError,
1501        ProviderSessionIdentity, ProviderSessionKey, TerminalFrame, TerminalMouseProtocolEncoding,
1502        FOREGROUND_PROCESS_NAME_MAX_BYTES, PROVIDER_INTERACTION_RESPONSE_MAX_BYTES,
1503    };
1504
1505    #[test]
1506    fn context_window_usage_ingress_requires_exact_bounded_segments() {
1507        let event = |usage| ProviderEvent::ContextWindowUsage { usage };
1508        let valid = ContextWindowUsage {
1509            uncached_input_tokens: 70,
1510            cache_read_tokens: 20,
1511            cache_write_tokens: 0,
1512            output_tokens: 10,
1513            unattributed_tokens: 5,
1514            used_tokens: 105,
1515            capacity_tokens: 100,
1516        };
1517        assert_eq!(event(valid).validate_ingress(), Ok(()));
1518        assert_eq!(
1519            event(ContextWindowUsage { capacity_tokens: 0, ..valid }).validate_ingress(),
1520            Err(ProviderEventValidationError::ZeroContextWindowCapacity)
1521        );
1522        assert_eq!(
1523            event(ContextWindowUsage { used_tokens: 104, ..valid }).validate_ingress(),
1524            Err(ProviderEventValidationError::ContextWindowSegmentsMismatch {
1525                segment_sum: 105,
1526                used_tokens: 104,
1527            })
1528        );
1529        assert_eq!(
1530            event(ContextWindowUsage {
1531                uncached_input_tokens: u64::MAX,
1532                cache_read_tokens: 1,
1533                cache_write_tokens: 0,
1534                output_tokens: 0,
1535                unattributed_tokens: 0,
1536                used_tokens: u64::MAX,
1537                capacity_tokens: 1,
1538            })
1539            .validate_ingress(),
1540            Err(ProviderEventValidationError::ContextWindowSegmentsOverflow)
1541        );
1542    }
1543
1544    #[test]
1545    fn provider_runtime_policy_enforces_semantic_invariants() {
1546        let raw = ProviderRuntimePolicy::raw_pty();
1547        assert!(raw.admits(ProviderRuntimeCapability::RawPtyLifecycle));
1548        assert!(!raw.admits(ProviderRuntimeCapability::SemanticReadiness));
1549        assert_eq!(raw.validate(), Ok(()));
1550
1551        assert_eq!(
1552            ProviderRuntimePolicy::new(false, true, false, false, false),
1553            Err(ProviderRuntimePolicyError::SemanticCapabilityRequiresRawPty),
1554        );
1555        assert_eq!(
1556            ProviderRuntimePolicy::new(true, false, true, false, false),
1557            Err(ProviderRuntimePolicyError::StructuredPromptRequiresReadiness),
1558        );
1559        assert_eq!(
1560            ProviderRuntimePolicy::new(true, true, true, false, true),
1561            Err(ProviderRuntimePolicyError::ResumeRequiresSessionIdentity),
1562        );
1563        assert!(ProviderRuntimePolicy::new(true, true, true, true, true).is_ok());
1564    }
1565
1566    #[test]
1567    fn provider_runtime_policy_serde_requires_every_field_and_revalidates() {
1568        let raw = ProviderRuntimePolicy::raw_pty();
1569        let encoded = serde_json::to_string(&raw).unwrap();
1570        assert_eq!(
1571            encoded,
1572            r#"{"raw_pty_lifecycle":true,"semantic_readiness":false,"structured_prompt":false,"provider_session_identity":false,"semantic_resume":false}"#,
1573        );
1574        assert_eq!(
1575            serde_json::from_str::<ProviderRuntimePolicy>(&encoded).unwrap(),
1576            raw,
1577        );
1578        assert!(serde_json::from_str::<ProviderRuntimePolicy>(
1579            r#"{"raw_pty_lifecycle":true,"semantic_readiness":false,"structured_prompt":false,"provider_session_identity":false}"#,
1580        )
1581        .is_err());
1582        assert!(serde_json::from_str::<ProviderRuntimePolicy>(
1583            r#"{"raw_pty_lifecycle":true,"semantic_readiness":false,"structured_prompt":true,"provider_session_identity":false,"semantic_resume":false}"#,
1584        )
1585        .is_err());
1586        assert!(serde_json::from_str::<super::ControlCommand>(
1587            r#"{"kind":"start","instance_id":1,"request":{"working_directory":"C:\\repo","terminal_size":{"rows":24,"columns":80}}}"#,
1588        )
1589        .is_err());
1590        assert!(serde_json::from_str::<super::ControlEffect>(
1591            r#"{"kind":"spawn","agent_id":"claude","transport":"pty","request":{"working_directory":"C:\\repo","terminal_size":{"rows":24,"columns":80}}}"#,
1592        )
1593        .is_err());
1594    }
1595
1596    #[test]
1597    fn foreground_process_is_bounded_and_agent_bound() {
1598        let claude = AgentId::new("claude").unwrap();
1599        let process = ForegroundProcess {
1600            root_process_id: 1,
1601            process_id: 2,
1602            process_name: "claude".to_owned(),
1603            kind: ForegroundProcessKind::Agent {
1604                agent_id: claude.clone(),
1605            },
1606        };
1607        assert!(process.is_valid_for(&claude));
1608        assert!(!process.is_valid_for(&AgentId::new("codex").unwrap()));
1609        assert!(!ForegroundProcess {
1610            process_name: "x".repeat(FOREGROUND_PROCESS_NAME_MAX_BYTES + 1),
1611            ..process
1612        }
1613        .is_valid_for(&claude));
1614    }
1615
1616    #[test]
1617    fn terminal_frame_metadata_defaults_for_older_serialized_frames() {
1618        let frame: TerminalFrame = serde_json::from_str(
1619            r#"{"sequence":1,"size":{"rows":24,"columns":80},"cursor_row":0,"cursor_column":0,"contents":"ready","formatted":[114]}"#,
1620        )
1621        .expect("legacy terminal frame");
1622
1623        assert!(frame.scrollback_formatted.is_empty());
1624        assert!(!frame.alternate_screen);
1625        assert!(!frame.mouse_protocol_enabled);
1626        assert_eq!(frame.mouse_protocol_encoding, TerminalMouseProtocolEncoding::Default);
1627    }
1628
1629    #[test]
1630    fn provider_interactions_require_bounded_identity_and_question_payloads() {
1631        let question = ProviderEvent::InteractionRequested {
1632            request_id: Some("question-1".to_owned()),
1633            interaction_kind: ProviderInteractionKind::Question,
1634            tool_name: "AskUserQuestion".to_owned(),
1635            prompt: "{\"question\":\"Continue?\"}".to_owned(),
1636            agent_id: Some("child-1".to_owned()),
1637        };
1638        assert_eq!(question.validate_ingress(), Ok(()));
1639
1640        assert!(matches!(
1641            ProviderEvent::InteractionRequested {
1642                request_id: Some("bad\nrequest".to_owned()),
1643                interaction_kind: ProviderInteractionKind::Approval,
1644                tool_name: "shell".to_owned(),
1645                prompt: String::new(),
1646                agent_id: None,
1647            }
1648            .validate_ingress(),
1649            Err(ProviderEventValidationError::InvalidField {
1650                field: "interaction request id",
1651                ..
1652            })
1653        ));
1654        assert!(matches!(
1655            ProviderEvent::InteractionRequested {
1656                request_id: None,
1657                interaction_kind: ProviderInteractionKind::Question,
1658                tool_name: "AskUserQuestion".to_owned(),
1659                prompt: String::new(),
1660                agent_id: None,
1661            }
1662            .validate_ingress(),
1663            Err(ProviderEventValidationError::Empty {
1664                field: "interaction prompt"
1665            })
1666        ));
1667
1668        for outcome in [
1669            ProviderInteractionOutcome::Approved,
1670            ProviderInteractionOutcome::Denied,
1671        ] {
1672            assert_eq!(
1673                ProviderEvent::InteractionResolved {
1674                    request_id: "approval-1".to_owned(),
1675                    outcome,
1676                }
1677                .validate_ingress(),
1678                Ok(())
1679            );
1680        }
1681        assert!(matches!(
1682            ProviderEvent::InteractionResolved {
1683                request_id: "bad\nrequest".to_owned(),
1684                outcome: ProviderInteractionOutcome::Approved,
1685            }
1686            .validate_ingress(),
1687            Err(ProviderEventValidationError::InvalidField {
1688                field: "interaction request id",
1689                ..
1690            })
1691        ));
1692        assert_eq!(
1693            ProviderEvent::InteractionResolved {
1694                request_id: "approval-1".to_owned(),
1695                outcome: ProviderInteractionOutcome::TurnEnded,
1696            }
1697            .validate_ingress(),
1698            Err(
1699                ProviderEventValidationError::InvalidInteractionResolutionOutcome {
1700                    outcome: ProviderInteractionOutcome::TurnEnded,
1701                }
1702            )
1703        );
1704    }
1705
1706    #[test]
1707    fn provider_interaction_responses_are_kind_checked_and_bounded() {
1708        assert_eq!(
1709            ProviderInteractionResponse::ApproveOnce
1710                .validate_for(ProviderInteractionKind::Approval),
1711            Ok(())
1712        );
1713        assert_eq!(
1714            ProviderInteractionResponse::Deny.validate_for(ProviderInteractionKind::Question),
1715            Ok(())
1716        );
1717        assert_eq!(
1718            ProviderInteractionResponse::Answer {
1719                text: "continue".to_owned(),
1720            }
1721            .validate_for(ProviderInteractionKind::Question),
1722            Ok(())
1723        );
1724        assert_eq!(
1725            ProviderInteractionResponse::ApproveOnce
1726                .validate_for(ProviderInteractionKind::Question),
1727            Err(ProviderInteractionResponseError::ApprovalRequiresApproval)
1728        );
1729        assert_eq!(
1730            ProviderInteractionResponse::Answer {
1731                text: String::new(),
1732            }
1733            .validate_for(ProviderInteractionKind::Question),
1734            Err(ProviderInteractionResponseError::EmptyAnswer)
1735        );
1736        assert_eq!(
1737            ProviderInteractionResponse::Answer {
1738                text: "x".repeat(PROVIDER_INTERACTION_RESPONSE_MAX_BYTES + 1),
1739            }
1740            .validate_for(ProviderInteractionKind::Question),
1741            Err(ProviderInteractionResponseError::InvalidAnswer {
1742                max: PROVIDER_INTERACTION_RESPONSE_MAX_BYTES,
1743            })
1744        );
1745    }
1746
1747    #[test]
1748    fn provider_ingress_allows_multiline_text_but_rejects_control_bytes() {
1749        ProviderEvent::Text {
1750            text: "first line\n\tsecond line".to_owned(),
1751            is_delta: false,
1752        }
1753        .validate_ingress()
1754        .unwrap();
1755
1756        assert!(matches!(
1757            ProviderEvent::Text {
1758                text: "unsafe\u{0000}text".to_owned(),
1759                is_delta: false,
1760            }
1761            .validate_ingress(),
1762            Err(ProviderEventValidationError::InvalidField { field: "text", .. })
1763        ));
1764        assert!(ProviderEvent::SessionStarted {
1765            session_id: "session\nother".to_owned(),
1766            model: "model".to_owned(),
1767            tools: Vec::new(),
1768        }
1769        .validate_ingress()
1770        .is_err());
1771    }
1772
1773    #[test]
1774    fn provider_session_identity_is_typed_and_bounded_at_ingress() {
1775        let valid = ProviderEvent::SessionIdentityObserved {
1776            identity: ProviderSessionIdentity {
1777                key: ProviderSessionKey::ConversationId,
1778                id: "conversation-1".to_owned(),
1779                transcript_path: Some("C:/sessions/conversation-1.jsonl".to_owned()),
1780            },
1781        };
1782        assert_eq!(valid.validate_ingress(), Ok(()));
1783
1784        for identity in [
1785            ProviderSessionIdentity {
1786                key: ProviderSessionKey::SessionId,
1787                id: "--help".to_owned(),
1788                transcript_path: None,
1789            },
1790            ProviderSessionIdentity {
1791                key: ProviderSessionKey::SessionId,
1792                id: "session-1".to_owned(),
1793                transcript_path: Some("bad\npath".to_owned()),
1794            },
1795        ] {
1796            assert!(ProviderEvent::SessionIdentityObserved { identity }
1797                .validate_ingress()
1798                .is_err());
1799        }
1800    }
1801}