返回 CodeWhale
op.rs
根目录 / crates / protocol / src / op.rs
1 //! `Op`-in API in `crates/protocol` (issue #5261, Phase A1 of the
2 //! core/protocol extraction spec).
3 //!
4 //! The TUI engine already had an internal channel (`Op` in
5 //! `crates/tui/src/core/ops.rs` with `tx_op` / `rx_op` and `tx_steer`).
6 //! This protocol file formalizes that channel so TUI, CLI, app-server, and
7 //! tests share one serializable API. The wire is `OpEnvelope` + `Op`;
8 //! transports that already speak JSON (app-server, tests) can send the
9 //! envelope directly, while in-process callers continue to use the typed
10 //! enum.
11 //!
12 //! Parity is compile-enforced from the engine side:
13 //! `crates/tui/src/core/protocol_parity.rs` matches every engine `Op`
14 //! variant exhaustively into a protocol `Op` (`protocol_covers_engine_ops`).
15 //!
16 //! What is deliberately stripped at this boundary:
17 //!
18 //! - `mpsc` / `oneshot` reply channels (`GetSubAgentSettlement`, `GetSessionSnapshot`,
19 //! `GetContextBudget`, `GetProviderRuntimeStatus`, `BootstrapMcp`, `RetryMcpServer`,
20 //! `ReloadMcp`). Over the wire the reply is an `EventMsg` or a response
21 //! frame, not a channel.
22 //! - `Arc<HookExecutor>` on `SendMessage`: hooks are host configuration, not
23 //! turn input.
24 //! - Resolved provider routes and clients. Only the non-secret receipt
25 //! (`model`, `model_provider`) crosses; the engine re-resolves.
26 //! - `Ruleset`, transcript `Message`s, and `SystemPrompt` cross as
27 //! `serde_json::Value` from their own `Serialize` until typed in a later
28 //! phase.
29 //!
30 //! `Steer` and `Cancel` have no engine `Op` twin on purpose: the engine
31 //! carries them on `tx_steer` and the cancellation token. They are protocol
32 //! operations regardless, because every out-of-process client needs them.
33
34 use std::path::PathBuf;
35
36 use serde::{Deserialize, Serialize};
37 use serde_json::Value;
38
39 /// Accept `engine_schedule_id` on the wire for compatibility, then discard it.
40 ///
41 /// Deserialization is the out-of-process boundary: the engine's own
42 /// `Op::ContinueGoal` travels an in-process channel and never reaches here, so
43 /// anything this sees was supplied by a caller who must not be able to set it.
44 /// Returning `None` sends such a request down the ordinary host-injected path
45 /// instead of letting it consume a pending engine schedule.
46 fn drop_engine_owned_schedule_id<'de, D>(deserializer: D) -> Result<Option<u64>, D::Error>
47 where
48 D: serde::Deserializer<'de>,
49 {
50 Option::<u64>::deserialize(deserializer)?;
51 Ok(None)
52 }
53
54 use crate::ids::{SessionId, ThreadId};
55 use crate::runtime::DynamicToolSpec;
56
57 /// Every `Op` is paired with the ids that route it. This is the
58 /// `Op`-in half of the `Op`-in / `EventMsg`-out contract.
59 #[derive(Debug, Clone, Serialize, Deserialize)]
60 pub struct OpEnvelope {
61 /// Monotonic `op:<n>` for dedup / tracing within a session.
62 pub op_id: String,
63 pub thread_id: ThreadId,
64 pub session_id: SessionId,
65 pub op: Op,
66 }
67
68 /// Token-limit facts for a route. `None` is unknown, never zero.
69 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
70 pub struct RouteLimits {
71 #[serde(default, skip_serializing_if = "Option::is_none")]
72 pub context_tokens: Option<u64>,
73 #[serde(default, skip_serializing_if = "Option::is_none")]
74 pub input_tokens: Option<u64>,
75 #[serde(default, skip_serializing_if = "Option::is_none")]
76 pub output_tokens: Option<u64>,
77 }
78
79 /// Compaction policy derived from a provider route. Twin of the engine's
80 /// `CompactionConfig`.
81 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
82 pub struct CompactionPolicy {
83 pub enabled: bool,
84 pub token_threshold: u64,
85 pub model: String,
86 /// `supported | unsupported | unknown`.
87 #[serde(default = "default_capability_state")]
88 pub image_input: String,
89 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub effective_context_window: Option<u32>,
91 #[serde(default)]
92 pub cache_summary: bool,
93 #[serde(default, skip_serializing_if = "Option::is_none")]
94 pub focus: Option<String>,
95 #[serde(default, skip_serializing_if = "Option::is_none")]
96 pub runtime_cost_owner: Option<String>,
97 #[serde(default, skip_serializing_if = "Option::is_none")]
98 pub workspace: Option<PathBuf>,
99 }
100
101 fn default_capability_state() -> String {
102 "unknown".to_string()
103 }
104
105 /// Per-turn authority payload carried by [`Op::SendMessage`]. Serializable
106 /// twin of `crates/tui/src/core/ops::TurnSpec`; extracted from the enum arm so
107 /// new per-turn fields accrete here instead of widening the variant. The
108 /// `SendMessage(TurnSpec)` newtype keeps the internally-tagged wire shape
109 /// byte-identical: `{"kind":"send_message", ...fields}`.
110 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
111 pub struct TurnSpec {
112 #[serde(
113 default,
114 rename = "maxOutputTokens",
115 alias = "max_output_tokens",
116 skip_serializing_if = "Option::is_none"
117 )]
118 pub max_output_tokens: Option<std::num::NonZeroU32>,
119 pub content: String,
120 #[serde(default, skip_serializing_if = "Vec::is_empty")]
121 pub images: Vec<crate::runtime::RuntimeImageInput>,
122 /// Effective mode for this turn (`"plan" | "agent" | "operate"` etc).
123 #[serde(default = "default_mode")]
124 pub mode: String,
125 /// Optional explicit route/model the caller resolved already (mirrors
126 /// `ResolvedRuntimeRoute` in `crates_tui::route_runtime`). `None` means
127 /// "use the thread's current route".
128 #[serde(default, skip_serializing_if = "Option::is_none")]
129 pub model: Option<String>,
130 #[serde(default, skip_serializing_if = "Option::is_none")]
131 pub model_provider: Option<String>,
132 /// Tool restriction from slash-command frontmatter.
133 #[serde(default)]
134 pub allowed_tools: Option<Vec<String>>,
135 /// Runtime-supplied dynamic tools for this turn only.
136 #[serde(default)]
137 pub dynamic_tools: Vec<DynamicToolSpec>,
138 /// Structural input provenance — only `external_user` may inherit
139 /// YOLO/auto-approval authority (mirrors `UserInputProvenance`).
140 #[serde(default = "default_provenance")]
141 pub provenance: String,
142 /// Compaction policy carried atomically with the route receipt.
143 /// Boxed only to keep the enum small; the wire shape is unchanged.
144 #[serde(default, skip_serializing_if = "Option::is_none")]
145 pub compaction: Option<Box<CompactionPolicy>>,
146 #[serde(default, skip_serializing_if = "Option::is_none")]
147 pub goal_objective: Option<String>,
148 #[serde(default, skip_serializing_if = "Option::is_none")]
149 pub goal_token_budget: Option<u32>,
150 /// `active | paused | complete | blocked`.
151 #[serde(default = "default_goal_status")]
152 pub goal_status: String,
153 /// `"off" | "low" | "medium" | "high" | "max"`; `None` = provider default.
154 #[serde(default, skip_serializing_if = "Option::is_none")]
155 pub reasoning_effort: Option<String>,
156 #[serde(default)]
157 pub reasoning_effort_auto: bool,
158 #[serde(default)]
159 pub auto_model: bool,
160 #[serde(default)]
161 pub allow_shell: bool,
162 #[serde(default)]
163 pub trust_mode: bool,
164 #[serde(default)]
165 pub auto_approve: bool,
166 /// `auto | bypass | suggest | never`.
167 #[serde(default = "default_approval_mode")]
168 pub approval_mode: String,
169 #[serde(default)]
170 pub translation_enabled: bool,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
172 pub verbosity: Option<String>,
173 }
174
175 /// Operations that can be submitted to the core engine. This is the
176 /// protocol view of `crates/tui/src/core/ops::Op` — same lifecycle,
177 /// same provenance gate — but serializable and free of `mpsc` / `oneshot`
178 /// fields. In-process callers convert at the boundary; out-of-process
179 /// callers send the JSON directly.
180 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
181 #[serde(tag = "kind", rename_all = "snake_case")]
182 pub enum Op {
183 /// Drive one model turn: `role=user` content plus the resolved route
184 /// receipt the engine will freeze at the client-freeze boundary. Headless
185 /// and TUI must produce byte-identical `MessageRequest`s for identical
186 /// `Op::SendMessage` payloads.
187 SendMessage(TurnSpec),
188
189 /// Steer an in-flight turn with additional user content (drains into
190 /// the turn loop's `rx_steer` channel).
191 Steer {
192 content: String,
193 },
194
195 /// Re-check and dispatch a goal continuation (synthetic turn that
196 /// continues the same logical goal run).
197 ContinueGoal {
198 #[serde(default)]
199 dynamic_tools: Vec<DynamicToolSpec>,
200 /// Engine-owned coalescing token; direct callers send `None`.
201 ///
202 /// Enforced, not merely documented: the engine mints these on an
203 /// in-process channel (`tx_op.try_send`) and they never cross serde,
204 /// so any value arriving through deserialization came from outside.
205 /// The IDs are predictable per-session counters, and a matching guess
206 /// is treated as an already-delayed internal token -- it would consume
207 /// the pending schedule and skip the host-injected quiet period. Wire
208 /// input is therefore always dropped to `None`.
209 #[serde(
210 default,
211 skip_serializing_if = "Option::is_none",
212 deserialize_with = "drop_engine_owned_schedule_id"
213 )]
214 engine_schedule_id: Option<u64>,
215 },
216
217 /// Execute a local composer shell command without a model turn.
218 RunShellCommand {
219 command: String,
220 #[serde(default = "default_mode")]
221 mode: String,
222 #[serde(default)]
223 allow_shell: bool,
224 #[serde(default)]
225 trust_mode: bool,
226 #[serde(default)]
227 auto_approve: bool,
228 #[serde(default = "default_approval_mode")]
229 approval_mode: String,
230 },
231
232 /// Set goal status without dispatching a model turn.
233 SetGoalStatus {
234 status: String,
235 #[serde(default)]
236 clear: bool,
237 #[serde(default, skip_serializing_if = "Option::is_none")]
238 goal_id: Option<String>,
239 },
240
241 /// Set (or replace) the active goal objective and start goal work.
242 SetGoalObjective {
243 objective: String,
244 #[serde(default, skip_serializing_if = "Option::is_none")]
245 token_budget: Option<u32>,
246 #[serde(default, skip_serializing_if = "Option::is_none")]
247 goal_id: Option<String>,
248 },
249
250 Cancel,
251 Shutdown,
252
253 /// Describe the exact request the next turn would send without sending it
254 /// (`/dryrun` / `/preview-request`, #1004). Headless and TUI must render
255 /// identical manifests for identical inputs.
256 PreviewOutboundRequest {
257 #[serde(default)]
258 json: bool,
259 #[serde(default)]
260 base_prompt_only: bool,
261 #[serde(default = "default_mode")]
262 mode: String,
263 #[serde(default)]
264 allow_shell: bool,
265 #[serde(default)]
266 trust_mode: bool,
267 #[serde(default)]
268 auto_approve: bool,
269 #[serde(default = "default_approval_mode")]
270 approval_mode: String,
271 #[serde(default)]
272 allowed_tools: Option<Vec<String>>,
273 #[serde(default)]
274 dynamic_tools: Vec<DynamicToolSpec>,
275 #[serde(default = "default_provenance")]
276 provenance: String,
277 /// The model selector the user chose (`auto` when auto routing).
278 #[serde(default)]
279 requested_model: String,
280 #[serde(default)]
281 requested_reasoning: String,
282 #[serde(default)]
283 auto_model: bool,
284 #[serde(default)]
285 hypothetical_prompt_supplied: bool,
286 /// Model-facing text of the hypothetical next message, when the host
287 /// planner resolved one.
288 #[serde(default, skip_serializing_if = "Option::is_none")]
289 hypothetical_prompt: Option<String>,
290 /// Why no exact next turn exists, when it does not.
291 #[serde(default, skip_serializing_if = "Option::is_none")]
292 unresolved: Option<String>,
293 },
294
295 ListSubAgents,
296 /// Inspect live children and pending handbacks at the Engine's idle
297 /// boundary. The host owns the response channel; it is not wire input.
298 GetSubAgentSettlement,
299 CancelSubAgent {
300 agent_id: String,
301 },
302 FollowUpSubAgent {
303 agent_id: String,
304 text: String,
305 },
306
307 ChangeMode {
308 #[serde(default = "default_mode")]
309 mode: String,
310 #[serde(default)]
311 allow_shell: bool,
312 #[serde(default)]
313 trust_mode: bool,
314 #[serde(default)]
315 auto_approve: bool,
316 #[serde(default = "default_approval_mode")]
317 approval_mode: String,
318 #[serde(default, skip_serializing_if = "Option::is_none")]
319 configured_sandbox_mode: Option<String>,
320 },
321
322 SetModel {
323 model: String,
324 #[serde(default = "default_mode")]
325 mode: String,
326 #[serde(default, skip_serializing_if = "Option::is_none")]
327 route_limits: Option<RouteLimits>,
328 },
329
330 SetCompaction {
331 config: CompactionPolicy,
332 },
333
334 /// Legacy serialized permission-update operation, retained for wire
335 /// compatibility. The in-process runtime now publishes directly through
336 /// its shared policy store instead of queuing a second replacement.
337 SetPermissionRuleset {
338 ruleset: Value,
339 },
340
341 SetStreamChunkTimeout {
342 timeout_secs: u64,
343 },
344
345 SetSubagentRuntimeConfig {
346 enabled: bool,
347 max_subagents: u64,
348 launch_concurrency: u64,
349 max_spawn_depth: u32,
350 api_timeout_secs: u64,
351 heartbeat_timeout_secs: u64,
352 },
353
354 /// `SearchProvider` in snake_case.
355 SetSearchProvider {
356 provider: String,
357 },
358
359 /// Replace the engine's merged Fleet roster. Only the roster's identity
360 /// crosses: member ids in precedence order plus the load state.
361 SetFleetRoster {
362 member_ids: Vec<String>,
363 #[serde(default)]
364 exact_selection: bool,
365 #[serde(default, skip_serializing_if = "Option::is_none")]
366 load_error: Option<String>,
367 },
368
369 /// Sync engine session state (resume/load). `messages` are transcript
370 /// `Message`s and `system_prompt` a `SystemPrompt`, both serialized.
371 SyncSession {
372 #[serde(default, skip_serializing_if = "Option::is_none")]
373 engine_session_id: Option<String>,
374 messages: Vec<Value>,
375 #[serde(default, skip_serializing_if = "Option::is_none")]
376 system_prompt: Option<Value>,
377 #[serde(default)]
378 system_prompt_override: bool,
379 model: String,
380 workspace: PathBuf,
381 #[serde(default = "default_mode")]
382 mode: String,
383 },
384
385 /// Rewind an unchanged conversation. `expected` is the complete observed
386 /// session snapshot; a receiver must compare it before changing history.
387 /// The success/refusal receipt travels out-of-band.
388 RewindConversation {
389 expected: Value,
390 messages: Vec<Value>,
391 },
392
393 /// Run context compaction on one exact provider route.
394 CompactContext {
395 id: String,
396 model: String,
397 model_provider: String,
398 compaction: CompactionPolicy,
399 },
400
401 CancelCompaction {
402 id: String,
403 },
404
405 /// Request a session snapshot; the reply travels out-of-band.
406 GetSessionSnapshot,
407 /// Request the live context-window budget for the session's route; the
408 /// reply travels out-of-band.
409 GetContextBudget,
410 /// Request provider concurrency state; the reply travels out-of-band.
411 GetProviderRuntimeStatus,
412 /// Populate the engine-owned MCP pool once at boot; reply out-of-band.
413 BootstrapMcp,
414 RetryMcpServer {
415 name: String,
416 },
417 ReloadMcp {
418 config_path: PathBuf,
419 },
420
421 PurgeContext,
422 EditLastTurn {
423 new_message: String,
424 },
425 SetAdvisorEnabled {
426 enabled: bool,
427 },
428 }
429
430 fn default_mode() -> String {
431 "agent".to_string()
432 }
433
434 fn default_provenance() -> String {
435 "external_user".to_string()
436 }
437
438 fn default_goal_status() -> String {
439 "active".to_string()
440 }
441
442 fn default_approval_mode() -> String {
443 "suggest".to_string()
444 }
445
446 /// Every wire tag `Op` can carry, in declaration order.
447 pub const OP_KINDS: &[&str] = &[
448 "send_message",
449 "steer",
450 "continue_goal",
451 "run_shell_command",
452 "set_goal_status",
453 "set_goal_objective",
454 "cancel",
455 "shutdown",
456 "preview_outbound_request",
457 "list_sub_agents",
458 "get_sub_agent_settlement",
459 "cancel_sub_agent",
460 "follow_up_sub_agent",
461 "change_mode",
462 "set_model",
463 "set_compaction",
464 "set_permission_ruleset",
465 "set_stream_chunk_timeout",
466 "set_subagent_runtime_config",
467 "set_search_provider",
468 "set_fleet_roster",
469 "sync_session",
470 "rewind_conversation",
471 "compact_context",
472 "cancel_compaction",
473 "get_session_snapshot",
474 "get_context_budget",
475 "get_provider_runtime_status",
476 "bootstrap_mcp",
477 "retry_mcp_server",
478 "reload_mcp",
479 "purge_context",
480 "edit_last_turn",
481 "set_advisor_enabled",
482 ];
483
484 impl Op {
485 #[must_use]
486 pub fn is_send_message(&self) -> bool {
487 matches!(self, Self::SendMessage(_))
488 }
489
490 #[must_use]
491 pub fn kind_str(&self) -> &'static str {
492 match self {
493 Self::SendMessage(_) => "send_message",
494 Self::Steer { .. } => "steer",
495 Self::ContinueGoal { .. } => "continue_goal",
496 Self::RunShellCommand { .. } => "run_shell_command",
497 Self::SetGoalStatus { .. } => "set_goal_status",
498 Self::SetGoalObjective { .. } => "set_goal_objective",
499 Self::Cancel => "cancel",
500 Self::Shutdown => "shutdown",
501 Self::PreviewOutboundRequest { .. } => "preview_outbound_request",
502 Self::ListSubAgents => "list_sub_agents",
503 Self::GetSubAgentSettlement => "get_sub_agent_settlement",
504 Self::CancelSubAgent { .. } => "cancel_sub_agent",
505 Self::FollowUpSubAgent { .. } => "follow_up_sub_agent",
506 Self::ChangeMode { .. } => "change_mode",
507 Self::SetModel { .. } => "set_model",
508 Self::SetCompaction { .. } => "set_compaction",
509 Self::SetPermissionRuleset { .. } => "set_permission_ruleset",
510 Self::SetStreamChunkTimeout { .. } => "set_stream_chunk_timeout",
511 Self::SetSubagentRuntimeConfig { .. } => "set_subagent_runtime_config",
512 Self::SetSearchProvider { .. } => "set_search_provider",
513 Self::SetFleetRoster { .. } => "set_fleet_roster",
514 Self::SyncSession { .. } => "sync_session",
515 Self::RewindConversation { .. } => "rewind_conversation",
516 Self::CompactContext { .. } => "compact_context",
517 Self::CancelCompaction { .. } => "cancel_compaction",
518 Self::GetSessionSnapshot => "get_session_snapshot",
519 Self::GetContextBudget => "get_context_budget",
520 Self::GetProviderRuntimeStatus => "get_provider_runtime_status",
521 Self::BootstrapMcp => "bootstrap_mcp",
522 Self::RetryMcpServer { .. } => "retry_mcp_server",
523 Self::ReloadMcp { .. } => "reload_mcp",
524 Self::PurgeContext => "purge_context",
525 Self::EditLastTurn { .. } => "edit_last_turn",
526 Self::SetAdvisorEnabled { .. } => "set_advisor_enabled",
527 }
528 }
529 }
530
531 /// Build a headless `SendMessage` envelope with fresh ids. This is the
532 /// one-line helper every headless caller (CLI `exec`, app-server, tests)
533 /// uses so TUI and headless start a session identically.
534 #[must_use]
535 pub fn headless_send_message_op(thread_id: ThreadId, content: impl Into<String>) -> OpEnvelope {
536 OpEnvelope {
537 op_id: format!("op-{}", uuid::Uuid::new_v4()),
538 thread_id: thread_id.clone(),
539 session_id: SessionId::new(),
540 op: Op::SendMessage(TurnSpec {
541 max_output_tokens: None,
542 content: content.into(),
543 images: Vec::new(),
544 mode: default_mode(),
545 model: None,
546 model_provider: None,
547 allowed_tools: None,
548 dynamic_tools: Vec::new(),
549 provenance: default_provenance(),
550 compaction: None,
551 goal_objective: None,
552 goal_token_budget: None,
553 goal_status: default_goal_status(),
554 reasoning_effort: None,
555 reasoning_effort_auto: false,
556 auto_model: false,
557 allow_shell: false,
558 trust_mode: false,
559 auto_approve: false,
560 approval_mode: default_approval_mode(),
561 translation_enabled: false,
562 verbosity: None,
563 }),
564 }
565 }
566
567 #[cfg(test)]
568 mod tests {
569 use super::*;
570 use serde_json::json;
571
572 fn policy() -> CompactionPolicy {
573 CompactionPolicy {
574 enabled: true,
575 token_threshold: 100_000,
576 model: "deepseek-chat".into(),
577 image_input: "unknown".into(),
578 effective_context_window: Some(128_000),
579 cache_summary: true,
580 focus: None,
581 runtime_cost_owner: None,
582 workspace: None,
583 }
584 }
585
586 /// One instance of every variant, in declaration order.
587 fn every_variant() -> Vec<Op> {
588 vec![
589 headless_send_message_op(ThreadId::new(), "hello").op,
590 Op::Steer {
591 content: "more".into(),
592 },
593 Op::ContinueGoal {
594 dynamic_tools: vec![],
595 // Engine-owned: deliberately not round-trippable. See
596 // `wire_supplied_engine_schedule_id_is_dropped`.
597 engine_schedule_id: None,
598 },
599 Op::RunShellCommand {
600 command: "ls".into(),
601 mode: "agent".into(),
602 allow_shell: true,
603 trust_mode: false,
604 auto_approve: false,
605 approval_mode: "suggest".into(),
606 },
607 Op::SetGoalStatus {
608 goal_id: None,
609 status: "paused".into(),
610 clear: false,
611 },
612 Op::SetGoalObjective {
613 goal_id: None,
614 objective: "ship".into(),
615 token_budget: Some(10),
616 },
617 Op::Cancel,
618 Op::Shutdown,
619 Op::PreviewOutboundRequest {
620 json: true,
621 base_prompt_only: false,
622 mode: "plan".into(),
623 allow_shell: false,
624 trust_mode: false,
625 auto_approve: false,
626 approval_mode: "auto".into(),
627 allowed_tools: Some(vec!["read_file".into()]),
628 dynamic_tools: vec![],
629 provenance: "external_user".into(),
630 requested_model: "auto".into(),
631 requested_reasoning: "high".into(),
632 auto_model: true,
633 hypothetical_prompt_supplied: true,
634 hypothetical_prompt: Some("hi".into()),
635 unresolved: None,
636 },
637 Op::ListSubAgents,
638 Op::GetSubAgentSettlement,
639 Op::CancelSubAgent {
640 agent_id: "a1".into(),
641 },
642 Op::FollowUpSubAgent {
643 agent_id: "a1".into(),
644 text: "go".into(),
645 },
646 Op::ChangeMode {
647 mode: "operate".into(),
648 allow_shell: true,
649 trust_mode: true,
650 auto_approve: false,
651 approval_mode: "bypass".into(),
652 configured_sandbox_mode: Some("workspace-write".into()),
653 },
654 Op::SetModel {
655 model: "m".into(),
656 mode: "agent".into(),
657 route_limits: Some(RouteLimits {
658 context_tokens: Some(1),
659 input_tokens: None,
660 output_tokens: None,
661 }),
662 },
663 Op::SetCompaction { config: policy() },
664 Op::SetPermissionRuleset {
665 ruleset: json!({"rules": []}),
666 },
667 Op::SetStreamChunkTimeout { timeout_secs: 30 },
668 Op::SetSubagentRuntimeConfig {
669 enabled: true,
670 max_subagents: 4,
671 launch_concurrency: 2,
672 max_spawn_depth: 1,
673 api_timeout_secs: 60,
674 heartbeat_timeout_secs: 10,
675 },
676 Op::SetSearchProvider {
677 provider: "brave".into(),
678 },
679 Op::SetFleetRoster {
680 member_ids: vec!["scout".into()],
681 exact_selection: false,
682 load_error: None,
683 },
684 Op::SyncSession {
685 engine_session_id: Some("s".into()),
686 messages: vec![json!({"role": "user", "content": []})],
687 system_prompt: None,
688 system_prompt_override: false,
689 model: "m".into(),
690 workspace: PathBuf::from("/ws"),
691 mode: "agent".into(),
692 },
693 Op::RewindConversation {
694 expected: json!({"session_id": "s", "messages": []}),
695 messages: vec![],
696 },
697 Op::CompactContext {
698 id: "cmp-1".into(),
699 model: "m".into(),
700 model_provider: "deepseek".into(),
701 compaction: policy(),
702 },
703 Op::CancelCompaction { id: "cmp-1".into() },
704 Op::GetSessionSnapshot,
705 Op::GetContextBudget,
706 Op::GetProviderRuntimeStatus,
707 Op::BootstrapMcp,
708 Op::RetryMcpServer { name: "fs".into() },
709 Op::ReloadMcp {
710 config_path: PathBuf::from("/tmp/mcp.json"),
711 },
712 Op::PurgeContext,
713 Op::EditLastTurn {
714 new_message: "again".into(),
715 },
716 Op::SetAdvisorEnabled { enabled: true },
717 ]
718 }
719
720 #[test]
721 fn every_variant_is_listed_once() {
722 let kinds: Vec<&str> = every_variant().iter().map(Op::kind_str).collect();
723 assert_eq!(kinds, OP_KINDS, "OP_KINDS must list every variant in order");
724 }
725
726 #[test]
727 fn wire_supplied_engine_schedule_id_is_dropped() {
728 // The engine mints these on an in-process channel, so a value arriving
729 // through serde came from an out-of-process caller. Honouring it would
730 // let a guessed counter consume the pending schedule and skip the
731 // host-injected quiet period.
732 let op: Op = serde_json::from_value(json!({
733 "kind": "continue_goal",
734 "dynamic_tools": [],
735 "engine_schedule_id": 3,
736 }))
737 .expect("continue_goal with a wire-supplied schedule id still parses");
738 match op {
739 Op::ContinueGoal {
740 engine_schedule_id, ..
741 } => assert_eq!(
742 engine_schedule_id, None,
743 "engine_schedule_id must never be settable from the wire"
744 ),
745 other => panic!("expected ContinueGoal, got {other:?}"),
746 }
747 }
748
749 #[test]
750 fn every_variant_round_trips_and_tags_by_kind() {
751 for op in every_variant() {
752 let value = serde_json::to_value(&op).unwrap();
753 assert_eq!(value["kind"], op.kind_str(), "{op:?}");
754 let back: Op = serde_json::from_value(value).unwrap();
755 assert_eq!(back, op);
756 }
757 }
758
759 #[test]
760 fn op_envelope_roundtrip() {
761 let env = headless_send_message_op(ThreadId::new(), "hello");
762 let json = serde_json::to_string(&env).unwrap();
763 let back: OpEnvelope = serde_json::from_str(&json).unwrap();
764 assert_eq!(back.thread_id, env.thread_id);
765 assert!(back.op.is_send_message());
766 }
767
768 #[test]
769 fn minimal_send_message_json_still_parses_with_defaults() {
770 // The pre-A1 wire shape: only `content` plus the tag.
771 let op: Op = serde_json::from_value(json!({
772 "kind": "send_message",
773 "content": "hello"
774 }))
775 .unwrap();
776 let Op::SendMessage(TurnSpec {
777 mode,
778 provenance,
779 goal_status,
780 approval_mode,
781 ..
782 }) = op
783 else {
784 panic!("expected send_message");
785 };
786 assert_eq!(mode, "agent");
787 assert_eq!(provenance, "external_user");
788 assert_eq!(goal_status, "active");
789 assert_eq!(approval_mode, "suggest");
790 }
791
792 #[test]
793 fn steer_roundtrip() {
794 let op = Op::Steer {
795 content: "more".into(),
796 };
797 let json = serde_json::to_string(&op).unwrap();
798 let back: Op = serde_json::from_str(&json).unwrap();
799 assert_eq!(back.kind_str(), "steer");
800 }
801 }
802
802 lines RUST