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