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 #[serde(default)]
217 pub profile_id: Option<String>,
218 pub provider_id: String,
219 pub source: String,
221 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
222 pub extra: serde_json::Map<String, Value>,
223}
224
225#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
245pub struct ToolResult {
246 pub kind: String,
247 pub command_id: String,
248 pub run_stream: StreamRef,
249 pub call_id: String,
251 pub text: String,
255 #[serde(default, skip_serializing_if = "Option::is_none")]
259 pub correlation_facts: Option<ToolCorrelationFacts>,
260 #[serde(default, skip_serializing_if = "Option::is_none")]
262 pub edit_facts: Option<Value>,
263 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
264 pub extra: serde_json::Map<String, Value>,
265}
266
267#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
273pub struct ToolCorrelationFacts {
274 #[serde(default, skip_serializing_if = "Option::is_none")]
277 pub tool_name: Option<String>,
278 #[serde(default, skip_serializing_if = "Option::is_none")]
280 pub outcome: Option<String>,
281 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
282 pub extra: serde_json::Map<String, Value>,
283}
284
285impl ToolResult {
286 pub fn outcome(&self) -> Option<&str> {
288 self.correlation_facts
289 .as_ref()
290 .and_then(|f| f.outcome.as_deref())
291 }
292
293 pub fn tool_name(&self) -> Option<&str> {
295 self.correlation_facts
296 .as_ref()
297 .and_then(|f| f.tool_name.as_deref())
298 }
299}
300
301impl ToolResult {
302 pub fn is_command_tool(&self) -> bool {
305 matches!(self.tool_name(), Some("bash" | "command"))
306 }
307
308 pub fn command_result(&self) -> Option<CommandResult> {
314 serde_json::from_str(&self.text).ok()
315 }
316
317 pub fn try_command_result(&self) -> Result<CommandResult, serde_json::Error> {
319 serde_json::from_str(&self.text)
320 }
321}
322
323#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
345pub struct CommandResult {
346 pub chunk_id: String,
347 pub command: String,
348 pub description: String,
349 pub exit_code: i32,
350 pub terminal_status: String,
351 pub output: String,
352 pub original_output_bytes: u64,
353 pub original_output_tokens: u64,
354 pub truncated: bool,
355 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
356 pub extra: serde_json::Map<String, Value>,
357}
358
359#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
363pub struct RunTerminal {
364 pub kind: String,
365 pub command_id: String,
366 pub run_stream: StreamRef,
367 pub terminal: String,
368 pub reason: Option<String>,
369 #[serde(default, skip_serializing_if = "Option::is_none")]
370 pub text: Option<String>,
371 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
372 pub extra: serde_json::Map<String, Value>,
373}
374
375#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
377pub struct TaskStreamLinked {
378 pub kind: String,
379 pub command_id: String,
380 pub run_stream: StreamRef,
381 pub task_id: String,
382 pub task_stream: StreamRef,
383 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
384 pub extra: serde_json::Map<String, Value>,
385}
386
387#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
389pub struct TaskLifecycle {
390 pub kind: String,
391 pub command_id: String,
392 pub run_stream: StreamRef,
393 pub task_id: String,
394 pub task_stream: StreamRef,
395 pub event: TaskLifecycleEvent,
396 #[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
397 pub extra: serde_json::Map<String, Value>,
398}
399
400#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
405#[serde(tag = "kind", rename_all = "snake_case")]
406pub enum TaskLifecycleEvent {
407 Proposed {
408 task_id: String,
409 task_kind: String,
412 },
413 Accepted {
414 task_id: String,
415 },
416 Started {
417 task_id: String,
418 #[serde(default, skip_serializing_if = "Option::is_none")]
420 span_id: Option<String>,
421 },
422 Scheduled {
423 task_id: String,
424 idempotency_key: String,
425 },
426 SideEffectIntent {
427 task_id: String,
428 idempotency_key: String,
429 operation: String,
430 policy_decision: String,
431 parent_task_id: Option<String>,
432 cancellation_handle: Option<Value>,
433 },
434 Status {
437 task_id: String,
438 message: String,
439 details: Value,
440 },
441 Output {
443 task_id: String,
444 chunk: String,
445 },
446 Completed {
447 task_id: String,
448 },
449 Cancelled {
450 task_id: String,
451 reason: String,
452 },
453 Rejected {
454 task_id: String,
455 reason: String,
456 },
457 Failed {
458 task_id: String,
459 reason: String,
460 },
461 #[serde(untagged)]
463 Unknown(Value),
464}
465
466#[cfg(test)]
467mod tests {
468 use super::*;
469 use serde_json::json;
470
471 #[test]
472 fn unknown_payload_type_is_preserved_not_error() {
473 let p = MusePayload::from_parts("subagent.lifecycle.spawned", json!({"x": 1})).unwrap();
474 match p {
475 MusePayload::Unknown {
476 payload_type,
477 payload,
478 } => {
479 assert_eq!(payload_type, "subagent.lifecycle.spawned");
480 assert_eq!(payload, json!({"x": 1}));
481 }
482 other => panic!("expected Unknown, got {other:?}"),
483 }
484 }
485
486 #[test]
487 fn task_lifecycle_failed_carries_reason() {
488 let e: TaskLifecycleEvent = serde_json::from_value(json!({
489 "kind": "failed",
490 "task_id": "t1",
491 "reason": "provider does not support base instructions"
492 }))
493 .unwrap();
494 assert!(matches!(e, TaskLifecycleEvent::Failed { ref reason, .. }
495 if reason.contains("base instructions")));
496 }
497
498 #[test]
499 fn command_result_parses_real_wire_shape_and_preserves_extensions() {
500 let result: ToolResult = serde_json::from_value(json!({
501 "kind": "tool_result",
502 "command_id": "cmd-1",
503 "run_stream": { "id": "run-1", "kind": "run" },
504 "call_id": "call-1",
505 "correlation_facts": { "outcome": "success", "tool_name": "bash" },
506 "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}"#
507 }))
508 .unwrap();
509
510 assert!(result.is_command_tool());
511 let command = result.command_result().expect("typed command result");
512 assert_eq!(command.command, "printf ok");
513 assert_eq!(command.output, "ok");
514 assert_eq!(command.exit_code, 0);
515 assert_eq!(command.extra["provider_extension"], true);
516 assert_eq!(
517 serde_json::to_value(command).unwrap()["provider_extension"],
518 true
519 );
520 }
521
522 #[test]
523 fn command_result_rejects_prose_and_recognizes_command_alias() {
524 let mut result: ToolResult = serde_json::from_value(json!({
525 "kind": "tool_result",
526 "command_id": "cmd-1",
527 "run_stream": { "id": "run-1", "kind": "run" },
528 "call_id": "call-1",
529 "correlation_facts": { "outcome": "failure", "tool_name": "command" },
530 "text": "tool failed before the command started"
531 }))
532 .unwrap();
533
534 assert!(result.is_command_tool());
535 assert!(result.command_result().is_none());
536 assert!(result.try_command_result().is_err());
537
538 result.correlation_facts = Some(ToolCorrelationFacts {
539 tool_name: Some("write_file".into()),
540 ..Default::default()
541 });
542 assert!(!result.is_command_tool());
543 }
544}