| 1 | //! The one normalized stream-event format both the `sse` and `events` |
| 2 | //! families speak. |
| 3 | //! |
| 4 | //! It is exactly the serde shape `codewhale_models::StreamEvent` already |
| 5 | //! deserializes (Anthropic-style `{"type": "content_block_delta", ...}`), so a |
| 6 | //! scripted provider script and a parser golden are the same language and a |
| 7 | //! TypeScript implementation needs no Rust type to read either. `StreamEvent` |
| 8 | //! is `Deserialize`-only; this module is its serializer, and |
| 9 | //! [`round_trips`] proves the two directions agree for every golden line. |
| 10 | |
| 11 | use codewhale_models::{ContentBlockStart, Delta, StreamEvent}; |
| 12 | use serde_json::{Value, json}; |
| 13 | |
| 14 | pub(super) fn to_json(event: &StreamEvent) -> Value { |
| 15 | match event { |
| 16 | StreamEvent::ToolProjectionWarning { |
| 17 | provider, |
| 18 | omitted_tool_names, |
| 19 | omitted_tool_count, |
| 20 | } => json!({ |
| 21 | "type": "tool_projection_warning", |
| 22 | "provider": provider, |
| 23 | "omitted_tool_names": omitted_tool_names, |
| 24 | "omitted_tool_count": omitted_tool_count, |
| 25 | }), |
| 26 | StreamEvent::MessageStart { message } => json!({ |
| 27 | "type": "message_start", |
| 28 | "message": serde_json::to_value(message).expect("MessageResponse serializes"), |
| 29 | }), |
| 30 | StreamEvent::ContentBlockStart { |
| 31 | index, |
| 32 | content_block, |
| 33 | } => json!({ |
| 34 | "type": "content_block_start", |
| 35 | "index": index, |
| 36 | "content_block": block_start_json(content_block), |
| 37 | }), |
| 38 | StreamEvent::ContentBlockDelta { index, delta } => json!({ |
| 39 | "type": "content_block_delta", |
| 40 | "index": index, |
| 41 | "delta": delta_json(delta), |
| 42 | }), |
| 43 | StreamEvent::ContentBlockStop { index } => json!({ |
| 44 | "type": "content_block_stop", |
| 45 | "index": index, |
| 46 | }), |
| 47 | StreamEvent::MessageDelta { delta, usage } => { |
| 48 | let mut value = json!({ |
| 49 | "type": "message_delta", |
| 50 | "delta": { |
| 51 | "stop_reason": delta.stop_reason, |
| 52 | "stop_sequence": delta.stop_sequence, |
| 53 | }, |
| 54 | }); |
| 55 | if let Some(usage) = usage { |
| 56 | value["usage"] = serde_json::to_value(usage).expect("Usage serializes"); |
| 57 | } |
| 58 | value |
| 59 | } |
| 60 | StreamEvent::MessageStop => json!({ "type": "message_stop" }), |
| 61 | StreamEvent::Ping => json!({ "type": "ping" }), |
| 62 | StreamEvent::Error { error } => json!({ "type": "error", "error": error }), |
| 63 | } |
| 64 | } |
| 65 | |
| 66 | fn block_start_json(block: &ContentBlockStart) -> Value { |
| 67 | match block { |
| 68 | ContentBlockStart::Text { text } => json!({ "type": "text", "text": text }), |
| 69 | ContentBlockStart::Thinking { thinking } => { |
| 70 | json!({ "type": "thinking", "thinking": thinking }) |
| 71 | } |
| 72 | ContentBlockStart::ToolUse { |
| 73 | id, |
| 74 | name, |
| 75 | input, |
| 76 | caller, |
| 77 | thought_signature, |
| 78 | } => { |
| 79 | let mut value = json!({ |
| 80 | "type": "tool_use", |
| 81 | "id": id, |
| 82 | "name": name, |
| 83 | "input": input, |
| 84 | }); |
| 85 | if let Some(caller) = caller { |
| 86 | value["caller"] = serde_json::to_value(caller).expect("ToolCaller serializes"); |
| 87 | } |
| 88 | if let Some(signature) = thought_signature { |
| 89 | value["thought_signature"] = json!(signature); |
| 90 | } |
| 91 | value |
| 92 | } |
| 93 | ContentBlockStart::ServerToolUse { id, name, input } => json!({ |
| 94 | "type": "server_tool_use", |
| 95 | "id": id, |
| 96 | "name": name, |
| 97 | "input": input, |
| 98 | }), |
| 99 | } |
| 100 | } |
| 101 | |
| 102 | fn delta_json(delta: &Delta) -> Value { |
| 103 | match delta { |
| 104 | Delta::TextDelta { text } => json!({ "type": "text_delta", "text": text }), |
| 105 | Delta::ThinkingDelta { thinking } => { |
| 106 | json!({ "type": "thinking_delta", "thinking": thinking }) |
| 107 | } |
| 108 | Delta::InputJsonDelta { partial_json } => { |
| 109 | json!({ "type": "input_json_delta", "partial_json": partial_json }) |
| 110 | } |
| 111 | Delta::SignatureDelta { signature } => { |
| 112 | json!({ "type": "signature_delta", "signature": signature }) |
| 113 | } |
| 114 | Delta::ReasoningStateDelta { state } => json!({ |
| 115 | "type": "reasoning_state_delta", |
| 116 | "state": serde_json::to_value(state).expect("OpaqueReasoningState serializes"), |
| 117 | }), |
| 118 | } |
| 119 | } |
| 120 | |
| 121 | pub(super) fn from_json(value: &Value) -> anyhow::Result<StreamEvent> { |
| 122 | Ok(serde_json::from_value(value.clone())?) |
| 123 | } |
| 124 | |
| 125 | /// `value` deserializes into a `StreamEvent` that serializes back to itself. |
| 126 | pub(super) fn round_trips(value: &Value) -> Result<(), String> { |
| 127 | let event = from_json(value).map_err(|error| format!("{value} does not parse: {error}"))?; |
| 128 | let back = to_json(&event); |
| 129 | if super::golden::canonical(&back) == super::golden::canonical(value) { |
| 130 | Ok(()) |
| 131 | } else { |
| 132 | Err(format!("{value} round-trips to {back}")) |
| 133 | } |
| 134 | } |
| 135 |