use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MuseRecord {
pub schema_version: u32,
pub id: String,
pub stream: StreamRef,
pub sequence: u64,
pub recorded_at: u64,
pub record_type: RecordType,
pub durability: Durability,
pub causation_id: String,
pub payload_type: String,
pub payload_schema_version: u32,
pub payload: Value,
}
impl MuseRecord {
pub fn typed_payload(&self) -> serde_json::Result<MusePayload> {
MusePayload::from_parts(&self.payload_type, self.payload.clone())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct StreamRef {
pub kind: StreamKind,
pub id: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum StreamKind {
Session,
Run,
Task,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RecordType {
Reconciliation,
Event,
Status,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Durability {
Durable,
Ephemeral,
}
#[derive(Debug, Clone, PartialEq)]
pub enum MusePayload {
CommandAccepted(CommandAccepted),
SessionRunLinked(SessionRunLinked),
TurnInputUser(TurnInputUser),
RunStarted(RunStarted),
ModelConfigured(ModelConfigured),
RunOutputDelta(RunOutputDelta),
ToolResult(ToolResult),
RunTerminal(RunTerminal),
TaskStreamLinked(TaskStreamLinked),
TaskLifecycle(TaskLifecycle),
Unknown {
payload_type: String,
payload: Value,
},
}
impl MusePayload {
pub fn from_parts(payload_type: &str, payload: Value) -> serde_json::Result<Self> {
Ok(match payload_type {
"runtime.command.accepted" => {
MusePayload::CommandAccepted(serde_json::from_value(payload)?)
}
"session.run.linked" => MusePayload::SessionRunLinked(serde_json::from_value(payload)?),
"turn.input.user" => MusePayload::TurnInputUser(serde_json::from_value(payload)?),
"run.lifecycle.started" => MusePayload::RunStarted(serde_json::from_value(payload)?),
"run.model.configured" => {
MusePayload::ModelConfigured(serde_json::from_value(payload)?)
}
"tool.result" => MusePayload::ToolResult(serde_json::from_value(payload)?),
"run.output.delta" => MusePayload::RunOutputDelta(serde_json::from_value(payload)?),
t if t.starts_with("run.terminal.") => {
MusePayload::RunTerminal(serde_json::from_value(payload)?)
}
"task.stream.linked" => MusePayload::TaskStreamLinked(serde_json::from_value(payload)?),
t if t.starts_with("task.lifecycle.") => {
MusePayload::TaskLifecycle(serde_json::from_value(payload)?)
}
other => MusePayload::Unknown {
payload_type: other.to_string(),
payload,
},
})
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CommandAccepted {
pub kind: String,
pub command_id: String,
pub command_kind: String,
pub client_id: Option<String>,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SessionRunLinked {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TurnInputUser {
pub kind: String,
pub command_id: String,
pub prompt: String,
pub run_stream: StreamRef,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunStarted {
pub kind: String,
pub command_id: String,
pub prompt: String,
pub run_stream: StreamRef,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunOutputDelta {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub text: String,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ModelConfigured {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub model_id: String,
pub display_label: String,
pub profile_id: String,
pub provider_id: String,
pub source: String,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolResult {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub call_id: String,
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub correlation_facts: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub edit_facts: Option<Value>,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
impl ToolResult {
pub fn is_command_tool(&self) -> bool {
matches!(
self.correlation_facts
.as_ref()
.and_then(|v| v.get("tool_name"))
.and_then(|v| v.as_str()),
Some("bash" | "command")
)
}
pub fn command_result(&self) -> Option<CommandResult> {
serde_json::from_str(&self.text).ok()
}
pub fn try_command_result(&self) -> Result<CommandResult, serde_json::Error> {
serde_json::from_str(&self.text)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CommandResult {
pub chunk_id: String,
pub command: String,
pub description: String,
pub exit_code: i32,
pub terminal_status: String,
pub output: String,
pub original_output_bytes: u64,
pub original_output_tokens: u64,
pub truncated: bool,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunTerminal {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub terminal: String,
pub reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub text: Option<String>,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TaskStreamLinked {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub task_id: String,
pub task_stream: StreamRef,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TaskLifecycle {
pub kind: String,
pub command_id: String,
pub run_stream: StreamRef,
pub task_id: String,
pub task_stream: StreamRef,
pub event: TaskLifecycleEvent,
#[serde(flatten, default, skip_serializing_if = "serde_json::Map::is_empty")]
pub extra: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum TaskLifecycleEvent {
Proposed {
task_id: String,
task_kind: String,
},
Accepted {
task_id: String,
},
Started {
task_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
span_id: Option<String>,
},
Scheduled {
task_id: String,
idempotency_key: String,
},
SideEffectIntent {
task_id: String,
idempotency_key: String,
operation: String,
policy_decision: String,
parent_task_id: Option<String>,
cancellation_handle: Option<Value>,
},
Status {
task_id: String,
message: String,
details: Value,
},
Output {
task_id: String,
chunk: String,
},
Completed {
task_id: String,
},
Cancelled {
task_id: String,
reason: String,
},
Rejected {
task_id: String,
reason: String,
},
Failed {
task_id: String,
reason: String,
},
#[serde(untagged)]
Unknown(Value),
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn unknown_payload_type_is_preserved_not_error() {
let p = MusePayload::from_parts("subagent.lifecycle.spawned", json!({"x": 1})).unwrap();
match p {
MusePayload::Unknown {
payload_type,
payload,
} => {
assert_eq!(payload_type, "subagent.lifecycle.spawned");
assert_eq!(payload, json!({"x": 1}));
}
other => panic!("expected Unknown, got {other:?}"),
}
}
#[test]
fn task_lifecycle_failed_carries_reason() {
let e: TaskLifecycleEvent = serde_json::from_value(json!({
"kind": "failed",
"task_id": "t1",
"reason": "provider does not support base instructions"
}))
.unwrap();
assert!(matches!(e, TaskLifecycleEvent::Failed { ref reason, .. }
if reason.contains("base instructions")));
}
#[test]
fn command_result_parses_real_wire_shape_and_preserves_extensions() {
let result: ToolResult = serde_json::from_value(json!({
"kind": "tool_result",
"command_id": "cmd-1",
"run_stream": { "id": "run-1", "kind": "run" },
"call_id": "call-1",
"correlation_facts": { "outcome": "success", "tool_name": "bash" },
"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}"#
}))
.unwrap();
assert!(result.is_command_tool());
let command = result.command_result().expect("typed command result");
assert_eq!(command.command, "printf ok");
assert_eq!(command.output, "ok");
assert_eq!(command.exit_code, 0);
assert_eq!(command.extra["provider_extension"], true);
assert_eq!(
serde_json::to_value(command).unwrap()["provider_extension"],
true
);
}
#[test]
fn command_result_rejects_prose_and_recognizes_command_alias() {
let mut result: ToolResult = serde_json::from_value(json!({
"kind": "tool_result",
"command_id": "cmd-1",
"run_stream": { "id": "run-1", "kind": "run" },
"call_id": "call-1",
"correlation_facts": { "outcome": "failure", "tool_name": "command" },
"text": "tool failed before the command started"
}))
.unwrap();
assert!(result.is_command_tool());
assert!(result.command_result().is_none());
assert!(result.try_command_result().is_err());
result.correlation_facts = Some(json!({ "tool_name": "write_file" }));
assert!(!result.is_command_tool());
}
}