1use serde::{Deserialize, Serialize};
16use serde_json::Value;
17
18#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
20pub struct MuseRecord {
21 pub schema_version: u32,
23 pub id: String,
25 pub stream: StreamRef,
27 pub sequence: u64,
29 pub recorded_at: u64,
31 pub record_type: RecordType,
32 pub durability: Durability,
33 pub causation_id: String,
35 pub payload_type: String,
37 pub payload_schema_version: u32,
39 pub payload: Value,
41}
42
43impl MuseRecord {
44 pub fn typed_payload(&self) -> serde_json::Result<MusePayload> {
50 MusePayload::from_parts(&self.payload_type, self.payload.clone())
51 }
52}
53
54#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
56pub struct StreamRef {
57 pub kind: StreamKind,
58 pub id: String,
59}
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
63#[serde(rename_all = "snake_case")]
64pub enum StreamKind {
65 Session,
66 Run,
67 Task,
68}
69
70#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
72#[serde(rename_all = "snake_case")]
73pub enum RecordType {
74 Reconciliation,
76 Event,
78 Status,
80}
81
82#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
84#[serde(rename_all = "snake_case")]
85pub enum Durability {
86 Durable,
87 Ephemeral,
88}
89
90#[derive(Debug, Clone, PartialEq)]
92pub enum MusePayload {
93 CommandAccepted(CommandAccepted),
95 SessionRunLinked(SessionRunLinked),
97 TurnInputUser(TurnInputUser),
99 RunStarted(RunStarted),
101 ModelConfigured(ModelConfigured),
103 RunOutputDelta(RunOutputDelta),
105 ToolResult(ToolResult),
107 RunTerminal(RunTerminal),
109 TaskStreamLinked(TaskStreamLinked),
111 TaskLifecycle(TaskLifecycle),
113 Unknown {
115 payload_type: String,
116 payload: Value,
117 },
118}
119
120impl MusePayload {
121 pub fn from_parts(payload_type: &str, payload: Value) -> serde_json::Result<Self> {
122 Ok(match payload_type {
123 "runtime.command.accepted" => {
124 MusePayload::CommandAccepted(serde_json::from_value(payload)?)
125 }
126 "session.run.linked" => MusePayload::SessionRunLinked(serde_json::from_value(payload)?),
127 "turn.input.user" => MusePayload::TurnInputUser(serde_json::from_value(payload)?),
128 "run.lifecycle.started" => MusePayload::RunStarted(serde_json::from_value(payload)?),
129 "run.model.configured" => {
130 MusePayload::ModelConfigured(serde_json::from_value(payload)?)
131 }
132 "tool.result" => MusePayload::ToolResult(serde_json::from_value(payload)?),
133 "run.output.delta" => MusePayload::RunOutputDelta(serde_json::from_value(payload)?),
134 t if t.starts_with("run.terminal.") => {
135 MusePayload::RunTerminal(serde_json::from_value(payload)?)
136 }
137 "task.stream.linked" => MusePayload::TaskStreamLinked(serde_json::from_value(payload)?),
138 t if t.starts_with("task.lifecycle.") => {
139 MusePayload::TaskLifecycle(serde_json::from_value(payload)?)
140 }
141 other => MusePayload::Unknown {
142 payload_type: other.to_string(),
143 payload,
144 },
145 })
146 }
147}
148
149#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
152pub struct CommandAccepted {
153 pub kind: String,
154 pub command_id: String,
155 pub command_kind: String,
156 pub client_id: Option<String>,
157 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
158 pub extra: serde_json::Map<String, Value>,
159}
160
161#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
163pub struct SessionRunLinked {
164 pub kind: String,
165 pub command_id: String,
166 pub run_stream: StreamRef,
167 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
168 pub extra: serde_json::Map<String, Value>,
169}
170
171#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
173pub struct TurnInputUser {
174 pub kind: String,
175 pub command_id: String,
176 pub prompt: String,
177 pub run_stream: StreamRef,
178 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
179 pub extra: serde_json::Map<String, Value>,
180}
181
182#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
184pub struct RunStarted {
185 pub kind: String,
186 pub command_id: String,
187 pub prompt: String,
188 pub run_stream: StreamRef,
189 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
190 pub extra: serde_json::Map<String, Value>,
191}
192
193#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
195pub struct RunOutputDelta {
196 pub kind: String,
197 pub command_id: String,
198 pub run_stream: StreamRef,
199 pub text: String,
200 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
201 pub extra: serde_json::Map<String, Value>,
202}
203
204#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
207pub struct ModelConfigured {
208 pub kind: String,
209 pub command_id: String,
210 pub run_stream: StreamRef,
211 pub model_id: String,
212 pub display_label: String,
213 pub profile_id: String,
214 pub provider_id: String,
215 pub source: String,
217 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
218 pub extra: serde_json::Map<String, Value>,
219}
220
221#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
241pub struct ToolResult {
242 pub kind: String,
243 pub command_id: String,
244 pub run_stream: StreamRef,
245 pub call_id: String,
247 pub text: String,
251 #[serde(default, skip_serializing_if = "Option::is_none")]
254 pub correlation_facts: Option<Value>,
255 #[serde(default, skip_serializing_if = "Option::is_none")]
257 pub edit_facts: Option<Value>,
258 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
259 pub extra: serde_json::Map<String, Value>,
260}
261
262impl ToolResult {
263 pub fn is_command_tool(&self) -> bool {
266 matches!(
267 self.correlation_facts
268 .as_ref()
269 .and_then(|v| v.get("tool_name"))
270 .and_then(|v| v.as_str()),
271 Some("bash" | "command")
272 )
273 }
274
275 pub fn command_result(&self) -> Option<CommandResult> {
281 serde_json::from_str(&self.text).ok()
282 }
283
284 pub fn try_command_result(&self) -> Result<CommandResult, serde_json::Error> {
286 serde_json::from_str(&self.text)
287 }
288}
289
290#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
312pub struct CommandResult {
313 pub chunk_id: String,
314 pub command: String,
315 pub description: String,
316 pub exit_code: i32,
317 pub terminal_status: String,
318 pub output: String,
319 pub original_output_bytes: u64,
320 pub original_output_tokens: u64,
321 pub truncated: bool,
322 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
323 pub extra: serde_json::Map<String, Value>,
324}
325
326#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
330pub struct RunTerminal {
331 pub kind: String,
332 pub command_id: String,
333 pub run_stream: StreamRef,
334 pub terminal: String,
335 pub reason: Option<String>,
336 #[serde(default, skip_serializing_if = "Option::is_none")]
337 pub text: Option<String>,
338 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
339 pub extra: serde_json::Map<String, Value>,
340}
341
342#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
344pub struct TaskStreamLinked {
345 pub kind: String,
346 pub command_id: String,
347 pub run_stream: StreamRef,
348 pub task_id: String,
349 pub task_stream: StreamRef,
350 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
351 pub extra: serde_json::Map<String, Value>,
352}
353
354#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
356pub struct TaskLifecycle {
357 pub kind: String,
358 pub command_id: String,
359 pub run_stream: StreamRef,
360 pub task_id: String,
361 pub task_stream: StreamRef,
362 pub event: TaskLifecycleEvent,
363 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
364 pub extra: serde_json::Map<String, Value>,
365}
366
367#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
372#[serde(tag = "kind", rename_all = "snake_case")]
373pub enum TaskLifecycleEvent {
374 Proposed {
375 task_id: String,
376 task_kind: String,
379 },
380 Accepted {
381 task_id: String,
382 },
383 Started {
384 task_id: String,
385 #[serde(default, skip_serializing_if = "Option::is_none")]
387 span_id: Option<String>,
388 },
389 Scheduled {
390 task_id: String,
391 idempotency_key: String,
392 },
393 SideEffectIntent {
394 task_id: String,
395 idempotency_key: String,
396 operation: String,
397 policy_decision: String,
398 parent_task_id: Option<String>,
399 cancellation_handle: Option<Value>,
400 },
401 Status {
404 task_id: String,
405 message: String,
406 details: Value,
407 },
408 Output {
410 task_id: String,
411 chunk: String,
412 },
413 Completed {
414 task_id: String,
415 },
416 Cancelled {
417 task_id: String,
418 reason: String,
419 },
420 Rejected {
421 task_id: String,
422 reason: String,
423 },
424 Failed {
425 task_id: String,
426 reason: String,
427 },
428 #[serde(untagged)]
430 Unknown(Value),
431}
432
433#[cfg(test)]
434mod tests {
435 use super::*;
436 use serde_json::json;
437
438 #[test]
439 fn unknown_payload_type_is_preserved_not_error() {
440 let p = MusePayload::from_parts("subagent.lifecycle.spawned", json!({"x": 1})).unwrap();
441 match p {
442 MusePayload::Unknown {
443 payload_type,
444 payload,
445 } => {
446 assert_eq!(payload_type, "subagent.lifecycle.spawned");
447 assert_eq!(payload, json!({"x": 1}));
448 }
449 other => panic!("expected Unknown, got {other:?}"),
450 }
451 }
452
453 #[test]
454 fn task_lifecycle_failed_carries_reason() {
455 let e: TaskLifecycleEvent = serde_json::from_value(json!({
456 "kind": "failed",
457 "task_id": "t1",
458 "reason": "provider does not support base instructions"
459 }))
460 .unwrap();
461 assert!(matches!(e, TaskLifecycleEvent::Failed { ref reason, .. }
462 if reason.contains("base instructions")));
463 }
464
465 #[test]
466 fn command_result_parses_real_wire_shape_and_preserves_extensions() {
467 let result: ToolResult = serde_json::from_value(json!({
468 "kind": "tool_result",
469 "command_id": "cmd-1",
470 "run_stream": { "id": "run-1", "kind": "run" },
471 "call_id": "call-1",
472 "correlation_facts": { "outcome": "success", "tool_name": "bash" },
473 "text": r#"{"chunk_id":"exec-12-1","command":"printf ok","description":"Print a value","exit_code":0,"terminal_status":"completed","output":"ok","original_output_bytes":2,"original_output_tokens":1,"truncated":false,"provider_extension":true}"#
474 }))
475 .unwrap();
476
477 assert!(result.is_command_tool());
478 let command = result.command_result().expect("typed command result");
479 assert_eq!(command.command, "printf ok");
480 assert_eq!(command.output, "ok");
481 assert_eq!(command.exit_code, 0);
482 assert_eq!(command.extra["provider_extension"], true);
483 assert_eq!(
484 serde_json::to_value(command).unwrap()["provider_extension"],
485 true
486 );
487 }
488
489 #[test]
490 fn command_result_rejects_prose_and_recognizes_command_alias() {
491 let mut result: ToolResult = serde_json::from_value(json!({
492 "kind": "tool_result",
493 "command_id": "cmd-1",
494 "run_stream": { "id": "run-1", "kind": "run" },
495 "call_id": "call-1",
496 "correlation_facts": { "outcome": "failure", "tool_name": "command" },
497 "text": "tool failed before the command started"
498 }))
499 .unwrap();
500
501 assert!(result.is_command_tool());
502 assert!(result.command_result().is_none());
503 assert!(result.try_command_result().is_err());
504
505 result.correlation_facts = Some(json!({ "tool_name": "write_file" }));
506 assert!(!result.is_command_tool());
507 }
508}