返回 CodeWhale
protocol_parity.rs
根目录 / crates / tui / src / core / protocol_parity.rs
1 //! Compile-enforced parity between the engine's internal `Op` / `Event` and
2 //! the protocol's `Op` / `EventMsg` (core/protocol extraction spec, Phase A1).
3 //!
4 //! Every function here is one exhaustive `match` with **no wildcard arm**.
5 //! Adding an engine `Event` or `Op` variant without a protocol twin fails to
6 //! compile right here — that is the `protocol_covers_engine_events` /
7 //! `protocol_covers_engine_ops` guard from spec §7. The `#[cfg(test)]`
8 //! block below only proves the projections agree with the protocol's own
9 //! wire-tag tables and that this file keeps its no-wildcard discipline.
10 //!
11 //! The reverse direction (protocol `Op` -> engine `Op`) is not total — engine
12 //! ops carry resolved routes, reply channels, and hook executors that the
13 //! host must supply — so it lands with the engine handle in Phase C/D, not
14 //! here.
15 //!
16 //! The foreground pet observer consumes the event projection, retaining only
17 //! lifecycle metadata. Other projections remain compile-time parity guards;
18 //! dead-code is allowed here rather than hiding those guards behind a test cfg
19 //! (which would let `cargo build` pass with an unmapped variant).
20 #![allow(dead_code)]
21
22 use std::collections::BTreeMap;
23
24 use codewhale_protocol::event_msg as wire;
25 use codewhale_protocol::ids::{SessionId, ThreadId};
26 use codewhale_protocol::op as wire_op;
27 use serde::Serialize;
28 use serde_json::Value;
29
30 use crate::agent_roster::{AgentRosterRow, RosterState};
31 use crate::compaction::CompactionConfig;
32 use crate::config::ProviderKind;
33 use crate::core::engine::preview::PreviewUnresolved;
34 use crate::core::events::{
35 Event, RouteBillingEnvelope, ToolGate, ToolGateVerdict, TurnOutcomeStatus, TurnRoute,
36 };
37 use crate::core::ops::Op;
38 use crate::cost_status::RouteBillingMode;
39 use crate::mcp::{McpManagerSnapshot, McpServerCapabilityMetadata};
40 use crate::model_profile::SupportState;
41 use crate::route_billing::RouteProduct;
42 use crate::tools::spec::ToolError;
43 use crate::tools::subagent::AgentWorkerStatus;
44 use crate::tools::user_input::UserInputRequest;
45 use codewhale_config::AppMode;
46 use codewhale_execpolicy::ApprovalMode;
47 use codewhale_models::Usage;
48 use codewhale_protocol::ResponseChannel;
49
50 /// Routing ids the engine does not carry on each event; the emitter supplies
51 /// them once per session.
52 #[derive(Debug, Clone)]
53 pub struct ProtocolIds {
54 pub thread_id: ThreadId,
55 pub session_id: SessionId,
56 }
57
58 fn to_value<T: Serialize>(value: &T) -> Value {
59 serde_json::to_value(value).unwrap_or(Value::Null)
60 }
61
62 fn count(value: usize) -> u64 {
63 u64::try_from(value).unwrap_or(u64::MAX)
64 }
65
66 /// Lossless mode label; round-trips through `AppMode::parse`.
67 #[must_use]
68 pub fn app_mode_str(mode: AppMode) -> &'static str {
69 match mode {
70 AppMode::Agent => "agent",
71 AppMode::Plan => "plan",
72 AppMode::Operate => "operate",
73 }
74 }
75
76 #[must_use]
77 pub fn approval_mode_str(mode: ApprovalMode) -> &'static str {
78 match mode {
79 ApprovalMode::Auto => "auto",
80 ApprovalMode::Bypass => "bypass",
81 ApprovalMode::Suggest => "suggest",
82 ApprovalMode::Never => "never",
83 }
84 }
85
86 #[must_use]
87 pub fn worker_status_str(status: AgentWorkerStatus) -> &'static str {
88 match status {
89 AgentWorkerStatus::Queued => "queued",
90 AgentWorkerStatus::Starting => "starting",
91 AgentWorkerStatus::Running => "running",
92 AgentWorkerStatus::WaitingForUser => "waiting_for_user",
93 AgentWorkerStatus::ModelWait => "model_wait",
94 AgentWorkerStatus::RunningTool => "running_tool",
95 AgentWorkerStatus::Completed => "completed",
96 AgentWorkerStatus::Failed => "failed",
97 AgentWorkerStatus::Cancelled => "cancelled",
98 AgentWorkerStatus::Interrupted => "interrupted",
99 }
100 }
101
102 fn roster_state_str(state: RosterState) -> &'static str {
103 match state {
104 RosterState::Running => "running",
105 RosterState::Waiting => "waiting",
106 RosterState::Parked => "parked",
107 RosterState::Done => "done",
108 RosterState::Failed => "failed",
109 RosterState::Cancelled => "cancelled",
110 }
111 }
112
113 fn billing_mode_str(mode: RouteBillingMode) -> &'static str {
114 match mode {
115 RouteBillingMode::Metered => "metered",
116 RouteBillingMode::Subscription => "subscription",
117 RouteBillingMode::Local => "local",
118 RouteBillingMode::Unknown => "unknown",
119 }
120 }
121
122 fn support_state_str(state: SupportState) -> &'static str {
123 match state {
124 SupportState::Supported => "supported",
125 SupportState::Unsupported => "unsupported",
126 SupportState::Unknown => "unknown",
127 }
128 }
129
130 fn capability_metadata_str(metadata: McpServerCapabilityMetadata) -> &'static str {
131 match metadata {
132 McpServerCapabilityMetadata::Advertised(_) => "advertised",
133 McpServerCapabilityMetadata::LegacyFallback => "legacy_fallback",
134 McpServerCapabilityMetadata::NotObserved => "not_observed",
135 }
136 }
137
138 fn provider_str(provider: ProviderKind) -> String {
139 provider.as_str().to_string()
140 }
141
142 fn usage_to_wire(usage: &Usage) -> wire::TokenUsage {
143 wire::TokenUsage {
144 input_tokens: usage.input_tokens,
145 output_tokens: usage.output_tokens,
146 prompt_cache_hit_tokens: usage.prompt_cache_hit_tokens,
147 prompt_cache_miss_tokens: usage.prompt_cache_miss_tokens,
148 prompt_cache_write_tokens: usage.prompt_cache_write_tokens,
149 reasoning_tokens: usage.reasoning_tokens,
150 reasoning_replay_tokens: usage.reasoning_replay_tokens,
151 code_execution_requests: usage
152 .server_tool_use
153 .as_ref()
154 .and_then(|server| server.code_execution_requests),
155 tool_search_requests: usage
156 .server_tool_use
157 .as_ref()
158 .and_then(|server| server.tool_search_requests),
159 }
160 }
161
162 fn route_product_to_wire(product: RouteProduct) -> wire::RouteProduct {
163 match product {
164 RouteProduct::Unproven => wire::RouteProduct::Unproven,
165 RouteProduct::Subscription(label) => wire::RouteProduct::Subscription {
166 label: label.to_string(),
167 },
168 RouteProduct::Metered => wire::RouteProduct::Metered,
169 }
170 }
171
172 fn billing_to_wire(billing: &RouteBillingEnvelope) -> wire::RouteBillingEnvelope {
173 wire::RouteBillingEnvelope {
174 openrouter_vendor: billing
175 .openrouter_vendor
176 .as_deref()
177 .map(crate::cost_status::sanitize_persisted_route_label),
178 billing_surface: billing.billing_surface.clone(),
179 endpoint_fingerprint: billing.endpoint_fingerprint.clone(),
180 provider_live_pricing: billing
181 .provider_live_pricing
182 .as_ref()
183 .map(to_value)
184 .filter(|v| !v.is_null()),
185 billing_mode: billing_mode_str(billing.billing_mode).to_string(),
186 dispatched_at: billing.dispatched_at,
187 }
188 }
189
190 fn route_to_wire(route: &TurnRoute) -> wire::TurnRoute {
191 wire::TurnRoute {
192 provider: provider_str(route.provider),
193 provider_identity: route.provider_identity.clone(),
194 model: route.model.clone(),
195 auto_model: route.auto_model,
196 receipt: route
197 .receipt
198 .as_ref()
199 .map(|receipt| wire::TurnRouteReceipt {
200 provider: provider_str(receipt.provider()),
201 provider_identity: receipt.provider_identity().to_string(),
202 wire_model: receipt.wire_model().to_string(),
203 endpoint_identity: receipt.endpoint_identity().to_string(),
204 credential_generation_present: !receipt.credential_generation().is_empty(),
205 }),
206 billing: route.billing.as_ref().map(billing_to_wire),
207 base_url: route.base_url.clone(),
208 billing_product: route_product_to_wire(route.billing_product),
209 }
210 }
211
212 fn tool_error_to_wire(error: &ToolError) -> wire::ToolCallError {
213 match error {
214 ToolError::InvalidInput { message } => wire::ToolCallError::InvalidInput {
215 message: message.clone(),
216 },
217 ToolError::MissingField { field } => wire::ToolCallError::MissingField {
218 field: field.clone(),
219 },
220 ToolError::PathEscape { path } => wire::ToolCallError::PathEscape { path: path.clone() },
221 ToolError::ExecutionFailed { message, .. } => wire::ToolCallError::ExecutionFailed {
222 message: message.clone(),
223 },
224 ToolError::Timeout { seconds } => wire::ToolCallError::Timeout { seconds: *seconds },
225 ToolError::Cancelled { message } => wire::ToolCallError::Cancelled {
226 message: message.clone(),
227 },
228 ToolError::NotAvailable { message } => wire::ToolCallError::NotAvailable {
229 message: message.clone(),
230 },
231 ToolError::PermissionDenied { message } => wire::ToolCallError::PermissionDenied {
232 message: message.clone(),
233 },
234 }
235 }
236
237 fn gate_to_wire(gate: ToolGate) -> wire::ToolGate {
238 match gate {
239 ToolGate::AutoReviewDeterministic => wire::ToolGate::AutoReviewDeterministic,
240 ToolGate::AutoReviewGuardian => wire::ToolGate::AutoReviewGuardian,
241 }
242 }
243
244 fn verdict_to_wire(verdict: ToolGateVerdict) -> wire::ToolGateVerdict {
245 match verdict {
246 ToolGateVerdict::Allowed => wire::ToolGateVerdict::Allowed,
247 ToolGateVerdict::Denied => wire::ToolGateVerdict::Denied,
248 ToolGateVerdict::Unavailable => wire::ToolGateVerdict::Unavailable,
249 }
250 }
251
252 fn outcome_status_to_wire(status: TurnOutcomeStatus) -> wire::TurnOutcomeStatus {
253 match status {
254 TurnOutcomeStatus::Completed => wire::TurnOutcomeStatus::Completed,
255 TurnOutcomeStatus::Interrupted => wire::TurnOutcomeStatus::Interrupted,
256 TurnOutcomeStatus::Failed => wire::TurnOutcomeStatus::Failed,
257 }
258 }
259
260 fn mcp_snapshot_to_wire(snapshot: &McpManagerSnapshot) -> wire::McpManagerSnapshot {
261 let item = |item: &crate::mcp::McpDiscoveredItem| wire::McpDiscoveredItem {
262 name: item.name.clone(),
263 model_name: item.model_name.clone(),
264 description: item.description.clone(),
265 };
266 wire::McpManagerSnapshot {
267 config_path: snapshot.config_path.clone(),
268 config_exists: snapshot.config_exists,
269 reload_required: snapshot.reload_required,
270 servers: snapshot
271 .servers
272 .iter()
273 .map(|server| wire::McpServerSnapshot {
274 name: server.name.clone(),
275 enabled: server.enabled,
276 required: server.required,
277 transport: server.transport.clone(),
278 command_or_url: redacted_command_or_url(&server.command_or_url),
279 connect_timeout: server.connect_timeout,
280 execute_timeout: server.execute_timeout,
281 read_timeout: server.read_timeout,
282 connected: server.connected,
283 error: server.error.clone(),
284 capability_metadata: capability_metadata_str(server.capability_metadata)
285 .to_string(),
286 tools: server.tools.iter().map(item).collect(),
287 resources: server.resources.iter().map(item).collect(),
288 prompts: server.prompts.iter().map(item).collect(),
289 })
290 .collect(),
291 }
292 }
293
294 /// Sanitize an MCP server's configured target before it goes on the wire.
295 ///
296 /// `command_or_url` is the raw configuration: either the configured URL, which
297 /// can carry userinfo (`https://user:token@host`) or query credentials, or a
298 /// stdio command line built as `command + " " + args.join(" ")`, whose args can
299 /// carry a token. That is acceptable in the local picker, where the only reader
300 /// is the person who configured it. This projection is not local -- it feeds
301 /// `EventMsg::McpSessionBoot`, which every SSE and stream-JSON consumer
302 /// receives and any frame-retaining log keeps.
303 ///
304 /// URLs reuse the existing masking in `client::redact_url_for_display`. Stdio
305 /// keeps the program name and elides the arguments, because a secret there can
306 /// be positional and is not reliably recognizable by key.
307 fn redacted_command_or_url(raw: &str) -> String {
308 // No `match` here on purpose: `projections_have_no_wildcard_arms` scans
309 // this whole file for wildcard arms, and a `_ =>` would trip it even in a
310 // helper that projects no engine variant.
311 let trimmed = raw.trim();
312 // A URL never contains whitespace; an argv line does. `://` alone is not
313 // enough to route here: `redact_url_for_display` returns its input
314 // verbatim when `Url::parse` fails, so a stdio command that merely
315 // mentions a URL — `docker run -e TOKEN=… img --url https://…` — would
316 // reach the wire with its positional secret intact.
317 if !trimmed.is_empty() && !trimmed.contains(char::is_whitespace) && trimmed.contains("://") {
318 return crate::client::redact_url_for_display(trimmed);
319 }
320 if let Some((program, rest)) = trimmed.split_once(char::is_whitespace)
321 && !rest.trim().is_empty()
322 {
323 return format!("{program} …");
324 }
325 trimmed.to_string()
326 }
327
328 fn roster_row_to_wire(row: &AgentRosterRow) -> wire::AgentRosterRow {
329 wire::AgentRosterRow {
330 worker_id: row.worker_id.clone(),
331 display_name: row.display_name.clone(),
332 model: row.model.clone(),
333 state: roster_state_str(row.state).to_string(),
334 status: worker_status_str(row.status).to_string(),
335 activity: row.activity.clone(),
336 millis: row.millis,
337 input_tokens: row.input_tokens,
338 output_tokens: row.output_tokens,
339 cost_microusd: row.cost_microusd,
340 steps_taken: row.steps_taken,
341 parent_run_id: row.parent_run_id.clone(),
342 run_id: row.run_id.clone(),
343 }
344 }
345
346 fn user_input_to_wire(request: &UserInputRequest) -> wire::UserInputRequest {
347 wire::UserInputRequest {
348 questions: request
349 .questions
350 .iter()
351 .map(|question| codewhale_protocol::UserInputQuestionEvent {
352 header: question.header.clone(),
353 id: question.id.clone(),
354 question: question.question.clone(),
355 options: question
356 .options
357 .iter()
358 .map(|option| codewhale_protocol::UserInputOptionEvent {
359 label: option.label.clone(),
360 description: option.description.clone(),
361 })
362 .collect(),
363 allow_free_text: question.allow_free_text,
364 multi_select: question.multi_select,
365 })
366 .collect(),
367 }
368 }
369
370 fn compaction_to_wire(config: &CompactionConfig) -> wire_op::CompactionPolicy {
371 wire_op::CompactionPolicy {
372 enabled: config.enabled,
373 token_threshold: count(config.token_threshold),
374 model: config.model.clone(),
375 image_input: support_state_str(config.image_input).to_string(),
376 effective_context_window: config.effective_context_window,
377 cache_summary: config.cache_summary,
378 focus: config.focus.clone(),
379 runtime_cost_owner: config.runtime_cost_owner.clone(),
380 workspace: config.workspace.clone(),
381 }
382 }
383
384 /// Project the engine's per-turn authority onto the wire `TurnSpec`. Host-only
385 /// fields (`initial_routed_usage`, `hook_executor`) and resolved routes are
386 /// stripped; only their non-secret receipts cross.
387 fn turn_spec_to_wire(spec: &crate::core::ops::TurnSpec) -> wire_op::TurnSpec {
388 // `spec.submission_id` is deliberately not projected: the token is
389 // host-process-local by design and the wire op has no correlation twin,
390 // so wire submitters observe `TurnStarted.submission_id` always absent
391 // and cannot correlate submissions on that channel.
392 wire_op::TurnSpec {
393 max_output_tokens: spec.max_output_tokens,
394 content: spec.content.clone(),
395 images: spec.images.clone(),
396 mode: app_mode_str(spec.mode).to_string(),
397 model: Some(spec.route.model.clone()),
398 model_provider: Some(spec.route.identity.key.to_string()),
399 allowed_tools: spec.allowed_tools.clone(),
400 dynamic_tools: spec.dynamic_tools.clone(),
401 provenance: spec.provenance.as_str().to_string(),
402 compaction: Some(Box::new(compaction_to_wire(&spec.compaction))),
403 goal_objective: spec.goal_objective.clone(),
404 goal_token_budget: spec.goal_token_budget,
405 goal_status: spec.goal_status.as_str().to_string(),
406 reasoning_effort: spec.reasoning_effort.clone(),
407 reasoning_effort_auto: spec.reasoning_effort_auto,
408 auto_model: spec.auto_model,
409 allow_shell: spec.allow_shell,
410 trust_mode: spec.trust_mode,
411 auto_approve: spec.auto_approve,
412 approval_mode: approval_mode_str(spec.approval_mode).to_string(),
413 translation_enabled: spec.translation_enabled,
414 verbosity: spec.verbosity.clone(),
415 }
416 }
417
418 fn preview_unresolved_str(unresolved: &PreviewUnresolved) -> String {
419 match unresolved {
420 PreviewUnresolved::AutoRouteNeedsPrompt => "auto_route_needs_prompt".to_string(),
421 PreviewUnresolved::AutoRouteClassificationNotExecuted => {
422 "auto_route_classification_not_executed".to_string()
423 }
424 PreviewUnresolved::NoPrompt => "no_prompt".to_string(),
425 PreviewUnresolved::PlanFailed(error) => format!("plan_failed: {error}"),
426 PreviewUnresolved::MessageSubmitHooksConfigured => {
427 "message_submit_hooks_configured".to_string()
428 }
429 PreviewUnresolved::PromptResolutionFailed(error) => {
430 format!("prompt_resolution_failed: {error}")
431 }
432 }
433 }
434
435 /// Project one engine event onto the protocol. Exhaustive: a new engine
436 /// variant without a protocol twin does not compile.
437 #[must_use]
438 pub fn event_to_protocol(event: &Event, ids: &ProtocolIds) -> wire::EventMsg {
439 let thread_id = ids.thread_id.clone();
440 let session_id = ids.session_id.clone();
441 match event {
442 Event::ToolProjectionWarning {
443 provider,
444 omitted_tool_names,
445 omitted_tool_count,
446 } => wire::EventMsg::ToolProjectionWarning {
447 thread_id,
448 session_id,
449 provider: provider.clone(),
450 omitted_tool_names: omitted_tool_names.clone(),
451 omitted_tool_count: count(*omitted_tool_count),
452 },
453 Event::SnapshotsDisabled { workspace, reason } => wire::EventMsg::SnapshotsDisabled {
454 thread_id,
455 session_id,
456 workspace: workspace.clone(),
457 reason: reason.clone(),
458 },
459 Event::MessageStarted { index } => wire::EventMsg::MessageStarted {
460 thread_id,
461 session_id,
462 index: count(*index),
463 },
464 Event::MessageDelta { index, content } => wire::EventMsg::ResponseDelta {
465 thread_id,
466 session_id,
467 index: count(*index),
468 delta: content.clone(),
469 channel: ResponseChannel::Text,
470 },
471 Event::MessageComplete { index } => wire::EventMsg::MessageComplete {
472 thread_id,
473 session_id,
474 index: count(*index),
475 },
476 Event::ThinkingStarted { index } => wire::EventMsg::ThinkingStarted {
477 thread_id,
478 session_id,
479 index: count(*index),
480 },
481 Event::ThinkingDelta { index, content } => wire::EventMsg::ResponseDelta {
482 thread_id,
483 session_id,
484 index: count(*index),
485 delta: content.clone(),
486 channel: ResponseChannel::Reasoning,
487 },
488 Event::ThinkingComplete { index } => wire::EventMsg::ThinkingComplete {
489 thread_id,
490 session_id,
491 index: count(*index),
492 },
493 Event::ToolExecutionStarted { id } => wire::EventMsg::ToolExecutionStarted {
494 thread_id,
495 session_id,
496 tool_call_id: id.clone(),
497 },
498 Event::ToolResultContent { id, blocks } => wire::EventMsg::ToolResultContent {
499 thread_id,
500 session_id,
501 tool_call_id: id.clone(),
502 blocks: to_value(blocks),
503 },
504 Event::ToolCallStarted {
505 id, name, input, ..
506 } => wire::EventMsg::ToolCallStarted {
507 thread_id,
508 session_id,
509 tool_call_id: id.clone(),
510 tool_name: name.clone(),
511 input: input.clone(),
512 },
513 Event::ToolCallHeartbeat => wire::EventMsg::ToolCallHeartbeat {
514 thread_id,
515 session_id,
516 },
517 Event::ToolCallComplete {
518 id, name, result, ..
519 } => wire::EventMsg::ToolCallComplete {
520 thread_id,
521 session_id,
522 tool_call_id: id.clone(),
523 tool_name: name.clone(),
524 result: match result {
525 Ok(result) => wire::ToolCallOutcome::Ok {
526 content: result.content.clone(),
527 success: result.success,
528 metadata: result.metadata.clone(),
529 },
530 Err(error) => wire::ToolCallOutcome::Err {
531 error: tool_error_to_wire(error),
532 },
533 },
534 },
535 Event::OperationActivityStarted {
536 span_id,
537 activity_kind,
538 } => wire::EventMsg::OperationActivityStarted {
539 thread_id,
540 session_id,
541 span_id: span_id.clone(),
542 activity_kind: *activity_kind,
543 },
544 Event::OperationActivityCompleted {
545 span_id,
546 activity_kind,
547 outcome,
548 } => wire::EventMsg::OperationActivityCompleted {
549 thread_id,
550 session_id,
551 span_id: span_id.clone(),
552 activity_kind: *activity_kind,
553 outcome: *outcome,
554 },
555 Event::TurnStarted {
556 turn_id,
557 created_at,
558 route,
559 submission_id,
560 } => wire::EventMsg::TurnStarted {
561 thread_id,
562 session_id,
563 turn_id: turn_id.clone(),
564 created_at: *created_at,
565 route: route.as_ref().map(route_to_wire),
566 submission_id: submission_id.clone(),
567 },
568 Event::ToolRequestSnapshot { snapshot } => wire::EventMsg::ToolRequestSnapshot {
569 thread_id,
570 session_id,
571 snapshot: to_value(snapshot),
572 },
573 Event::WorkspaceSnapshotTaken { snapshot } => wire::EventMsg::WorkspaceSnapshotTaken {
574 thread_id,
575 session_id,
576 snapshot: to_value(snapshot),
577 },
578 Event::RouteDispatched { turn_id, route } => wire::EventMsg::RouteDispatched {
579 thread_id,
580 session_id,
581 turn_id: turn_id.clone(),
582 route: route_to_wire(route),
583 },
584 Event::TurnComplete {
585 usage,
586 parent_route_usage,
587 routed_usage_dropped_records,
588 status,
589 error,
590 tool_catalog,
591 base_url,
592 } => wire::EventMsg::TurnComplete {
593 thread_id,
594 session_id,
595 turn_id: None,
596 status: outcome_status_to_wire(*status),
597 error: error.clone(),
598 usage: usage_to_wire(usage),
599 parent_route_usage: Some(usage_to_wire(parent_route_usage)),
600 routed_usage_dropped_records: *routed_usage_dropped_records,
601 tool_catalog: tool_catalog
602 .as_ref()
603 .map(|tools| tools.iter().map(to_value).collect()),
604 base_url: base_url.clone(),
605 },
606 Event::TurnUsage {
607 max_output_tokens,
608 usage,
609 duration_ms,
610 first_token_ms,
611 request_ms,
612 } => wire::EventMsg::TurnUsage {
613 max_output_tokens: *max_output_tokens,
614 thread_id,
615 session_id,
616 usage: usage_to_wire(usage),
617 duration_ms: *duration_ms,
618 first_token_ms: *first_token_ms,
619 request_ms: *request_ms,
620 },
621 Event::RoutedTurnUsage {
622 usage,
623 duration_ms,
624 first_token_ms,
625 request_ms,
626 } => wire::EventMsg::RoutedTurnUsage {
627 thread_id,
628 session_id,
629 usage: usage_to_wire(usage),
630 duration_ms: *duration_ms,
631 first_token_ms: *first_token_ms,
632 request_ms: *request_ms,
633 },
634 Event::GoalUpdated { snapshot } => wire::EventMsg::GoalUpdated {
635 thread_id,
636 session_id,
637 snapshot: to_value(snapshot),
638 },
639 Event::GoalContinuationWaiting { delay_seconds } => {
640 wire::EventMsg::GoalContinuationWaiting {
641 thread_id,
642 session_id,
643 delay_seconds: *delay_seconds,
644 }
645 }
646 Event::GoalContinuationWaitEnded { interrupted } => {
647 wire::EventMsg::GoalContinuationWaitEnded {
648 thread_id,
649 session_id,
650 interrupted: *interrupted,
651 }
652 }
653 Event::CompactionStarted { id, auto, message } => wire::EventMsg::CompactionStarted {
654 thread_id,
655 session_id,
656 id: id.clone(),
657 auto: *auto,
658 message: message.clone(),
659 },
660 Event::CompactionCompleted {
661 id,
662 auto,
663 message,
664 messages_before,
665 messages_after,
666 summary_prompt,
667 post_input_tokens,
668 } => wire::EventMsg::CompactionCompleted {
669 thread_id,
670 session_id,
671 id: id.clone(),
672 auto: *auto,
673 message: message.clone(),
674 messages_before: messages_before.map(count),
675 messages_after: messages_after.map(count),
676 summary_prompt: summary_prompt.clone(),
677 post_input_tokens: *post_input_tokens,
678 },
679 Event::CompactionCancelled { id, auto, message } => wire::EventMsg::CompactionCancelled {
680 thread_id,
681 session_id,
682 id: id.clone(),
683 auto: *auto,
684 message: message.clone(),
685 },
686 Event::PurgeStarted { message } => wire::EventMsg::PurgeStarted {
687 thread_id,
688 session_id,
689 message: message.clone(),
690 },
691 Event::PurgeCompleted {
692 messages_before,
693 messages_after,
694 removed_count,
695 replaced_count,
696 message,
697 } => wire::EventMsg::PurgeCompleted {
698 thread_id,
699 session_id,
700 messages_before: count(*messages_before),
701 messages_after: count(*messages_after),
702 removed_count: count(*removed_count),
703 replaced_count: count(*replaced_count),
704 message: message.clone(),
705 },
706 Event::PurgeFailed { message } => wire::EventMsg::PurgeFailed {
707 thread_id,
708 session_id,
709 message: message.clone(),
710 },
711 Event::CompactionFailed { id, auto, message } => wire::EventMsg::CompactionFailed {
712 thread_id,
713 session_id,
714 id: id.clone(),
715 auto: *auto,
716 message: message.clone(),
717 },
718 Event::AgentSpawned {
719 owner_session_id,
720 id,
721 prompt,
722 worker_status,
723 parent_run_id,
724 spawn_depth,
725 model,
726 route_source,
727 display_name,
728 } => wire::EventMsg::AgentSpawned {
729 thread_id,
730 session_id,
731 owner_session_id: owner_session_id.clone(),
732 id: id.clone(),
733 prompt: prompt.clone(),
734 worker_status: worker_status.map(|status| worker_status_str(status).to_string()),
735 parent_run_id: parent_run_id.clone(),
736 spawn_depth: *spawn_depth,
737 model: model.clone(),
738 route_source: route_source.clone(),
739 display_name: display_name.clone(),
740 },
741 Event::AgentProgress {
742 owner_session_id,
743 id,
744 status,
745 activity,
746 parent_run_id,
747 spawn_depth,
748 } => wire::EventMsg::AgentProgress {
749 thread_id,
750 session_id,
751 owner_session_id: owner_session_id.clone(),
752 id: id.clone(),
753 status: status.clone(),
754 activity: wire::AgentProgressActivity {
755 worker_status: worker_status_str(activity.worker_status).to_string(),
756 step: activity.step,
757 tool_name: activity.tool_name.clone(),
758 },
759 parent_run_id: parent_run_id.clone(),
760 spawn_depth: *spawn_depth,
761 },
762 Event::AgentComplete {
763 owner_session_id,
764 id,
765 result,
766 outcome,
767 parent_run_id,
768 spawn_depth,
769 continuable,
770 // Child usage stays off the wire: no protocol client consumes
771 // it, and metrics reads the persisted runtime payload (#6315).
772 usage: _,
773 display_name,
774 } => wire::EventMsg::AgentComplete {
775 thread_id,
776 session_id,
777 owner_session_id: owner_session_id.clone(),
778 id: id.clone(),
779 result: result.clone(),
780 worker_status: outcome
781 .as_ref()
782 .map(|status| crate::tools::subagent::subagent_status_name(status).to_string()),
783 parent_run_id: parent_run_id.clone(),
784 spawn_depth: *spawn_depth,
785 continuable: *continuable,
786 display_name: display_name.clone(),
787 },
788 Event::SubAgentFollowUp {
789 owner_session_id,
790 agent_id,
791 outcome,
792 } => wire::EventMsg::SubAgentFollowUp {
793 thread_id,
794 session_id,
795 owner_session_id: owner_session_id.clone(),
796 agent_id: agent_id.clone(),
797 outcome: match outcome {
798 Ok(outcome) => wire::SubAgentFollowUpOutcome::Ok {
799 agent_id: outcome.agent_id.clone(),
800 target_agent_id: outcome.target_agent_id.clone(),
801 delivered: outcome.delivered,
802 resumed: outcome.resumed,
803 note: outcome.note.clone(),
804 },
805 Err(reason) => wire::SubAgentFollowUpOutcome::Err {
806 reason: reason.clone(),
807 },
808 },
809 },
810 Event::AgentList {
811 owner_session_id,
812 agents,
813 coordination,
814 queued_follow_ups,
815 roster,
816 } => wire::EventMsg::AgentList {
817 thread_id,
818 session_id,
819 owner_session_id: owner_session_id.clone(),
820 agents: agents.iter().map(to_value).collect(),
821 coordination: to_value(coordination),
822 queued_follow_ups: queued_follow_ups
823 .iter()
824 .map(|(agent_id, queued)| (agent_id.clone(), count(*queued)))
825 .collect::<BTreeMap<_, _>>(),
826 roster: roster.iter().map(roster_row_to_wire).collect(),
827 },
828 Event::SubAgentMailbox {
829 owner_session_id,
830 turn_id,
831 seq,
832 message,
833 } => wire::EventMsg::SubAgentMailbox {
834 thread_id,
835 session_id,
836 owner_session_id: owner_session_id.clone(),
837 turn_id: turn_id.clone(),
838 seq: *seq,
839 message: to_value(message),
840 },
841 Event::WorkflowUi {
842 owner_session_id,
843 run_id,
844 event,
845 } => wire::EventMsg::WorkflowUi {
846 thread_id,
847 session_id,
848 owner_session_id: owner_session_id.clone(),
849 run_id: run_id.clone(),
850 ui_event: event.clone(),
851 },
852 Event::Error {
853 envelope,
854 recoverable,
855 } => wire::EventMsg::Error {
856 thread_id,
857 session_id,
858 category: envelope.category.to_string(),
859 severity: envelope.severity.to_string(),
860 recoverable: *recoverable,
861 code: envelope.code.clone(),
862 message: envelope.message.clone(),
863 },
864 Event::Status { message } => wire::EventMsg::Status {
865 thread_id,
866 session_id,
867 message: message.clone(),
868 },
869 Event::McpSessionBoot {
870 generation,
871 snapshot,
872 connecting,
873 finished,
874 } => wire::EventMsg::McpSessionBoot {
875 thread_id,
876 session_id,
877 generation: *generation,
878 snapshot: mcp_snapshot_to_wire(snapshot),
879 connecting: connecting.clone(),
880 finished: *finished,
881 },
882 Event::RequestManifestReady { rendered } => wire::EventMsg::RequestManifestReady {
883 thread_id,
884 session_id,
885 rendered: rendered.clone(),
886 },
887 // The in-process `ack` notifier is an engine handle, not wire data.
888 Event::PauseEvents { ack: _ } => wire::EventMsg::PauseEvents {
889 thread_id,
890 session_id,
891 },
892 Event::ResumeEvents => wire::EventMsg::ResumeEvents {
893 thread_id,
894 session_id,
895 },
896 Event::ApprovalRequired {
897 id,
898 tool_name,
899 description,
900 input,
901 approval_key,
902 approval_grouping_key,
903 intent_summary,
904 approval_force_prompt,
905 } => wire::EventMsg::ApprovalRequired {
906 thread_id,
907 session_id,
908 id: id.clone(),
909 tool_name: tool_name.clone(),
910 description: description.clone(),
911 input: input.clone(),
912 approval_key: approval_key.clone(),
913 approval_grouping_key: approval_grouping_key.clone(),
914 intent_summary: intent_summary.clone(),
915 approval_force_prompt: *approval_force_prompt,
916 },
917 Event::ApprovalWithdrawn { id } => wire::EventMsg::ApprovalWithdrawn {
918 thread_id,
919 session_id,
920 id: id.clone(),
921 },
922 Event::UserInputRequired { id, request } => wire::EventMsg::UserInputRequired {
923 thread_id,
924 session_id,
925 id: id.clone(),
926 request: user_input_to_wire(request),
927 },
928 Event::SessionUpdated {
929 session_id: engine_session_id,
930 messages,
931 system_prompt,
932 model,
933 workspace,
934 } => wire::EventMsg::SessionUpdated {
935 thread_id,
936 session_id,
937 engine_session_id: engine_session_id.clone(),
938 messages: messages.iter().map(to_value).collect(),
939 system_prompt: system_prompt.as_ref().map(to_value),
940 model: model.clone(),
941 workspace: workspace.clone(),
942 },
943 Event::ElevationRequired {
944 tool_id,
945 tool_name,
946 command,
947 denial_reason,
948 blocked_network,
949 blocked_write,
950 } => wire::EventMsg::ElevationRequired {
951 thread_id,
952 session_id,
953 tool_id: tool_id.clone(),
954 tool_name: tool_name.clone(),
955 command: command.clone(),
956 denial_reason: denial_reason.clone(),
957 blocked_network: *blocked_network,
958 blocked_write: *blocked_write,
959 },
960 Event::LspRepairUpdate {
961 diagnostics_found,
962 files,
963 injected,
964 } => wire::EventMsg::LspRepairUpdate {
965 thread_id,
966 session_id,
967 diagnostics_found: count(*diagnostics_found),
968 files: count(*files),
969 injected: *injected,
970 },
971 Event::ToolGateDecision {
972 agent_id,
973 tool_id,
974 tool_name,
975 gate,
976 decision,
977 risk,
978 reason,
979 } => wire::EventMsg::ToolGateDecision {
980 thread_id,
981 session_id,
982 agent_id: agent_id.clone(),
983 tool_id: tool_id.clone(),
984 tool_name: tool_name.clone(),
985 gate: gate_to_wire(*gate),
986 decision: verdict_to_wire(*decision),
987 risk: risk.clone(),
988 reason: reason.clone(),
989 },
990 Event::AdvisoryNote {
991 turn_id,
992 note,
993 tool_call_count,
994 } => wire::EventMsg::AdvisoryNote {
995 thread_id,
996 session_id,
997 turn_id: turn_id.clone(),
998 note: note.clone(),
999 tool_call_count: *tool_call_count,
1000 },
1001 Event::PrefixCacheChange {
1002 description,
1003 system_prompt_changed,
1004 tools_changed,
1005 stability_pct,
1006 changed,
1007 pinned_combined_hash,
1008 pin_reason,
1009 last_miss_reason,
1010 context_updates,
1011 } => wire::EventMsg::PrefixCacheChange {
1012 thread_id,
1013 session_id,
1014 description: description.clone(),
1015 system_prompt_changed: *system_prompt_changed,
1016 tools_changed: *tools_changed,
1017 stability_pct: *stability_pct,
1018 changed: *changed,
1019 pinned_combined_hash: pinned_combined_hash.clone(),
1020 pin_reason: pin_reason.clone(),
1021 last_miss_reason: last_miss_reason.clone(),
1022 context_updates: *context_updates,
1023 },
1024 }
1025 }
1026
1027 /// Project one engine op onto the protocol. Exhaustive: a new engine variant
1028 /// without a protocol twin does not compile. Reply channels, hook executors,
1029 /// and resolved clients are stripped; only their non-secret receipts cross.
1030 #[must_use]
1031 pub fn op_to_protocol(op: &Op) -> wire_op::Op {
1032 match op {
1033 Op::SendMessage(spec) => wire_op::Op::SendMessage(turn_spec_to_wire(spec)),
1034 Op::ContinueGoal {
1035 dynamic_tools,
1036 engine_schedule_id,
1037 } => wire_op::Op::ContinueGoal {
1038 dynamic_tools: dynamic_tools.clone(),
1039 engine_schedule_id: *engine_schedule_id,
1040 },
1041 Op::RunShellCommand {
1042 command,
1043 mode,
1044 allow_shell,
1045 trust_mode,
1046 auto_approve,
1047 approval_mode,
1048 } => wire_op::Op::RunShellCommand {
1049 command: command.clone(),
1050 mode: app_mode_str(*mode).to_string(),
1051 allow_shell: *allow_shell,
1052 trust_mode: *trust_mode,
1053 auto_approve: *auto_approve,
1054 approval_mode: approval_mode_str(*approval_mode).to_string(),
1055 },
1056 Op::SetGoalStatus {
1057 status,
1058 clear,
1059 goal_id,
1060 } => wire_op::Op::SetGoalStatus {
1061 goal_id: goal_id.clone(),
1062 status: status.as_str().to_string(),
1063 clear: *clear,
1064 },
1065 Op::SetGoalObjective {
1066 objective,
1067 token_budget,
1068 goal_id,
1069 } => wire_op::Op::SetGoalObjective {
1070 goal_id: goal_id.clone(),
1071 objective: objective.clone(),
1072 token_budget: *token_budget,
1073 },
1074 Op::PreviewOutboundRequest {
1075 inputs,
1076 json,
1077 base_prompt_only,
1078 } => wire_op::Op::PreviewOutboundRequest {
1079 json: *json,
1080 base_prompt_only: *base_prompt_only,
1081 mode: app_mode_str(inputs.mode).to_string(),
1082 allow_shell: inputs.allow_shell,
1083 trust_mode: inputs.trust_mode,
1084 auto_approve: inputs.auto_approve,
1085 approval_mode: approval_mode_str(inputs.approval_mode).to_string(),
1086 allowed_tools: inputs.allowed_tools.clone(),
1087 dynamic_tools: inputs.dynamic_tools.clone(),
1088 provenance: inputs.provenance.as_str().to_string(),
1089 requested_model: inputs.requested_model.clone(),
1090 requested_reasoning: inputs.requested_reasoning.clone(),
1091 auto_model: inputs.auto_model,
1092 hypothetical_prompt_supplied: inputs.hypothetical_prompt_supplied,
1093 hypothetical_prompt: inputs.next_turn.as_ref().map(|turn| turn.content.clone()),
1094 unresolved: if inputs.next_turn.is_some() {
1095 None
1096 } else {
1097 Some(preview_unresolved_str(&inputs.unresolved))
1098 },
1099 },
1100 Op::ListSubAgents => wire_op::Op::ListSubAgents,
1101 Op::GetSubAgentSettlement { tx: _ } => wire_op::Op::GetSubAgentSettlement,
1102 Op::CancelSubAgent { agent_id } => wire_op::Op::CancelSubAgent {
1103 agent_id: agent_id.clone(),
1104 },
1105 Op::FollowUpSubAgent { agent_id, text } => wire_op::Op::FollowUpSubAgent {
1106 agent_id: agent_id.clone(),
1107 text: text.clone(),
1108 },
1109 Op::ChangeMode {
1110 mode,
1111 allow_shell,
1112 trust_mode,
1113 auto_approve,
1114 approval_mode,
1115 configured_sandbox_mode,
1116 } => wire_op::Op::ChangeMode {
1117 mode: app_mode_str(*mode).to_string(),
1118 allow_shell: *allow_shell,
1119 trust_mode: *trust_mode,
1120 auto_approve: *auto_approve,
1121 approval_mode: approval_mode_str(*approval_mode).to_string(),
1122 configured_sandbox_mode: configured_sandbox_mode.clone(),
1123 },
1124 Op::SetModel {
1125 model,
1126 mode,
1127 route_limits,
1128 } => wire_op::Op::SetModel {
1129 model: model.clone(),
1130 mode: app_mode_str(*mode).to_string(),
1131 route_limits: route_limits.map(|limits| wire_op::RouteLimits {
1132 context_tokens: limits.context_tokens,
1133 input_tokens: limits.input_tokens,
1134 output_tokens: limits.output_tokens,
1135 }),
1136 },
1137 Op::SetCompaction { config } => wire_op::Op::SetCompaction {
1138 config: compaction_to_wire(config),
1139 },
1140 Op::SetStreamChunkTimeout { timeout_secs } => wire_op::Op::SetStreamChunkTimeout {
1141 timeout_secs: *timeout_secs,
1142 },
1143 Op::SetSubagentRuntimeConfig {
1144 enabled,
1145 max_subagents,
1146 launch_concurrency,
1147 max_spawn_depth,
1148 api_timeout_secs,
1149 heartbeat_timeout_secs,
1150 } => wire_op::Op::SetSubagentRuntimeConfig {
1151 enabled: *enabled,
1152 max_subagents: count(*max_subagents),
1153 launch_concurrency: count(*launch_concurrency),
1154 max_spawn_depth: *max_spawn_depth,
1155 api_timeout_secs: *api_timeout_secs,
1156 heartbeat_timeout_secs: *heartbeat_timeout_secs,
1157 },
1158 Op::SetSearchProvider { provider } => wire_op::Op::SetSearchProvider {
1159 provider: provider.as_str().to_string(),
1160 },
1161 Op::SetFleetRoster { roster } => wire_op::Op::SetFleetRoster {
1162 member_ids: roster
1163 .members()
1164 .iter()
1165 .map(|member| member.id.clone())
1166 .collect(),
1167 exact_selection: roster.is_exact_selection(),
1168 load_error: roster.load_error().map(str::to_string),
1169 },
1170 Op::SyncSession {
1171 session_id,
1172 messages,
1173 system_prompt,
1174 system_prompt_override,
1175 model,
1176 workspace,
1177 mode,
1178 } => wire_op::Op::SyncSession {
1179 engine_session_id: session_id.clone(),
1180 messages: messages.iter().map(to_value).collect(),
1181 system_prompt: system_prompt.as_ref().map(to_value),
1182 system_prompt_override: *system_prompt_override,
1183 model: model.clone(),
1184 workspace: workspace.clone(),
1185 mode: app_mode_str(*mode).to_string(),
1186 },
1187 Op::RewindConversation {
1188 expected,
1189 messages,
1190 tx: _,
1191 } => wire_op::Op::RewindConversation {
1192 expected: to_value(expected),
1193 messages: messages.iter().map(to_value).collect(),
1194 },
1195 Op::CompactContext {
1196 id,
1197 route,
1198 compaction,
1199 } => wire_op::Op::CompactContext {
1200 id: id.clone(),
1201 model: route.model.clone(),
1202 model_provider: route.identity.key.to_string(),
1203 compaction: compaction_to_wire(compaction),
1204 },
1205 Op::CancelCompaction { id } => wire_op::Op::CancelCompaction { id: id.clone() },
1206 // Reply channels never cross the wire: the answer is a frame.
1207 Op::GetSessionSnapshot { tx: _ } => wire_op::Op::GetSessionSnapshot,
1208 Op::GetContextBudget { tx: _ } => wire_op::Op::GetContextBudget,
1209 Op::GetProviderRuntimeStatus { tx: _ } => wire_op::Op::GetProviderRuntimeStatus,
1210 Op::BootstrapMcp { tx: _ } => wire_op::Op::BootstrapMcp,
1211 Op::RetryMcpServer { name, tx: _ } => wire_op::Op::RetryMcpServer { name: name.clone() },
1212 Op::ReloadMcp { config_path, tx: _ } => wire_op::Op::ReloadMcp {
1213 config_path: config_path.clone(),
1214 },
1215 Op::PurgeContext => wire_op::Op::PurgeContext,
1216 Op::EditLastTurn {
1217 new_message,
1218 submission_id: _,
1219 } => wire_op::Op::EditLastTurn {
1220 new_message: new_message.clone(),
1221 },
1222 Op::SetAdvisorEnabled { enabled } => wire_op::Op::SetAdvisorEnabled { enabled: *enabled },
1223 Op::Shutdown => wire_op::Op::Shutdown,
1224 }
1225 }
1226
1227 impl Event {
1228 /// Protocol projection of this event. See [`event_to_protocol`].
1229 #[must_use]
1230 pub fn to_protocol(&self, ids: &ProtocolIds) -> wire::EventMsg {
1231 event_to_protocol(self, ids)
1232 }
1233 }
1234
1235 impl Op {
1236 /// Protocol projection of this op. See [`op_to_protocol`].
1237 #[must_use]
1238 pub fn to_protocol(&self) -> wire_op::Op {
1239 op_to_protocol(self)
1240 }
1241 }
1242
1243 #[cfg(test)]
1244 mod tests {
1245 #[test]
1246 fn mcp_command_or_url_is_redacted_before_it_reaches_the_wire() {
1247 // Local display may show the configured value; this projection feeds
1248 // EventMsg::McpSessionBoot, which every SSE/stream-JSON consumer sees.
1249 assert_eq!(
1250 redacted_command_or_url("https://user:tok@mcp.example.com/sse"),
1251 "https://***:***@mcp.example.com/sse"
1252 );
1253 let masked = redacted_command_or_url("https://mcp.example.com/sse?api_key=SECRET");
1254 assert!(!masked.contains("SECRET"), "{masked}");
1255 // stdio: keep the program, drop the args -- a token there can be
1256 // positional, so key-based masking is not enough.
1257 assert_eq!(
1258 redacted_command_or_url("npx -y server --token SECRET"),
1259 "npx …"
1260 );
1261 // Nothing to hide, nothing changed.
1262 assert_eq!(redacted_command_or_url("npx"), "npx");
1263 // A stdio command that merely mentions a URL must still collapse to
1264 // the program name. Routing it to the URL redactor returns it verbatim
1265 // (`Url::parse` rejects the spaces), leaking the positional token.
1266 let mentions_url = redacted_command_or_url(
1267 "docker run -e TOKEN=sk-live-abc img --url https://mcp.example.com",
1268 );
1269 assert_eq!(mentions_url, "docker …");
1270 assert!(
1271 !mentions_url.contains("sk-live-abc"),
1272 "argv secret reached the wire: {mentions_url}"
1273 );
1274 }
1275
1276 use super::*;
1277 use crate::error_taxonomy::{ErrorCategory, ErrorEnvelope, ErrorSeverity};
1278 use crate::tools::goal::GoalStatus;
1279 use crate::tools::spec::ToolResult;
1280 use serde_json::json;
1281
1282 const SOURCE: &str = include_str!("protocol_parity.rs");
1283
1284 fn ids() -> ProtocolIds {
1285 ProtocolIds {
1286 thread_id: ThreadId::new(),
1287 session_id: SessionId::new(),
1288 }
1289 }
1290
1291 #[test]
1292 fn worker_lifecycle_wire_preserves_typed_outcomes_without_parsing_result_text() {
1293 use crate::tools::subagent::SubAgentStatus;
1294 let ids = ids();
1295 for (outcome, expected) in [
1296 (Some(SubAgentStatus::Completed), Some("completed")),
1297 (
1298 Some(SubAgentStatus::Failed("private failure".into())),
1299 Some("failed"),
1300 ),
1301 (
1302 Some(SubAgentStatus::Interrupted("private reason".into())),
1303 Some("interrupted"),
1304 ),
1305 (Some(SubAgentStatus::Cancelled), Some("cancelled")),
1306 (
1307 Some(SubAgentStatus::BudgetExhausted),
1308 Some("budget_exhausted"),
1309 ),
1310 (None, None),
1311 ] {
1312 let event = Event::AgentComplete {
1313 display_name: None,
1314 owner_session_id: "owner".into(),
1315 id: "worker".into(),
1316 result: "Completed successfully".into(),
1317 outcome,
1318 parent_run_id: Some("parent".into()),
1319 spawn_depth: Some(2),
1320 continuable: Some(false),
1321 usage: None,
1322 };
1323 let wire = serde_json::to_value(event_to_protocol(&event, &ids)).unwrap();
1324 assert_eq!(wire["worker_status"].as_str(), expected);
1325 assert_eq!(wire["parent_run_id"], "parent");
1326 assert_eq!(wire["spawn_depth"], 2);
1327 assert_eq!(wire["continuable"], false);
1328 assert!(!wire.to_string().contains("private"));
1329 }
1330 for status in [
1331 AgentWorkerStatus::Queued,
1332 AgentWorkerStatus::Starting,
1333 AgentWorkerStatus::Running,
1334 AgentWorkerStatus::WaitingForUser,
1335 AgentWorkerStatus::ModelWait,
1336 AgentWorkerStatus::RunningTool,
1337 AgentWorkerStatus::Completed,
1338 AgentWorkerStatus::Failed,
1339 AgentWorkerStatus::Cancelled,
1340 AgentWorkerStatus::Interrupted,
1341 ] {
1342 // Runtime serializes the producer enum; stream-json uses this
1343 // exhaustive adapter. Their discriminants must remain identical.
1344 assert_eq!(
1345 serde_json::to_value(status).unwrap(),
1346 worker_status_str(status)
1347 );
1348 }
1349 }
1350
1351 #[test]
1352 fn wire_accounting_preserves_parent_total_and_distinct_routed_telemetry() {
1353 let ids = ids();
1354 let total = Usage {
1355 input_tokens: 49,
1356 output_tokens: 19,
1357 ..Usage::default()
1358 };
1359 let parent = Usage {
1360 input_tokens: 7,
1361 output_tokens: 5,
1362 ..Usage::default()
1363 };
1364 let complete = event_to_protocol(
1365 &Event::TurnComplete {
1366 usage: total.clone(),
1367 parent_route_usage: parent.clone(),
1368 routed_usage_dropped_records: 3,
1369 status: TurnOutcomeStatus::Completed,
1370 error: None,
1371 tool_catalog: None,
1372 base_url: None,
1373 },
1374 &ids,
1375 );
1376 let json = serde_json::to_value(&complete).unwrap();
1377 assert_eq!(json["usage"]["input_tokens"], 49);
1378 assert_eq!(json["usage"]["output_tokens"], 19);
1379 assert_eq!(json["parent_route_usage"]["input_tokens"], 7);
1380 assert_eq!(json["parent_route_usage"]["output_tokens"], 5);
1381 assert_eq!(json["routed_usage_dropped_records"], 3);
1382 assert_eq!(
1383 serde_json::from_value::<wire::EventMsg>(json.clone()).unwrap(),
1384 complete
1385 );
1386
1387 // A legacy terminal receipt has no parent subset, which differs from
1388 // an explicitly reported zero parent on a compaction-only turn.
1389 let mut legacy = json;
1390 legacy.as_object_mut().unwrap().remove("parent_route_usage");
1391 legacy
1392 .as_object_mut()
1393 .unwrap()
1394 .remove("routed_usage_dropped_records");
1395 assert!(matches!(
1396 serde_json::from_value::<wire::EventMsg>(legacy).unwrap(),
1397 wire::EventMsg::TurnComplete {
1398 parent_route_usage: None,
1399 routed_usage_dropped_records: 0,
1400 ..
1401 }
1402 ));
1403
1404 for (event, tag) in [
1405 (
1406 Event::TurnUsage {
1407 max_output_tokens: None,
1408 usage: parent,
1409 duration_ms: 12,
1410 first_token_ms: Some(2),
1411 request_ms: Some(10),
1412 },
1413 "turn_usage",
1414 ),
1415 (
1416 Event::RoutedTurnUsage {
1417 usage: total,
1418 duration_ms: 27,
1419 first_token_ms: None,
1420 request_ms: None,
1421 },
1422 "routed_turn_usage",
1423 ),
1424 ] {
1425 let projected = event_to_protocol(&event, &ids);
1426 let json = serde_json::to_value(&projected).unwrap();
1427 assert_eq!(json["event"], tag);
1428 assert_eq!(
1429 serde_json::from_value::<wire::EventMsg>(json).unwrap(),
1430 projected
1431 );
1432 }
1433 }
1434
1435 /// The guard is the exhaustive `match` in `event_to_protocol`: this test
1436 /// exists so the guard has a name in the test log and so the projection
1437 /// is proven to agree with the protocol's wire-tag table.
1438 #[test]
1439 fn protocol_covers_engine_events() {
1440 let ids = ids();
1441 let usage = Usage {
1442 input_tokens: 3,
1443 output_tokens: 4,
1444 ..Usage::default()
1445 };
1446 let events = vec![
1447 Event::MessageStarted { index: 0 },
1448 Event::MessageDelta {
1449 index: 0,
1450 content: "hello".into(),
1451 },
1452 Event::ThinkingDelta {
1453 index: 1,
1454 content: "hmm".into(),
1455 },
1456 Event::ToolCallStarted {
1457 model_call: None,
1458 id: "c1".into(),
1459 name: "read_file".into(),
1460 input: json!({"path": "x"}),
1461 },
1462 Event::ToolCallHeartbeat,
1463 Event::ToolCallComplete {
1464 model_call: None,
1465 id: "c1".into(),
1466 name: "read_file".into(),
1467 result: Ok(ToolResult::success("ok")),
1468 },
1469 Event::ToolCallComplete {
1470 model_call: None,
1471 id: "c2".into(),
1472 name: "bash".into(),
1473 result: Err(ToolError::Timeout { seconds: 9 }),
1474 },
1475 Event::TurnStarted {
1476 turn_id: "turn-1".into(),
1477 created_at: chrono::Utc::now(),
1478 route: None,
1479 // A host-stamped token, not `None`: the projection must carry
1480 // it verbatim or the assertion on `events[7]` below fails.
1481 submission_id: Some("sub-host-1".into()),
1482 },
1483 Event::TurnComplete {
1484 usage: usage.clone(),
1485 parent_route_usage: usage.clone(),
1486 routed_usage_dropped_records: 0,
1487 status: TurnOutcomeStatus::Interrupted,
1488 error: Some("stopped".into()),
1489 tool_catalog: None,
1490 base_url: Some("https://example.invalid".into()),
1491 },
1492 Event::WorkspaceSnapshotTaken {
1493 snapshot: crate::snapshot::WorkspaceSnapshotRef {
1494 kind: crate::snapshot::WorkspaceSnapshotKind::Tool,
1495 snapshot_id: "a".repeat(40),
1496 tree_id: "b".repeat(40),
1497 session_id: "thr_1".into(),
1498 tool_call_id: Some("c1".into()),
1499 write_paths: Some(vec!["src/lib.rs".into()]),
1500 changed_paths: Some(Vec::new()),
1501 },
1502 },
1503 Event::RoutedTurnUsage {
1504 usage: usage.clone(),
1505 duration_ms: 12,
1506 first_token_ms: Some(3),
1507 request_ms: None,
1508 },
1509 Event::TurnUsage {
1510 max_output_tokens: None,
1511 usage,
1512 duration_ms: 12,
1513 first_token_ms: Some(3),
1514 request_ms: None,
1515 },
1516 Event::Error {
1517 envelope: ErrorEnvelope {
1518 category: ErrorCategory::RateLimit,
1519 severity: ErrorSeverity::Warning,
1520 recoverable: true,
1521 code: "E429".into(),
1522 message: "slow down".into(),
1523 },
1524 recoverable: true,
1525 },
1526 Event::status("ready"),
1527 Event::PauseEvents { ack: None },
1528 Event::ResumeEvents,
1529 Event::ToolGateDecision {
1530 agent_id: None,
1531 tool_id: "c3".into(),
1532 tool_name: "bash".into(),
1533 gate: ToolGate::AutoReviewGuardian,
1534 decision: ToolGateVerdict::Unavailable,
1535 risk: None,
1536 reason: "timeout".into(),
1537 },
1538 Event::WorkflowUi {
1539 owner_session_id: "owner".into(),
1540 run_id: "run".into(),
1541 event: json!({"type": "task_started"}),
1542 },
1543 ];
1544
1545 for event in &events {
1546 let msg = event.to_protocol(&ids);
1547 assert!(
1548 wire::EVENT_KINDS.contains(&msg.kind_str()),
1549 "{} is not in EVENT_KINDS",
1550 msg.kind_str()
1551 );
1552 assert_eq!(msg.thread_id(), &ids.thread_id);
1553 assert_eq!(msg.session_id(), &ids.session_id);
1554 let value = serde_json::to_value(&msg).unwrap();
1555 assert_eq!(value["event"], msg.kind_str());
1556 let back: wire::EventMsg = serde_json::from_value(value).unwrap();
1557 assert_eq!(back, msg);
1558 }
1559
1560 let delta = events[1].to_protocol(&ids);
1561 let thinking = events[2].to_protocol(&ids);
1562 assert!(matches!(
1563 delta,
1564 wire::EventMsg::ResponseDelta {
1565 channel: ResponseChannel::Text,
1566 ..
1567 }
1568 ));
1569 assert!(matches!(
1570 thinking,
1571 wire::EventMsg::ResponseDelta {
1572 channel: ResponseChannel::Reasoning,
1573 ..
1574 }
1575 ));
1576 assert_eq!(
1577 serde_json::to_value(events[6].to_protocol(&ids)).unwrap()["result"],
1578 json!({"outcome": "err", "error": {"kind": "timeout", "seconds": 9}})
1579 );
1580 // The host correlation token must cross the projection verbatim: the
1581 // app forwarder binds its submit-window actions to this echo, so a
1582 // silently dropped mapping would defeat the contract.
1583 assert_eq!(
1584 serde_json::to_value(events[7].to_protocol(&ids)).unwrap()["submission_id"],
1585 json!("sub-host-1")
1586 );
1587 assert_eq!(
1588 serde_json::to_value(events[8].to_protocol(&ids)).unwrap()["status"],
1589 "interrupted"
1590 );
1591 let error = events
1592 .iter()
1593 .find(|event| matches!(event, Event::Error { .. }))
1594 .unwrap();
1595 let error = serde_json::to_value(error.to_protocol(&ids)).unwrap();
1596 assert_eq!(error["category"], "rate_limit");
1597 assert_eq!(error["severity"], "warning");
1598 }
1599
1600 #[test]
1601 fn conditional_rewind_projection_preserves_expected_state_and_reply_owner() {
1602 let expected = crate::core::ops::SessionSnapshot {
1603 session_id: "observed-conversation".into(),
1604 messages: vec![],
1605 total_tokens: 9,
1606 model: "observed-model".into(),
1607 model_provider: "deepseek".into(),
1608 model_provider_id: Some("configured-provider".into()),
1609 workspace: std::path::PathBuf::from("/observed/workspace"),
1610 system_prompt: None,
1611 mode: "agent".into(),
1612 };
1613 let (tx, mut receive) = tokio::sync::oneshot::channel();
1614 let op = Op::RewindConversation {
1615 expected: Box::new(expected.clone()),
1616 messages: vec![],
1617 tx,
1618 };
1619 let value = serde_json::to_value(op.to_protocol()).unwrap();
1620 assert_eq!(value["kind"], "rewind_conversation");
1621 assert!(value["expected"] == to_value(&expected));
1622 assert!(value.get("tx").is_none());
1623 assert!(matches!(
1624 receive.try_recv(),
1625 Err(tokio::sync::oneshot::error::TryRecvError::Empty)
1626 ));
1627 }
1628
1629 #[test]
1630 fn protocol_covers_engine_ops() {
1631 let (tx, _rx) = tokio::sync::oneshot::channel();
1632 let settlement_reply = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
1633 let ops = vec![
1634 Op::SetGoalStatus {
1635 goal_id: None,
1636 status: GoalStatus::Paused,
1637 clear: false,
1638 },
1639 Op::SetGoalObjective {
1640 goal_id: None,
1641 objective: "ship".into(),
1642 token_budget: Some(7),
1643 },
1644 Op::ListSubAgents,
1645 Op::CancelSubAgent {
1646 agent_id: "a1".into(),
1647 },
1648 Op::FollowUpSubAgent {
1649 agent_id: "a1".into(),
1650 text: "go".into(),
1651 },
1652 Op::SetStreamChunkTimeout { timeout_secs: 30 },
1653 Op::SetSubagentRuntimeConfig {
1654 enabled: true,
1655 max_subagents: 4,
1656 launch_concurrency: 2,
1657 max_spawn_depth: 1,
1658 api_timeout_secs: 60,
1659 heartbeat_timeout_secs: 10,
1660 },
1661 Op::CancelCompaction { id: "cmp".into() },
1662 Op::GetSessionSnapshot {
1663 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1664 },
1665 Op::GetContextBudget {
1666 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1667 },
1668 Op::GetProviderRuntimeStatus {
1669 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1670 },
1671 Op::BootstrapMcp {
1672 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1673 },
1674 Op::RetryMcpServer {
1675 name: "fs".into(),
1676 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1677 },
1678 Op::ReloadMcp {
1679 config_path: std::path::PathBuf::from("/tmp/mcp.json"),
1680 tx: std::sync::Arc::new(std::sync::Mutex::new(None)),
1681 },
1682 Op::PurgeContext,
1683 Op::EditLastTurn {
1684 new_message: "again".into(),
1685 submission_id: None,
1686 },
1687 Op::SetAdvisorEnabled { enabled: true },
1688 Op::GetSubAgentSettlement {
1689 tx: std::sync::Arc::clone(&settlement_reply),
1690 },
1691 Op::Shutdown,
1692 ];
1693
1694 for op in &ops {
1695 let msg = op.to_protocol();
1696 assert!(
1697 wire_op::OP_KINDS.contains(&msg.kind_str()),
1698 "{} is not in OP_KINDS",
1699 msg.kind_str()
1700 );
1701 let value = serde_json::to_value(&msg).unwrap();
1702 assert_eq!(value["kind"], msg.kind_str());
1703 let back: wire_op::Op = serde_json::from_value(value).unwrap();
1704 assert_eq!(back, msg);
1705 }
1706
1707 assert_eq!(
1708 serde_json::to_value(ops[0].to_protocol()).unwrap(),
1709 json!({"kind": "set_goal_status", "status": "paused", "clear": false})
1710 );
1711 assert_eq!(
1712 serde_json::to_value(ops[8].to_protocol()).unwrap(),
1713 json!({"kind": "get_session_snapshot"}),
1714 "reply channels must not leak onto the wire"
1715 );
1716 let settlement = ops
1717 .iter()
1718 .find(|op| matches!(op, Op::GetSubAgentSettlement { .. }))
1719 .unwrap();
1720 assert_eq!(
1721 serde_json::to_value(settlement.to_protocol()).unwrap(),
1722 json!({"kind": "get_sub_agent_settlement"}),
1723 "the settlement operation must retain its own channel-free protocol twin"
1724 );
1725 assert!(
1726 settlement_reply.lock().unwrap().is_some(),
1727 "projection must not consume the host's live response sender"
1728 );
1729 }
1730
1731 #[test]
1732 fn mode_labels_round_trip_through_app_mode_parse() {
1733 for mode in [AppMode::Agent, AppMode::Plan, AppMode::Operate] {
1734 assert_eq!(AppMode::parse(app_mode_str(mode)), Some(mode), "{mode:?}");
1735 }
1736 for mode in [
1737 ApprovalMode::Auto,
1738 ApprovalMode::Bypass,
1739 ApprovalMode::Suggest,
1740 ApprovalMode::Never,
1741 ] {
1742 assert_eq!(
1743 ApprovalMode::from_config_value(approval_mode_str(mode)),
1744 Some(mode)
1745 );
1746 }
1747 }
1748
1749 /// The projections are only a guard while they stay exhaustive. A
1750 /// wildcard arm would let a new engine variant slip through unmapped.
1751 #[test]
1752 fn projections_have_no_wildcard_arms() {
1753 let wildcard_arms: Vec<&str> = SOURCE
1754 .lines()
1755 .filter(|line| {
1756 let trimmed = line.trim_start();
1757 trimmed.starts_with("_ =>")
1758 || trimmed.starts_with("_=>")
1759 || trimmed.starts_with("Event::_")
1760 || trimmed.starts_with("Op::_")
1761 || (trimmed.contains(" => ") && trimmed.starts_with("other =>"))
1762 })
1763 .collect();
1764 assert!(
1765 wildcard_arms.is_empty(),
1766 "protocol_parity.rs must match engine variants exhaustively; found {wildcard_arms:?}"
1767 );
1768 }
1769 #[test]
1770 fn acp_optional_tool_facts_round_trip_with_core_routing_identity() {
1771 let ids = ids();
1772 for event in [
1773 Event::ToolExecutionStarted {
1774 id: "core-call".into(),
1775 },
1776 Event::ToolResultContent {
1777 id: "core-call".into(),
1778 blocks: vec![codewhale_tools::ToolResultContentBlock::Image {
1779 mime_type: "image/png".into(),
1780 data: "QUJD".into(),
1781 }],
1782 },
1783 ] {
1784 let projected = event_to_protocol(&event, &ids);
1785 assert_eq!(projected.thread_id(), &ids.thread_id);
1786 assert_eq!(projected.session_id(), &ids.session_id);
1787 let value = serde_json::to_value(&projected).unwrap();
1788 assert_eq!(
1789 serde_json::from_value::<wire::EventMsg>(value).unwrap(),
1790 projected
1791 );
1792 }
1793 }
1794 }
1795
1795 lines RUST