返回 CodeWhale
event_msg.rs
根目录 / crates / protocol / src / event_msg.rs
1 //! `EventMsg`-out API in `crates/protocol` (issue #5261, Phase A1 of the
2 //! core/protocol extraction spec).
3 //!
4 //! Mirrors `crates/tui/src/core/events::Event` variant-by-variant as a
5 //! serializable protocol. The TUI's `rx_event` / `Event` channel, the
6 //! app-server's SSE stream, and the CLI's `stream-json` output all speak this
7 //! one type so headless and TUI observe byte-identical event shapes for the
8 //! same `Op`.
9 //!
10 //! Parity is compile-enforced from the engine side:
11 //! `crates/tui/src/core/protocol_parity.rs` matches every engine `Event`
12 //! variant exhaustively into an `EventMsg` (`protocol_covers_engine_events`).
13 //! Adding an engine variant without a twin here fails to compile.
14 //!
15 //! Payload fidelity rules for this phase:
16 //!
17 //! - Scalars, ids, lifecycle enums, route/billing receipts, tool outcomes,
18 //! MCP snapshots, approvals, and gate decisions are typed here.
19 //! - Deep domain payloads whose canonical serde type still lives above this
20 //! crate (goal snapshots, sub-agent results, coordination projections,
21 //! mailbox messages, transcript messages, tool catalogs, tool inspection
22 //! snapshots, workflow UI events) cross as `serde_json::Value` produced by
23 //! that type's own `Serialize`. They are typed in later phases; the variant
24 //! and its field names are already stable.
25 //! - Engine-only handles (`oneshot`/`Notify`, `Arc<HookExecutor>`) never
26 //! cross. A variant that carried one is projected without it.
27
28 use std::collections::BTreeMap;
29 use std::path::PathBuf;
30
31 use chrono::{DateTime, Utc};
32 use serde::{Deserialize, Serialize};
33 use serde_json::Value;
34
35 use crate::ResponseChannel;
36 use crate::UserInputQuestionEvent;
37 use crate::ids::{SessionId, ThreadId};
38
39 /// Final status for a turn (`TurnComplete.status`).
40 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
41 #[serde(rename_all = "snake_case")]
42 pub enum TurnOutcomeStatus {
43 Completed,
44 Interrupted,
45 Failed,
46 }
47
48 impl TurnOutcomeStatus {
49 #[must_use]
50 pub fn as_str(self) -> &'static str {
51 match self {
52 Self::Completed => "completed",
53 Self::Interrupted => "interrupted",
54 Self::Failed => "failed",
55 }
56 }
57 }
58
59 /// Token usage reported by a provider for one model call or one whole turn.
60 /// Field-for-field twin of the engine's `Usage`; `None` means the provider
61 /// did not report the fact, never zero.
62 #[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
63 pub struct TokenUsage {
64 pub input_tokens: u32,
65 pub output_tokens: u32,
66 #[serde(default, skip_serializing_if = "Option::is_none")]
67 pub prompt_cache_hit_tokens: Option<u32>,
68 #[serde(default, skip_serializing_if = "Option::is_none")]
69 pub prompt_cache_miss_tokens: Option<u32>,
70 #[serde(default, skip_serializing_if = "Option::is_none")]
71 pub prompt_cache_write_tokens: Option<u32>,
72 #[serde(default, skip_serializing_if = "Option::is_none")]
73 pub reasoning_tokens: Option<u32>,
74 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub reasoning_replay_tokens: Option<u32>,
76 #[serde(default, skip_serializing_if = "Option::is_none")]
77 pub code_execution_requests: Option<u32>,
78 #[serde(default, skip_serializing_if = "Option::is_none")]
79 pub tool_search_requests: Option<u32>,
80 }
81
82 /// Secret-free proof of the base route a turn's client was installed on.
83 /// The credential generation digest is redacted by design on the engine side
84 /// and never crosses; only its presence does.
85 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
86 pub struct TurnRouteReceipt {
87 pub provider: String,
88 pub provider_identity: String,
89 pub wire_model: String,
90 pub endpoint_identity: String,
91 pub credential_generation_present: bool,
92 }
93
94 /// Credential/pay-mode product truth captured at the client-freeze boundary.
95 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
96 #[serde(tag = "kind", rename_all = "snake_case")]
97 pub enum RouteProduct {
98 /// No product fact was captured. Not a licence to guess.
99 Unproven,
100 /// Subscription-backed with this user-facing quota label.
101 Subscription { label: String },
102 /// Bills per token.
103 Metered,
104 }
105
106 /// Billing evidence captured at application admission before the provider permit.
107 /// This does not attest network delivery. Absent before admission.
108 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
109 pub struct RouteBillingEnvelope {
110 #[serde(default, skip_serializing_if = "Option::is_none")]
111 pub openrouter_vendor: Option<String>,
112 #[serde(default, skip_serializing_if = "Option::is_none")]
113 pub billing_surface: Option<String>,
114 #[serde(default, skip_serializing_if = "Option::is_none")]
115 pub endpoint_fingerprint: Option<String>,
116 /// Validated frozen provider-live quote serialized by the runtime owner.
117 #[serde(default, skip_serializing_if = "Option::is_none")]
118 pub provider_live_pricing: Option<Value>,
119 /// `RouteBillingMode` in snake_case.
120 pub billing_mode: String,
121 pub dispatched_at: DateTime<Utc>,
122 }
123
124 /// Provider/model route resolved for a model-backed turn.
125 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
126 pub struct TurnRoute {
127 /// `ProviderKind` key (`deepseek`, `openai`, `custom`, ...).
128 pub provider: String,
129 /// Exact non-secret configured route key.
130 pub provider_identity: String,
131 pub model: String,
132 pub auto_model: bool,
133 #[serde(default, skip_serializing_if = "Option::is_none")]
134 pub receipt: Option<TurnRouteReceipt>,
135 #[serde(default, skip_serializing_if = "Option::is_none")]
136 pub billing: Option<RouteBillingEnvelope>,
137 /// Endpoint the client was frozen against, verbatim. Empty when unknown.
138 pub base_url: String,
139 pub billing_product: RouteProduct,
140 }
141
142 /// Structured error surfaced by a tool execution.
143 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144 #[serde(tag = "kind", rename_all = "snake_case")]
145 pub enum ToolCallError {
146 InvalidInput { message: String },
147 MissingField { field: String },
148 PathEscape { path: PathBuf },
149 ExecutionFailed { message: String },
150 Timeout { seconds: u64 },
151 Cancelled { message: String },
152 NotAvailable { message: String },
153 PermissionDenied { message: String },
154 }
155
156 /// Outcome of a tool call: the engine's `Result<ToolResult, ToolError>`.
157 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
158 #[serde(tag = "outcome", rename_all = "snake_case")]
159 pub enum ToolCallOutcome {
160 Ok {
161 content: String,
162 success: bool,
163 #[serde(default, skip_serializing_if = "Option::is_none")]
164 metadata: Option<Value>,
165 },
166 Err {
167 error: ToolCallError,
168 },
169 }
170
171 /// Lifecycle metadata paired with a human-readable `AgentProgress` message.
172 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
173 pub struct AgentProgressActivity {
174 /// `AgentWorkerStatus` in snake_case.
175 pub worker_status: String,
176 #[serde(default, skip_serializing_if = "Option::is_none")]
177 pub step: Option<u32>,
178 #[serde(default, skip_serializing_if = "Option::is_none")]
179 pub tool_name: Option<String>,
180 }
181
182 /// Receipt for an operator follow-up to a child agent.
183 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
184 #[serde(tag = "outcome", rename_all = "snake_case")]
185 pub enum SubAgentFollowUpOutcome {
186 Ok {
187 agent_id: String,
188 target_agent_id: String,
189 delivered: bool,
190 resumed: bool,
191 note: String,
192 },
193 Err {
194 reason: String,
195 },
196 }
197
198 /// One row of the receipts-only agent roster.
199 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
200 pub struct AgentRosterRow {
201 pub worker_id: String,
202 pub display_name: String,
203 pub model: String,
204 /// Coarse rail state:
205 /// `running | waiting | parked | done | failed | cancelled`.
206 pub state: String,
207 /// `AgentWorkerStatus` in snake_case.
208 pub status: String,
209 #[serde(default, skip_serializing_if = "Option::is_none")]
210 pub activity: Option<String>,
211 #[serde(default, skip_serializing_if = "Option::is_none")]
212 pub millis: Option<u64>,
213 #[serde(default, skip_serializing_if = "Option::is_none")]
214 pub input_tokens: Option<u64>,
215 #[serde(default, skip_serializing_if = "Option::is_none")]
216 pub output_tokens: Option<u64>,
217 #[serde(default, skip_serializing_if = "Option::is_none")]
218 pub cost_microusd: Option<u64>,
219 pub steps_taken: u32,
220 #[serde(default, skip_serializing_if = "Option::is_none")]
221 pub parent_run_id: Option<String>,
222 pub run_id: String,
223 }
224
225 /// One discovered MCP tool / resource / prompt.
226 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227 pub struct McpDiscoveredItem {
228 pub name: String,
229 pub model_name: String,
230 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub description: Option<String>,
232 }
233
234 /// One configured MCP server as seen by the engine-owned pool.
235 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
236 pub struct McpServerSnapshot {
237 pub name: String,
238 pub enabled: bool,
239 pub required: bool,
240 pub transport: String,
241 pub command_or_url: String,
242 pub connect_timeout: u64,
243 pub execute_timeout: u64,
244 pub read_timeout: u64,
245 pub connected: bool,
246 #[serde(default, skip_serializing_if = "Option::is_none")]
247 pub error: Option<String>,
248 /// `advertised | legacy_fallback | not_observed`.
249 pub capability_metadata: String,
250 pub tools: Vec<McpDiscoveredItem>,
251 pub resources: Vec<McpDiscoveredItem>,
252 pub prompts: Vec<McpDiscoveredItem>,
253 }
254
255 /// Engine-owned MCP pool snapshot.
256 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
257 pub struct McpManagerSnapshot {
258 pub config_path: PathBuf,
259 pub config_exists: bool,
260 pub reload_required: bool,
261 pub servers: Vec<McpServerSnapshot>,
262 }
263
264 /// Structured clarification request (`request_user_input`).
265 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
266 pub struct UserInputRequest {
267 pub questions: Vec<UserInputQuestionEvent>,
268 }
269
270 /// Which permission gate produced a `ToolGateDecision`.
271 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
272 #[serde(rename_all = "snake_case")]
273 pub enum ToolGate {
274 AutoReviewDeterministic,
275 AutoReviewGuardian,
276 }
277
278 /// What a permission gate decided for one proposed tool call.
279 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
280 #[serde(rename_all = "snake_case")]
281 pub enum ToolGateVerdict {
282 Allowed,
283 Denied,
284 Unavailable,
285 }
286
287 /// One event emitted by the core engine to every consumer (TUI, CLI,
288 /// app-server, tests). This is the `EventMsg`-out half of the `Op`-in /
289 /// `EventMsg`-out contract: a projection of every internal engine `Event`
290 /// variant plus the thread/session ids that route it.
291 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
292 #[serde(tag = "event", rename_all = "snake_case")]
293 pub enum EventMsg {
294 /// A route compatibility check omitted tools from the provider request.
295 ToolProjectionWarning {
296 thread_id: ThreadId,
297 session_id: SessionId,
298 provider: String,
299 omitted_tool_names: Vec<String>,
300 omitted_tool_count: u64,
301 },
302
303 /// Workspace snapshots (undo) are off for this workspace. `reason` is one
304 /// rendered, localized line: the consequence, the gate that refused, and
305 /// the recovery that actually lifts *that* gate (the size cap's config key
306 /// appears only for the size gate).
307 SnapshotsDisabled {
308 thread_id: ThreadId,
309 session_id: SessionId,
310 workspace: String,
311 reason: String,
312 },
313
314 // === Streaming ===
315 MessageStarted {
316 thread_id: ThreadId,
317 session_id: SessionId,
318 index: u64,
319 },
320 /// Incremental content delta on the `text` (message) or `reasoning`
321 /// (thinking) channel.
322 ResponseDelta {
323 thread_id: ThreadId,
324 session_id: SessionId,
325 index: u64,
326 delta: String,
327 #[serde(default, skip_serializing_if = "ResponseChannel::is_text")]
328 channel: ResponseChannel,
329 },
330 MessageComplete {
331 thread_id: ThreadId,
332 session_id: SessionId,
333 index: u64,
334 },
335 ThinkingStarted {
336 thread_id: ThreadId,
337 session_id: SessionId,
338 index: u64,
339 },
340 ThinkingComplete {
341 thread_id: ThreadId,
342 session_id: SessionId,
343 index: u64,
344 },
345
346 // === Tools ===
347 ToolCallStarted {
348 thread_id: ThreadId,
349 session_id: SessionId,
350 tool_call_id: String,
351 tool_name: String,
352 input: Value,
353 },
354 /// Liveness pulse while a tool future remains pending. Carries no output.
355 ToolCallHeartbeat {
356 thread_id: ThreadId,
357 session_id: SessionId,
358 },
359 ToolExecutionStarted {
360 thread_id: ThreadId,
361 session_id: SessionId,
362 tool_call_id: String,
363 },
364 ToolResultContent {
365 thread_id: ThreadId,
366 session_id: SessionId,
367 tool_call_id: String,
368 blocks: Value,
369 },
370 ToolCallComplete {
371 thread_id: ThreadId,
372 session_id: SessionId,
373 tool_call_id: String,
374 tool_name: String,
375 result: ToolCallOutcome,
376 },
377 /// Trusted Engine-owned activity for an operation that passed dispatch
378 /// and authority checks. No tool name, arguments, command, or result is
379 /// included in the pet-facing activity contract.
380 ///
381 /// Reserved on the wire: `protocol_parity` maps the engine event, but no
382 /// runtime thread emits the pair yet, so Runtime API and GPUI clients do
383 /// not receive it. The shared pet (`pet_watch`) is the only consumer today.
384 OperationActivityStarted {
385 thread_id: ThreadId,
386 session_id: SessionId,
387 span_id: String,
388 activity_kind: crate::engine_owner::OwnerActivityKind,
389 },
390 OperationActivityCompleted {
391 thread_id: ThreadId,
392 session_id: SessionId,
393 span_id: String,
394 activity_kind: crate::engine_owner::OwnerActivityKind,
395 outcome: crate::engine_owner::OwnerOperationOutcome,
396 },
397
398 // === Turn lifecycle ===
399 TurnStarted {
400 thread_id: ThreadId,
401 session_id: SessionId,
402 turn_id: String,
403 created_at: DateTime<Utc>,
404 #[serde(default, skip_serializing_if = "Option::is_none")]
405 route: Option<TurnRoute>,
406 /// Host submission correlation echo: the token an in-process host
407 /// stamped on its `SendMessage`/`EditLastTurn` op, echoed verbatim by
408 /// that turn's start; `None` for every engine self-started turn.
409 /// Additive and default-absent on the wire. Only in-process engine
410 /// handles can stamp a token — the wire `Op` carries no correlation
411 /// field — so wire submitters only ever observe `None` here.
412 #[serde(default, skip_serializing_if = "Option::is_none")]
413 submission_id: Option<String>,
414 },
415 /// Bounded tool-field projection from a prepared model-client request
416 /// (`ToolInspectionSnapshot` serialized).
417 ToolRequestSnapshot {
418 thread_id: ThreadId,
419 session_id: SessionId,
420 snapshot: Value,
421 },
422 /// A workspace snapshot the engine took for the running turn
423 /// (`WorkspaceSnapshotRef` serialized: `kind`, `snapshot_id`, `tree_id`,
424 /// `session_id`, optional `tool_call_id`, `write_paths` and
425 /// `changed_paths`).
426 WorkspaceSnapshotTaken {
427 thread_id: ThreadId,
428 session_id: SessionId,
429 snapshot: Value,
430 },
431 /// Immutable billing route captured at application admission.
432 RouteDispatched {
433 thread_id: ThreadId,
434 session_id: SessionId,
435 turn_id: String,
436 route: TurnRoute,
437 },
438 TurnComplete {
439 thread_id: ThreadId,
440 session_id: SessionId,
441 /// The engine's `TurnComplete` carries no turn id; the emitter fills
442 /// it from the envelope when it knows it.
443 #[serde(default, skip_serializing_if = "Option::is_none")]
444 turn_id: Option<String>,
445 status: TurnOutcomeStatus,
446 #[serde(default, skip_serializing_if = "Option::is_none")]
447 error: Option<String>,
448 usage: TokenUsage,
449 /// Parent-route subset; absent in legacy events, whose split is unknown.
450 #[serde(default, skip_serializing_if = "Option::is_none")]
451 parent_route_usage: Option<TokenUsage>,
452 #[serde(default)]
453 routed_usage_dropped_records: u64,
454 /// Tool catalog sent with this turn's model request (`Tool` serialized).
455 #[serde(default, skip_serializing_if = "Option::is_none")]
456 tool_catalog: Option<Vec<Value>>,
457 #[serde(default, skip_serializing_if = "Option::is_none")]
458 base_url: Option<String>,
459 },
460 /// Usage for one model call within the turn.
461 TurnUsage {
462 #[serde(
463 default,
464 rename = "maxOutputTokens",
465 skip_serializing_if = "Option::is_none"
466 )]
467 max_output_tokens: Option<u32>,
468 thread_id: ThreadId,
469 session_id: SessionId,
470 usage: TokenUsage,
471 duration_ms: u64,
472 #[serde(default, skip_serializing_if = "Option::is_none")]
473 first_token_ms: Option<u64>,
474 #[serde(default, skip_serializing_if = "Option::is_none")]
475 request_ms: Option<u64>,
476 },
477
478 /// Child-call telemetry; cost belongs to its own routed receipt.
479 RoutedTurnUsage {
480 thread_id: ThreadId,
481 session_id: SessionId,
482 usage: TokenUsage,
483 duration_ms: u64,
484 #[serde(default, skip_serializing_if = "Option::is_none")]
485 first_token_ms: Option<u64>,
486 #[serde(default, skip_serializing_if = "Option::is_none")]
487 request_ms: Option<u64>,
488 },
489
490 // === Goals ===
491 /// Runtime goal state changed (`GoalSnapshot` serialized).
492 GoalUpdated {
493 thread_id: ThreadId,
494 session_id: SessionId,
495 snapshot: Value,
496 },
497 GoalContinuationWaiting {
498 thread_id: ThreadId,
499 session_id: SessionId,
500 delay_seconds: u64,
501 },
502 GoalContinuationWaitEnded {
503 thread_id: ThreadId,
504 session_id: SessionId,
505 interrupted: bool,
506 },
507
508 // === Compaction / purge ===
509 CompactionStarted {
510 thread_id: ThreadId,
511 session_id: SessionId,
512 id: String,
513 auto: bool,
514 message: String,
515 },
516 CompactionCompleted {
517 thread_id: ThreadId,
518 session_id: SessionId,
519 id: String,
520 auto: bool,
521 message: String,
522 #[serde(default, skip_serializing_if = "Option::is_none")]
523 messages_before: Option<u64>,
524 #[serde(default, skip_serializing_if = "Option::is_none")]
525 messages_after: Option<u64>,
526 #[serde(default, skip_serializing_if = "Option::is_none")]
527 summary_prompt: Option<String>,
528 #[serde(default, skip_serializing_if = "Option::is_none")]
529 post_input_tokens: Option<u64>,
530 },
531 CompactionCancelled {
532 thread_id: ThreadId,
533 session_id: SessionId,
534 id: String,
535 auto: bool,
536 message: String,
537 },
538 CompactionFailed {
539 thread_id: ThreadId,
540 session_id: SessionId,
541 id: String,
542 auto: bool,
543 message: String,
544 },
545 PurgeStarted {
546 thread_id: ThreadId,
547 session_id: SessionId,
548 message: String,
549 },
550 PurgeCompleted {
551 thread_id: ThreadId,
552 session_id: SessionId,
553 messages_before: u64,
554 messages_after: u64,
555 removed_count: u64,
556 replaced_count: u64,
557 message: String,
558 },
559 PurgeFailed {
560 thread_id: ThreadId,
561 session_id: SessionId,
562 message: String,
563 },
564
565 // === Sub-agents ===
566 AgentSpawned {
567 thread_id: ThreadId,
568 session_id: SessionId,
569 owner_session_id: String,
570 id: String,
571 prompt: String,
572 #[serde(default, skip_serializing_if = "Option::is_none")]
573 worker_status: Option<String>,
574 #[serde(default, skip_serializing_if = "Option::is_none")]
575 parent_run_id: Option<String>,
576 spawn_depth: u32,
577 model: String,
578 #[serde(default, skip_serializing_if = "Option::is_none")]
579 route_source: Option<String>,
580 /// The name the agent goes by (workflow task label, dispatch name, or
581 /// role). Absent from older producers; never the raw id.
582 #[serde(default, skip_serializing_if = "Option::is_none")]
583 display_name: Option<String>,
584 },
585 AgentProgress {
586 thread_id: ThreadId,
587 session_id: SessionId,
588 owner_session_id: String,
589 id: String,
590 status: String,
591 activity: AgentProgressActivity,
592 #[serde(default, skip_serializing_if = "Option::is_none")]
593 parent_run_id: Option<String>,
594 spawn_depth: u32,
595 },
596 AgentComplete {
597 thread_id: ThreadId,
598 session_id: SessionId,
599 owner_session_id: String,
600 id: String,
601 result: String,
602 #[serde(default, skip_serializing_if = "Option::is_none")]
603 worker_status: Option<String>,
604 #[serde(default, skip_serializing_if = "Option::is_none")]
605 parent_run_id: Option<String>,
606 #[serde(default, skip_serializing_if = "Option::is_none")]
607 spawn_depth: Option<u32>,
608 #[serde(default, skip_serializing_if = "Option::is_none")]
609 continuable: Option<bool>,
610 #[serde(default, skip_serializing_if = "Option::is_none")]
611 display_name: Option<String>,
612 },
613 SubAgentFollowUp {
614 thread_id: ThreadId,
615 session_id: SessionId,
616 owner_session_id: String,
617 agent_id: String,
618 outcome: SubAgentFollowUpOutcome,
619 },
620 /// Sub-agent listing. `agents` are `SubAgentResult`s and `coordination`
621 /// is the `CoordinationDetailProjection`, both serialized.
622 AgentList {
623 thread_id: ThreadId,
624 session_id: SessionId,
625 owner_session_id: String,
626 agents: Vec<Value>,
627 coordination: Value,
628 /// `agent_id` -> queued follow-up count; only non-zero entries.
629 #[serde(default)]
630 queued_follow_ups: BTreeMap<String, u64>,
631 roster: Vec<AgentRosterRow>,
632 },
633 /// Structured sub-agent mailbox envelope (`MailboxMessage` serialized).
634 /// Deduplicate on `(turn_id, seq)`, never `seq` alone.
635 SubAgentMailbox {
636 thread_id: ThreadId,
637 session_id: SessionId,
638 owner_session_id: String,
639 turn_id: String,
640 seq: u64,
641 message: Value,
642 },
643 /// Live workflow UI event. `ui_event` is the flattened
644 /// `{"type": ..., "at_ms": ..., ...}` object (named `event` on the
645 /// engine side; renamed here because `event` is the wire tag).
646 WorkflowUi {
647 thread_id: ThreadId,
648 session_id: SessionId,
649 owner_session_id: String,
650 run_id: String,
651 ui_event: Value,
652 },
653
654 // === System ===
655 Error {
656 thread_id: ThreadId,
657 session_id: SessionId,
658 /// `ErrorCategory` in snake_case.
659 category: String,
660 /// `ErrorSeverity` in snake_case.
661 severity: String,
662 recoverable: bool,
663 code: String,
664 message: String,
665 },
666 Status {
667 thread_id: ThreadId,
668 session_id: SessionId,
669 message: String,
670 },
671 McpSessionBoot {
672 thread_id: ThreadId,
673 session_id: SessionId,
674 generation: u64,
675 snapshot: McpManagerSnapshot,
676 connecting: Vec<String>,
677 finished: bool,
678 },
679 /// Rendered `/preview-request` manifest.
680 RequestManifestReady {
681 thread_id: ThreadId,
682 session_id: SessionId,
683 rendered: String,
684 },
685 /// Pause terminal input events for an interactive subprocess. The engine's
686 /// in-process acknowledgement handle does not cross the wire.
687 PauseEvents {
688 thread_id: ThreadId,
689 session_id: SessionId,
690 },
691 ResumeEvents {
692 thread_id: ThreadId,
693 session_id: SessionId,
694 },
695 ApprovalRequired {
696 thread_id: ThreadId,
697 session_id: SessionId,
698 id: String,
699 tool_name: String,
700 description: String,
701 input: Value,
702 approval_key: String,
703 approval_grouping_key: String,
704 #[serde(default, skip_serializing_if = "Option::is_none")]
705 intent_summary: Option<String>,
706 approval_force_prompt: bool,
707 },
708 ApprovalWithdrawn {
709 thread_id: ThreadId,
710 session_id: SessionId,
711 id: String,
712 },
713 UserInputRequired {
714 thread_id: ThreadId,
715 session_id: SessionId,
716 id: String,
717 request: UserInputRequest,
718 },
719 /// Authoritative API conversation state (`Message`s / `SystemPrompt`
720 /// serialized).
721 SessionUpdated {
722 thread_id: ThreadId,
723 session_id: SessionId,
724 engine_session_id: String,
725 messages: Vec<Value>,
726 #[serde(default, skip_serializing_if = "Option::is_none")]
727 system_prompt: Option<Value>,
728 model: String,
729 workspace: PathBuf,
730 },
731 ElevationRequired {
732 thread_id: ThreadId,
733 session_id: SessionId,
734 tool_id: String,
735 tool_name: String,
736 #[serde(default, skip_serializing_if = "Option::is_none")]
737 command: Option<String>,
738 denial_reason: String,
739 blocked_network: bool,
740 blocked_write: bool,
741 },
742 LspRepairUpdate {
743 thread_id: ThreadId,
744 session_id: SessionId,
745 diagnostics_found: u64,
746 files: u64,
747 injected: bool,
748 },
749 ToolGateDecision {
750 thread_id: ThreadId,
751 session_id: SessionId,
752 #[serde(default, skip_serializing_if = "Option::is_none")]
753 agent_id: Option<String>,
754 tool_id: String,
755 tool_name: String,
756 gate: ToolGate,
757 decision: ToolGateVerdict,
758 #[serde(default, skip_serializing_if = "Option::is_none")]
759 risk: Option<String>,
760 reason: String,
761 },
762 AdvisoryNote {
763 thread_id: ThreadId,
764 session_id: SessionId,
765 turn_id: String,
766 note: String,
767 tool_call_count: u32,
768 },
769
770 // === Prefix cache ===
771 PrefixCacheChange {
772 thread_id: ThreadId,
773 session_id: SessionId,
774 description: String,
775 system_prompt_changed: bool,
776 tools_changed: bool,
777 stability_pct: u32,
778 changed: bool,
779 pinned_combined_hash: String,
780 pin_reason: String,
781 last_miss_reason: String,
782 context_updates: u64,
783 },
784 }
785
786 /// Envelope that carries an `EventMsg` over the wire / channel with a
787 /// monotonic seq so consumers can detect drops. Mirrors the existing
788 /// `RuntimeEventEnvelope` but typed to `EventMsg`.
789 #[derive(Debug, Clone, Serialize, Deserialize)]
790 pub struct EventEnvelope {
791 pub seq: u64,
792 pub thread_id: ThreadId,
793 pub session_id: SessionId,
794 pub turn_id: Option<String>,
795 pub event: EventMsg,
796 }
797
798 /// Every wire tag `EventMsg` can carry, in declaration order. Kept next to
799 /// the enum so a new variant is added here in the same edit; the test below
800 /// proves the list and `kind_str` agree.
801 pub const EVENT_KINDS: &[&str] = &[
802 "tool_projection_warning",
803 "message_started",
804 "response_delta",
805 "message_complete",
806 "thinking_started",
807 "thinking_complete",
808 "tool_call_started",
809 "tool_call_heartbeat",
810 "tool_call_complete",
811 "operation_activity_started",
812 "operation_activity_completed",
813 "turn_started",
814 "tool_request_snapshot",
815 "workspace_snapshot_taken",
816 "route_dispatched",
817 "turn_complete",
818 "turn_usage",
819 "routed_turn_usage",
820 "goal_updated",
821 "goal_continuation_waiting",
822 "goal_continuation_wait_ended",
823 "compaction_started",
824 "compaction_completed",
825 "compaction_cancelled",
826 "compaction_failed",
827 "purge_started",
828 "purge_completed",
829 "purge_failed",
830 "agent_spawned",
831 "agent_progress",
832 "agent_complete",
833 "sub_agent_follow_up",
834 "agent_list",
835 "sub_agent_mailbox",
836 "workflow_ui",
837 "error",
838 "status",
839 "mcp_session_boot",
840 "request_manifest_ready",
841 "pause_events",
842 "resume_events",
843 "approval_required",
844 "approval_withdrawn",
845 "user_input_required",
846 "session_updated",
847 "elevation_required",
848 "lsp_repair_update",
849 "tool_gate_decision",
850 "advisory_note",
851 "prefix_cache_change",
852 ];
853
854 impl EventMsg {
855 #[must_use]
856 pub fn kind_str(&self) -> &'static str {
857 match self {
858 Self::ToolProjectionWarning { .. } => "tool_projection_warning",
859 Self::SnapshotsDisabled { .. } => "snapshots_disabled",
860 Self::MessageStarted { .. } => "message_started",
861 Self::ResponseDelta { .. } => "response_delta",
862 Self::MessageComplete { .. } => "message_complete",
863 Self::ThinkingStarted { .. } => "thinking_started",
864 Self::ThinkingComplete { .. } => "thinking_complete",
865 Self::ToolCallStarted { .. } => "tool_call_started",
866 Self::ToolCallHeartbeat { .. } => "tool_call_heartbeat",
867 Self::ToolExecutionStarted { .. } => "tool_execution_started",
868 Self::ToolResultContent { .. } => "tool_result_content",
869 Self::ToolCallComplete { .. } => "tool_call_complete",
870 Self::OperationActivityStarted { .. } => "operation_activity_started",
871 Self::OperationActivityCompleted { .. } => "operation_activity_completed",
872 Self::TurnStarted { .. } => "turn_started",
873 Self::ToolRequestSnapshot { .. } => "tool_request_snapshot",
874 Self::WorkspaceSnapshotTaken { .. } => "workspace_snapshot_taken",
875 Self::RouteDispatched { .. } => "route_dispatched",
876 Self::TurnComplete { .. } => "turn_complete",
877 Self::TurnUsage { .. } => "turn_usage",
878 Self::RoutedTurnUsage { .. } => "routed_turn_usage",
879 Self::GoalUpdated { .. } => "goal_updated",
880 Self::GoalContinuationWaiting { .. } => "goal_continuation_waiting",
881 Self::GoalContinuationWaitEnded { .. } => "goal_continuation_wait_ended",
882 Self::CompactionStarted { .. } => "compaction_started",
883 Self::CompactionCompleted { .. } => "compaction_completed",
884 Self::CompactionCancelled { .. } => "compaction_cancelled",
885 Self::CompactionFailed { .. } => "compaction_failed",
886 Self::PurgeStarted { .. } => "purge_started",
887 Self::PurgeCompleted { .. } => "purge_completed",
888 Self::PurgeFailed { .. } => "purge_failed",
889 Self::AgentSpawned { .. } => "agent_spawned",
890 Self::AgentProgress { .. } => "agent_progress",
891 Self::AgentComplete { .. } => "agent_complete",
892 Self::SubAgentFollowUp { .. } => "sub_agent_follow_up",
893 Self::AgentList { .. } => "agent_list",
894 Self::SubAgentMailbox { .. } => "sub_agent_mailbox",
895 Self::WorkflowUi { .. } => "workflow_ui",
896 Self::Error { .. } => "error",
897 Self::Status { .. } => "status",
898 Self::McpSessionBoot { .. } => "mcp_session_boot",
899 Self::RequestManifestReady { .. } => "request_manifest_ready",
900 Self::PauseEvents { .. } => "pause_events",
901 Self::ResumeEvents { .. } => "resume_events",
902 Self::ApprovalRequired { .. } => "approval_required",
903 Self::ApprovalWithdrawn { .. } => "approval_withdrawn",
904 Self::UserInputRequired { .. } => "user_input_required",
905 Self::SessionUpdated { .. } => "session_updated",
906 Self::ElevationRequired { .. } => "elevation_required",
907 Self::LspRepairUpdate { .. } => "lsp_repair_update",
908 Self::ToolGateDecision { .. } => "tool_gate_decision",
909 Self::AdvisoryNote { .. } => "advisory_note",
910 Self::PrefixCacheChange { .. } => "prefix_cache_change",
911 }
912 }
913
914 #[must_use]
915 pub fn thread_id(&self) -> &ThreadId {
916 match self {
917 Self::ToolProjectionWarning { thread_id, .. }
918 | Self::SnapshotsDisabled { thread_id, .. }
919 | Self::MessageStarted { thread_id, .. }
920 | Self::ResponseDelta { thread_id, .. }
921 | Self::MessageComplete { thread_id, .. }
922 | Self::ThinkingStarted { thread_id, .. }
923 | Self::ThinkingComplete { thread_id, .. }
924 | Self::ToolCallStarted { thread_id, .. }
925 | Self::ToolCallHeartbeat { thread_id, .. }
926 | Self::ToolExecutionStarted { thread_id, .. }
927 | Self::ToolResultContent { thread_id, .. }
928 | Self::ToolCallComplete { thread_id, .. }
929 | Self::OperationActivityStarted { thread_id, .. }
930 | Self::OperationActivityCompleted { thread_id, .. }
931 | Self::TurnStarted { thread_id, .. }
932 | Self::ToolRequestSnapshot { thread_id, .. }
933 | Self::WorkspaceSnapshotTaken { thread_id, .. }
934 | Self::RouteDispatched { thread_id, .. }
935 | Self::TurnComplete { thread_id, .. }
936 | Self::TurnUsage { thread_id, .. }
937 | Self::RoutedTurnUsage { thread_id, .. }
938 | Self::GoalUpdated { thread_id, .. }
939 | Self::GoalContinuationWaiting { thread_id, .. }
940 | Self::GoalContinuationWaitEnded { thread_id, .. }
941 | Self::CompactionStarted { thread_id, .. }
942 | Self::CompactionCompleted { thread_id, .. }
943 | Self::CompactionCancelled { thread_id, .. }
944 | Self::CompactionFailed { thread_id, .. }
945 | Self::PurgeStarted { thread_id, .. }
946 | Self::PurgeCompleted { thread_id, .. }
947 | Self::PurgeFailed { thread_id, .. }
948 | Self::AgentSpawned { thread_id, .. }
949 | Self::AgentProgress { thread_id, .. }
950 | Self::AgentComplete { thread_id, .. }
951 | Self::SubAgentFollowUp { thread_id, .. }
952 | Self::AgentList { thread_id, .. }
953 | Self::SubAgentMailbox { thread_id, .. }
954 | Self::WorkflowUi { thread_id, .. }
955 | Self::Error { thread_id, .. }
956 | Self::Status { thread_id, .. }
957 | Self::McpSessionBoot { thread_id, .. }
958 | Self::RequestManifestReady { thread_id, .. }
959 | Self::PauseEvents { thread_id, .. }
960 | Self::ResumeEvents { thread_id, .. }
961 | Self::ApprovalRequired { thread_id, .. }
962 | Self::ApprovalWithdrawn { thread_id, .. }
963 | Self::UserInputRequired { thread_id, .. }
964 | Self::SessionUpdated { thread_id, .. }
965 | Self::ElevationRequired { thread_id, .. }
966 | Self::LspRepairUpdate { thread_id, .. }
967 | Self::ToolGateDecision { thread_id, .. }
968 | Self::AdvisoryNote { thread_id, .. }
969 | Self::PrefixCacheChange { thread_id, .. } => thread_id,
970 }
971 }
972
973 #[must_use]
974 pub fn session_id(&self) -> &SessionId {
975 match self {
976 Self::ToolProjectionWarning { session_id, .. }
977 | Self::SnapshotsDisabled { session_id, .. }
978 | Self::MessageStarted { session_id, .. }
979 | Self::ResponseDelta { session_id, .. }
980 | Self::MessageComplete { session_id, .. }
981 | Self::ThinkingStarted { session_id, .. }
982 | Self::ThinkingComplete { session_id, .. }
983 | Self::ToolCallStarted { session_id, .. }
984 | Self::ToolCallHeartbeat { session_id, .. }
985 | Self::ToolExecutionStarted { session_id, .. }
986 | Self::ToolResultContent { session_id, .. }
987 | Self::ToolCallComplete { session_id, .. }
988 | Self::OperationActivityStarted { session_id, .. }
989 | Self::OperationActivityCompleted { session_id, .. }
990 | Self::TurnStarted { session_id, .. }
991 | Self::ToolRequestSnapshot { session_id, .. }
992 | Self::WorkspaceSnapshotTaken { session_id, .. }
993 | Self::RouteDispatched { session_id, .. }
994 | Self::TurnComplete { session_id, .. }
995 | Self::TurnUsage { session_id, .. }
996 | Self::RoutedTurnUsage { session_id, .. }
997 | Self::GoalUpdated { session_id, .. }
998 | Self::GoalContinuationWaiting { session_id, .. }
999 | Self::GoalContinuationWaitEnded { session_id, .. }
1000 | Self::CompactionStarted { session_id, .. }
1001 | Self::CompactionCompleted { session_id, .. }
1002 | Self::CompactionCancelled { session_id, .. }
1003 | Self::CompactionFailed { session_id, .. }
1004 | Self::PurgeStarted { session_id, .. }
1005 | Self::PurgeCompleted { session_id, .. }
1006 | Self::PurgeFailed { session_id, .. }
1007 | Self::AgentSpawned { session_id, .. }
1008 | Self::AgentProgress { session_id, .. }
1009 | Self::AgentComplete { session_id, .. }
1010 | Self::SubAgentFollowUp { session_id, .. }
1011 | Self::AgentList { session_id, .. }
1012 | Self::SubAgentMailbox { session_id, .. }
1013 | Self::WorkflowUi { session_id, .. }
1014 | Self::Error { session_id, .. }
1015 | Self::Status { session_id, .. }
1016 | Self::McpSessionBoot { session_id, .. }
1017 | Self::RequestManifestReady { session_id, .. }
1018 | Self::PauseEvents { session_id, .. }
1019 | Self::ResumeEvents { session_id, .. }
1020 | Self::ApprovalRequired { session_id, .. }
1021 | Self::ApprovalWithdrawn { session_id, .. }
1022 | Self::UserInputRequired { session_id, .. }
1023 | Self::SessionUpdated { session_id, .. }
1024 | Self::ElevationRequired { session_id, .. }
1025 | Self::LspRepairUpdate { session_id, .. }
1026 | Self::ToolGateDecision { session_id, .. }
1027 | Self::AdvisoryNote { session_id, .. }
1028 | Self::PrefixCacheChange { session_id, .. } => session_id,
1029 }
1030 }
1031 }
1032
1033 #[cfg(test)]
1034 mod tests {
1035 use super::*;
1036 use serde_json::json;
1037
1038 fn ids() -> (ThreadId, SessionId) {
1039 (ThreadId::new(), SessionId::new())
1040 }
1041
1042 /// One instance of every variant, in declaration order. A new variant
1043 /// must be added here too, or `every_variant_is_listed_once` fails.
1044 fn every_variant() -> Vec<EventMsg> {
1045 let (t, s) = ids();
1046 let usage = TokenUsage {
1047 input_tokens: 1,
1048 output_tokens: 2,
1049 ..TokenUsage::default()
1050 };
1051 let route = TurnRoute {
1052 provider: "deepseek".into(),
1053 provider_identity: "deepseek".into(),
1054 model: "deepseek-chat".into(),
1055 auto_model: false,
1056 receipt: Some(TurnRouteReceipt {
1057 provider: "deepseek".into(),
1058 provider_identity: "deepseek".into(),
1059 wire_model: "deepseek-chat".into(),
1060 endpoint_identity: "api.deepseek.com".into(),
1061 credential_generation_present: true,
1062 }),
1063 billing: Some(RouteBillingEnvelope {
1064 openrouter_vendor: None,
1065 billing_surface: None,
1066 endpoint_fingerprint: Some("fp".into()),
1067 provider_live_pricing: None,
1068 billing_mode: "metered".into(),
1069 dispatched_at: DateTime::<Utc>::from_timestamp(0, 0).unwrap(),
1070 }),
1071 base_url: "https://api.deepseek.com".into(),
1072 billing_product: RouteProduct::Subscription {
1073 label: "pro".into(),
1074 },
1075 };
1076 vec![
1077 EventMsg::ToolProjectionWarning {
1078 thread_id: t.clone(),
1079 session_id: s.clone(),
1080 provider: "openai".into(),
1081 omitted_tool_names: vec!["a".into()],
1082 omitted_tool_count: 3,
1083 },
1084 EventMsg::MessageStarted {
1085 thread_id: t.clone(),
1086 session_id: s.clone(),
1087 index: 0,
1088 },
1089 EventMsg::ResponseDelta {
1090 thread_id: t.clone(),
1091 session_id: s.clone(),
1092 index: 0,
1093 delta: "hi".into(),
1094 channel: ResponseChannel::Reasoning,
1095 },
1096 EventMsg::MessageComplete {
1097 thread_id: t.clone(),
1098 session_id: s.clone(),
1099 index: 0,
1100 },
1101 EventMsg::ThinkingStarted {
1102 thread_id: t.clone(),
1103 session_id: s.clone(),
1104 index: 1,
1105 },
1106 EventMsg::ThinkingComplete {
1107 thread_id: t.clone(),
1108 session_id: s.clone(),
1109 index: 1,
1110 },
1111 EventMsg::ToolCallStarted {
1112 thread_id: t.clone(),
1113 session_id: s.clone(),
1114 tool_call_id: "c1".into(),
1115 tool_name: "read_file".into(),
1116 input: json!({"path": "x"}),
1117 },
1118 EventMsg::ToolCallHeartbeat {
1119 thread_id: t.clone(),
1120 session_id: s.clone(),
1121 },
1122 EventMsg::ToolCallComplete {
1123 thread_id: t.clone(),
1124 session_id: s.clone(),
1125 tool_call_id: "c1".into(),
1126 tool_name: "read_file".into(),
1127 result: ToolCallOutcome::Err {
1128 error: ToolCallError::Timeout { seconds: 3 },
1129 },
1130 },
1131 EventMsg::OperationActivityStarted {
1132 thread_id: t.clone(),
1133 session_id: s.clone(),
1134 span_id: "span-1".into(),
1135 activity_kind: crate::engine_owner::OwnerActivityKind::Reading,
1136 },
1137 EventMsg::OperationActivityCompleted {
1138 thread_id: t.clone(),
1139 session_id: s.clone(),
1140 span_id: "span-1".into(),
1141 activity_kind: crate::engine_owner::OwnerActivityKind::Reading,
1142 outcome: crate::engine_owner::OwnerOperationOutcome::Succeeded,
1143 },
1144 EventMsg::TurnStarted {
1145 thread_id: t.clone(),
1146 session_id: s.clone(),
1147 turn_id: "turn-1".into(),
1148 created_at: DateTime::<Utc>::from_timestamp(1, 0).unwrap(),
1149 route: Some(route.clone()),
1150 submission_id: None,
1151 },
1152 EventMsg::ToolRequestSnapshot {
1153 thread_id: t.clone(),
1154 session_id: s.clone(),
1155 snapshot: json!({"tool_count": 2}),
1156 },
1157 EventMsg::WorkspaceSnapshotTaken {
1158 thread_id: t.clone(),
1159 session_id: s.clone(),
1160 snapshot: json!({"kind": "pre_turn", "tree_id": "t"}),
1161 },
1162 EventMsg::RouteDispatched {
1163 thread_id: t.clone(),
1164 session_id: s.clone(),
1165 turn_id: "turn-1".into(),
1166 route,
1167 },
1168 EventMsg::TurnComplete {
1169 thread_id: t.clone(),
1170 session_id: s.clone(),
1171 turn_id: None,
1172 status: TurnOutcomeStatus::Failed,
1173 error: Some("boom".into()),
1174 usage: usage.clone(),
1175 parent_route_usage: Some(usage.clone()),
1176 routed_usage_dropped_records: 0,
1177 tool_catalog: Some(vec![json!({"name": "read_file"})]),
1178 base_url: None,
1179 },
1180 EventMsg::TurnUsage {
1181 max_output_tokens: None,
1182 thread_id: t.clone(),
1183 session_id: s.clone(),
1184 usage: usage.clone(),
1185 duration_ms: 10,
1186 first_token_ms: Some(2),
1187 request_ms: None,
1188 },
1189 EventMsg::RoutedTurnUsage {
1190 thread_id: t.clone(),
1191 session_id: s.clone(),
1192 usage,
1193 duration_ms: 10,
1194 first_token_ms: Some(2),
1195 request_ms: None,
1196 },
1197 EventMsg::GoalUpdated {
1198 thread_id: t.clone(),
1199 session_id: s.clone(),
1200 snapshot: json!({"status": "active"}),
1201 },
1202 EventMsg::GoalContinuationWaiting {
1203 thread_id: t.clone(),
1204 session_id: s.clone(),
1205 delay_seconds: 5,
1206 },
1207 EventMsg::GoalContinuationWaitEnded {
1208 thread_id: t.clone(),
1209 session_id: s.clone(),
1210 interrupted: true,
1211 },
1212 EventMsg::CompactionStarted {
1213 thread_id: t.clone(),
1214 session_id: s.clone(),
1215 id: "cmp-1".into(),
1216 auto: true,
1217 message: "m".into(),
1218 },
1219 EventMsg::CompactionCompleted {
1220 thread_id: t.clone(),
1221 session_id: s.clone(),
1222 id: "cmp-1".into(),
1223 auto: true,
1224 message: "m".into(),
1225 messages_before: Some(10),
1226 messages_after: Some(2),
1227 summary_prompt: None,
1228 post_input_tokens: Some(100),
1229 },
1230 EventMsg::CompactionCancelled {
1231 thread_id: t.clone(),
1232 session_id: s.clone(),
1233 id: "cmp-1".into(),
1234 auto: false,
1235 message: "m".into(),
1236 },
1237 EventMsg::CompactionFailed {
1238 thread_id: t.clone(),
1239 session_id: s.clone(),
1240 id: "cmp-1".into(),
1241 auto: false,
1242 message: "m".into(),
1243 },
1244 EventMsg::PurgeStarted {
1245 thread_id: t.clone(),
1246 session_id: s.clone(),
1247 message: "m".into(),
1248 },
1249 EventMsg::PurgeCompleted {
1250 thread_id: t.clone(),
1251 session_id: s.clone(),
1252 messages_before: 4,
1253 messages_after: 2,
1254 removed_count: 2,
1255 replaced_count: 0,
1256 message: "m".into(),
1257 },
1258 EventMsg::PurgeFailed {
1259 thread_id: t.clone(),
1260 session_id: s.clone(),
1261 message: "m".into(),
1262 },
1263 EventMsg::AgentSpawned {
1264 thread_id: t.clone(),
1265 session_id: s.clone(),
1266 owner_session_id: "owner".into(),
1267 id: "a1".into(),
1268 prompt: "p".into(),
1269 worker_status: Some("starting".into()),
1270 parent_run_id: None,
1271 spawn_depth: 1,
1272 model: "m".into(),
1273 route_source: Some("task.model".into()),
1274 display_name: Some("audit docs".into()),
1275 },
1276 EventMsg::AgentProgress {
1277 thread_id: t.clone(),
1278 session_id: s.clone(),
1279 owner_session_id: "owner".into(),
1280 id: "a1".into(),
1281 status: "running".into(),
1282 activity: AgentProgressActivity {
1283 worker_status: "running_tool".into(),
1284 step: Some(2),
1285 tool_name: Some("bash".into()),
1286 },
1287 parent_run_id: None,
1288 spawn_depth: 1,
1289 },
1290 EventMsg::AgentComplete {
1291 thread_id: t.clone(),
1292 session_id: s.clone(),
1293 owner_session_id: "owner".into(),
1294 id: "a1".into(),
1295 result: "done".into(),
1296 worker_status: Some("completed".into()),
1297 parent_run_id: None,
1298 spawn_depth: Some(1),
1299 continuable: Some(false),
1300 display_name: Some("audit docs".into()),
1301 },
1302 EventMsg::SubAgentFollowUp {
1303 thread_id: t.clone(),
1304 session_id: s.clone(),
1305 owner_session_id: "owner".into(),
1306 agent_id: "a1".into(),
1307 outcome: SubAgentFollowUpOutcome::Ok {
1308 agent_id: "a1".into(),
1309 target_agent_id: "a2".into(),
1310 delivered: false,
1311 resumed: true,
1312 note: "resumed".into(),
1313 },
1314 },
1315 EventMsg::AgentList {
1316 thread_id: t.clone(),
1317 session_id: s.clone(),
1318 owner_session_id: "owner".into(),
1319 agents: vec![json!({"id": "a1"})],
1320 coordination: json!({}),
1321 queued_follow_ups: BTreeMap::from([("a1".to_string(), 1)]),
1322 roster: vec![AgentRosterRow {
1323 worker_id: "w1".into(),
1324 display_name: "scout".into(),
1325 model: "m".into(),
1326 state: "running".into(),
1327 status: "running".into(),
1328 activity: None,
1329 millis: Some(5),
1330 input_tokens: None,
1331 output_tokens: None,
1332 cost_microusd: None,
1333 steps_taken: 1,
1334 parent_run_id: None,
1335 run_id: "r1".into(),
1336 }],
1337 },
1338 EventMsg::SubAgentMailbox {
1339 thread_id: t.clone(),
1340 session_id: s.clone(),
1341 owner_session_id: "owner".into(),
1342 turn_id: "turn-1".into(),
1343 seq: 7,
1344 message: json!({"kind": "spawned"}),
1345 },
1346 EventMsg::WorkflowUi {
1347 thread_id: t.clone(),
1348 session_id: s.clone(),
1349 owner_session_id: "owner".into(),
1350 run_id: "run-1".into(),
1351 ui_event: json!({"type": "task_started"}),
1352 },
1353 EventMsg::Error {
1354 thread_id: t.clone(),
1355 session_id: s.clone(),
1356 category: "network".into(),
1357 severity: "error".into(),
1358 recoverable: true,
1359 code: "E1".into(),
1360 message: "m".into(),
1361 },
1362 EventMsg::Status {
1363 thread_id: t.clone(),
1364 session_id: s.clone(),
1365 message: "m".into(),
1366 },
1367 EventMsg::McpSessionBoot {
1368 thread_id: t.clone(),
1369 session_id: s.clone(),
1370 generation: 1,
1371 snapshot: McpManagerSnapshot {
1372 config_path: PathBuf::from("/tmp/mcp.json"),
1373 config_exists: true,
1374 reload_required: false,
1375 servers: vec![McpServerSnapshot {
1376 name: "fs".into(),
1377 enabled: true,
1378 required: false,
1379 transport: "stdio".into(),
1380 command_or_url: "npx".into(),
1381 connect_timeout: 1,
1382 execute_timeout: 2,
1383 read_timeout: 3,
1384 connected: true,
1385 error: None,
1386 capability_metadata: "advertised".into(),
1387 tools: vec![McpDiscoveredItem {
1388 name: "read".into(),
1389 model_name: "fs__read".into(),
1390 description: None,
1391 }],
1392 resources: vec![],
1393 prompts: vec![],
1394 }],
1395 },
1396 connecting: vec!["slow".into()],
1397 finished: false,
1398 },
1399 EventMsg::RequestManifestReady {
1400 thread_id: t.clone(),
1401 session_id: s.clone(),
1402 rendered: "manifest".into(),
1403 },
1404 EventMsg::PauseEvents {
1405 thread_id: t.clone(),
1406 session_id: s.clone(),
1407 },
1408 EventMsg::ResumeEvents {
1409 thread_id: t.clone(),
1410 session_id: s.clone(),
1411 },
1412 EventMsg::ApprovalRequired {
1413 thread_id: t.clone(),
1414 session_id: s.clone(),
1415 id: "c1".into(),
1416 tool_name: "bash".into(),
1417 description: "rm".into(),
1418 input: json!({"command": "rm"}),
1419 approval_key: "k".into(),
1420 approval_grouping_key: "g".into(),
1421 intent_summary: None,
1422 approval_force_prompt: true,
1423 },
1424 EventMsg::ApprovalWithdrawn {
1425 thread_id: t.clone(),
1426 session_id: s.clone(),
1427 id: "c1".into(),
1428 },
1429 EventMsg::UserInputRequired {
1430 thread_id: t.clone(),
1431 session_id: s.clone(),
1432 id: "c2".into(),
1433 request: UserInputRequest {
1434 questions: vec![UserInputQuestionEvent {
1435 header: "h".into(),
1436 id: "q1".into(),
1437 question: "?".into(),
1438 options: vec![],
1439 allow_free_text: true,
1440 multi_select: false,
1441 }],
1442 },
1443 },
1444 EventMsg::SessionUpdated {
1445 thread_id: t.clone(),
1446 session_id: s.clone(),
1447 engine_session_id: "sess".into(),
1448 messages: vec![json!({"role": "user", "content": []})],
1449 system_prompt: Some(json!("sys")),
1450 model: "m".into(),
1451 workspace: PathBuf::from("/ws"),
1452 },
1453 EventMsg::ElevationRequired {
1454 thread_id: t.clone(),
1455 session_id: s.clone(),
1456 tool_id: "c3".into(),
1457 tool_name: "bash".into(),
1458 command: Some("curl".into()),
1459 denial_reason: "net".into(),
1460 blocked_network: true,
1461 blocked_write: false,
1462 },
1463 EventMsg::LspRepairUpdate {
1464 thread_id: t.clone(),
1465 session_id: s.clone(),
1466 diagnostics_found: 1,
1467 files: 1,
1468 injected: true,
1469 },
1470 EventMsg::ToolGateDecision {
1471 thread_id: t.clone(),
1472 session_id: s.clone(),
1473 agent_id: None,
1474 tool_id: "c4".into(),
1475 tool_name: "bash".into(),
1476 gate: ToolGate::AutoReviewGuardian,
1477 decision: ToolGateVerdict::Denied,
1478 risk: Some("high".into()),
1479 reason: "no".into(),
1480 },
1481 EventMsg::AdvisoryNote {
1482 thread_id: t.clone(),
1483 session_id: s.clone(),
1484 turn_id: "turn-1".into(),
1485 note: "n".into(),
1486 tool_call_count: 2,
1487 },
1488 EventMsg::PrefixCacheChange {
1489 thread_id: t,
1490 session_id: s,
1491 description: "d".into(),
1492 system_prompt_changed: false,
1493 tools_changed: true,
1494 stability_pct: 90,
1495 changed: true,
1496 pinned_combined_hash: "h".into(),
1497 pin_reason: "initial".into(),
1498 last_miss_reason: String::new(),
1499 context_updates: 0,
1500 },
1501 ]
1502 }
1503
1504 #[test]
1505 fn every_variant_is_listed_once() {
1506 let kinds: Vec<&str> = every_variant().iter().map(EventMsg::kind_str).collect();
1507 assert_eq!(
1508 kinds, EVENT_KINDS,
1509 "EVENT_KINDS must list every variant in order"
1510 );
1511 }
1512
1513 #[test]
1514 fn every_variant_round_trips_and_tags_by_event() {
1515 for msg in every_variant() {
1516 let value = serde_json::to_value(&msg).unwrap();
1517 assert_eq!(value["event"], msg.kind_str(), "{msg:?}");
1518 assert_eq!(value["thread_id"], msg.thread_id().to_string(), "{msg:?}");
1519 assert_eq!(value["session_id"], msg.session_id().to_string(), "{msg:?}");
1520 let back: EventMsg = serde_json::from_value(value).unwrap();
1521 assert_eq!(back, msg);
1522 }
1523 }
1524
1525 #[test]
1526 fn agent_events_from_producers_without_a_display_name_still_load() {
1527 // #6565 added `display_name`; payloads written before it omit it.
1528 let mut every = every_variant();
1529 every.retain(|msg| {
1530 matches!(
1531 msg,
1532 EventMsg::AgentSpawned { .. } | EventMsg::AgentComplete { .. }
1533 )
1534 });
1535 assert_eq!(every.len(), 2);
1536 for msg in every {
1537 let mut value = serde_json::to_value(&msg).unwrap();
1538 assert_eq!(value["display_name"], "audit docs");
1539 value.as_object_mut().unwrap().remove("display_name");
1540 let back: EventMsg = serde_json::from_value(value.clone()).unwrap();
1541 match &back {
1542 EventMsg::AgentSpawned { display_name, .. }
1543 | EventMsg::AgentComplete { display_name, .. } => assert_eq!(*display_name, None),
1544 other => panic!("unexpected {other:?}"),
1545 }
1546 assert!(
1547 serde_json::to_value(&back)
1548 .unwrap()
1549 .get("display_name")
1550 .is_none()
1551 );
1552 }
1553 }
1554
1555 #[test]
1556 fn event_msg_roundtrip() {
1557 let msg = EventMsg::TurnComplete {
1558 thread_id: ThreadId::new(),
1559 session_id: SessionId::new(),
1560 turn_id: Some("turn-1".into()),
1561 status: TurnOutcomeStatus::Completed,
1562 error: None,
1563 usage: TokenUsage::default(),
1564 parent_route_usage: None,
1565 routed_usage_dropped_records: 0,
1566 tool_catalog: None,
1567 base_url: None,
1568 };
1569 let json = serde_json::to_string(&msg).unwrap();
1570 let back: EventMsg = serde_json::from_str(&json).unwrap();
1571 assert_eq!(back.kind_str(), "turn_complete");
1572 assert!(json.contains(r#""status":"completed""#));
1573 }
1574
1575 #[test]
1576 fn text_channel_is_elided_on_the_wire() {
1577 let msg = EventMsg::ResponseDelta {
1578 thread_id: ThreadId::new(),
1579 session_id: SessionId::new(),
1580 index: 0,
1581 delta: "x".into(),
1582 channel: ResponseChannel::Text,
1583 };
1584 let json = serde_json::to_string(&msg).unwrap();
1585 assert!(!json.contains("channel"), "{json}");
1586 let back: EventMsg = serde_json::from_str(&json).unwrap();
1587 assert_eq!(back, msg);
1588 }
1589
1590 /// A host-stamped `submission_id` crosses the wire verbatim, a `None`
1591 /// token stays absent, and a payload from a producer that predates the
1592 /// field still deserializes (`serde(default)`).
1593 #[test]
1594 fn turn_started_submission_id_is_additive_and_default_absent() {
1595 let msg = EventMsg::TurnStarted {
1596 thread_id: ThreadId::new(),
1597 session_id: SessionId::new(),
1598 turn_id: "turn-1".into(),
1599 created_at: DateTime::<Utc>::from_timestamp(1, 0).unwrap(),
1600 route: None,
1601 submission_id: Some("sub-host-1".into()),
1602 };
1603 let json = serde_json::to_string(&msg).unwrap();
1604 assert!(
1605 json.contains(r#""submission_id":"sub-host-1""#),
1606 "a host-stamped token must be serialized: {json}"
1607 );
1608 let back: EventMsg = serde_json::from_str(&json).unwrap();
1609 assert_eq!(back, msg);
1610
1611 let none = EventMsg::TurnStarted {
1612 thread_id: ThreadId::new(),
1613 session_id: SessionId::new(),
1614 turn_id: "turn-2".into(),
1615 created_at: DateTime::<Utc>::from_timestamp(2, 0).unwrap(),
1616 route: None,
1617 submission_id: None,
1618 };
1619 let value = serde_json::to_value(&none).unwrap();
1620 assert!(
1621 value.get("submission_id").is_none(),
1622 "a self-started turn's None token must stay absent on the wire: {value}"
1623 );
1624 // An older producer that predates the field: the absent key defaults.
1625 let back: EventMsg = serde_json::from_value(value).unwrap();
1626 assert_eq!(back, none);
1627 }
1628 }
1629
1629 lines RUST