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#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
373#[serde(tag = "kind", rename_all = "kebab-case")]
374pub enum ForegroundRequirement {
375 Any,
377 Agent { agent_id: AgentId },
379 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}