claude_code_sdk_rust/internal/
parser.rs1use crate::error::{ClaudeSDKError, MessageParseError, Result};
2use crate::types::{
3 HookEventMessage, Message, MirrorErrorMessage, TaskNotificationMessage, TaskProgressMessage,
4 TaskStartedMessage, TaskUpdatedMessage,
5};
6
7const KNOWN_MESSAGE_TYPES: &[&str] = &[
8 "user",
9 "assistant",
10 "system",
11 "result",
12 "stream_event",
13 "rate_limit_event",
14];
15
16pub fn parse_message_line(line: &str) -> Result<Option<Message>> {
17 let value = serde_json::from_str::<serde_json::Value>(line)?;
18 parse_message_value(value)
19}
20
21pub fn parse_message_value(value: serde_json::Value) -> Result<Option<Message>> {
22 let message_type = value.get("type").and_then(|v| v.as_str()).ok_or_else(|| {
23 let data = value.as_object().cloned();
24 let mut error = MessageParseError::new("Message missing 'type' field");
25 if let Some(data) = data {
26 error = error.with_data(data);
27 }
28 ClaudeSDKError::MessageParse(error)
29 })?;
30
31 if !KNOWN_MESSAGE_TYPES.contains(&message_type) {
32 return Ok(None);
33 }
34
35 if message_type == "system" {
36 return parse_system_message_value(value);
37 }
38
39 match serde_json::from_value::<Message>(value.clone()) {
40 Ok(message) => Ok(Some(message)),
41 Err(err) => Err(parse_error_with_payload(err, &value)),
42 }
43}
44
45fn parse_error_with_payload(err: serde_json::Error, value: &serde_json::Value) -> ClaudeSDKError {
48 let payload = value.to_string();
49 let payload = if payload.len() > 600 {
50 let cut = payload
51 .char_indices()
52 .take_while(|(idx, _)| *idx <= 600)
53 .last()
54 .map(|(idx, ch)| idx + ch.len_utf8())
55 .unwrap_or(payload.len());
56 format!("{}...", &payload[..cut])
57 } else {
58 payload
59 };
60 let mut error = MessageParseError::new(format!(
61 "Failed to parse CLI message: {err}; payload: {payload}"
62 ));
63 if let Some(data) = value.as_object() {
64 error = error.with_data(data.clone());
65 }
66 ClaudeSDKError::MessageParse(error)
67}
68
69fn parse_system_message_value(value: serde_json::Value) -> Result<Option<Message>> {
70 let subtype = value.get("subtype").and_then(|v| v.as_str());
71 match subtype {
72 Some("task_started") => parse_task_started(value)
73 .map(Message::TaskStartedMsg)
74 .map(Some),
75 Some("task_progress") => parse_task_progress(value)
76 .map(Message::TaskProgressMsg)
77 .map(Some),
78 Some("task_notification") => parse_task_notification(value)
79 .map(Message::TaskNotificationMsg)
80 .map(Some),
81 Some("task_updated") => parse_task_updated(value)
82 .map(Message::TaskUpdatedMsg)
83 .map(Some),
84 Some("hook_started" | "hook_response") => {
85 parse_hook_event(value).map(Message::HookEventMsg).map(Some)
86 }
87 Some("mirror_error") => parse_mirror_error(value)
88 .map(Message::MirrorErrorMsg)
89 .map(Some),
90 _ => serde_json::from_value::<Message>(value)
91 .map(Some)
92 .map_err(ClaudeSDKError::Serialization),
93 }
94}
95
96fn parse_mirror_error(value: serde_json::Value) -> Result<MirrorErrorMessage> {
97 let mut data = value.as_object().cloned().ok_or_else(|| {
98 ClaudeSDKError::MessageParse(MessageParseError::new("System message must be an object"))
99 })?;
100 data.remove("type");
101 let key = data.get("key").and_then(|value| value.as_object()).cloned();
102 let error = data
103 .get("error")
104 .and_then(|value| value.as_str())
105 .unwrap_or_default()
106 .to_string();
107 Ok(MirrorErrorMessage { key, error, data })
108}
109
110fn parse_task_started(value: serde_json::Value) -> Result<TaskStartedMessage> {
111 serde_json::from_value::<TaskStartedMessage>(strip_system_fields(value)?)
112 .map_err(ClaudeSDKError::Serialization)
113}
114
115fn parse_task_progress(value: serde_json::Value) -> Result<TaskProgressMessage> {
116 serde_json::from_value::<TaskProgressMessage>(strip_system_fields(value)?)
117 .map_err(ClaudeSDKError::Serialization)
118}
119
120fn parse_task_notification(value: serde_json::Value) -> Result<TaskNotificationMessage> {
121 serde_json::from_value::<TaskNotificationMessage>(strip_system_fields(value)?)
122 .map_err(ClaudeSDKError::Serialization)
123}
124
125fn parse_task_updated(value: serde_json::Value) -> Result<TaskUpdatedMessage> {
131 let data = value.as_object().ok_or_else(|| {
132 ClaudeSDKError::MessageParse(MessageParseError::new("System message must be an object"))
133 })?;
134 let task_id = data
135 .get("task_id")
136 .and_then(|v| v.as_str())
137 .unwrap_or_default()
138 .to_string();
139 let patch = data
140 .get("patch")
141 .and_then(|v| v.as_object())
142 .cloned()
143 .unwrap_or_default();
144 let status = patch
145 .get("status")
146 .and_then(|v| serde_json::from_value::<crate::types::TaskUpdatedStatus>(v.clone()).ok());
147 let session_id = data
148 .get("session_id")
149 .and_then(|v| v.as_str())
150 .map(|s| s.to_string());
151 let uuid = data
152 .get("uuid")
153 .and_then(|v| v.as_str())
154 .map(|s| s.to_string());
155 Ok(TaskUpdatedMessage {
156 task_id,
157 patch,
158 status,
159 session_id,
160 uuid,
161 })
162}
163
164fn parse_hook_event(value: serde_json::Value) -> Result<HookEventMessage> {
165 let mut data = value.as_object().cloned().ok_or_else(|| {
166 ClaudeSDKError::MessageParse(MessageParseError::new("System message must be an object"))
167 })?;
168 let subtype = data
169 .get("subtype")
170 .and_then(|value| value.as_str())
171 .unwrap_or_default()
172 .to_string();
173 let hook_event_name = data
174 .get("hook_event")
175 .or_else(|| data.get("hook_name"))
176 .and_then(|value| value.as_str())
177 .map(ToString::to_string);
178 let session_id = data
179 .get("session_id")
180 .and_then(|value| value.as_str())
181 .map(ToString::to_string);
182 let uuid = data
183 .get("uuid")
184 .and_then(|value| value.as_str())
185 .map(ToString::to_string);
186 data.remove("type");
187 Ok(HookEventMessage {
188 subtype,
189 hook_event_name,
190 session_id,
191 uuid,
192 data,
193 })
194}
195
196fn strip_system_fields(value: serde_json::Value) -> Result<serde_json::Value> {
197 let mut data = value.as_object().cloned().ok_or_else(|| {
198 ClaudeSDKError::MessageParse(MessageParseError::new("System message must be an object"))
199 })?;
200 data.remove("type");
201 data.remove("subtype");
202 Ok(serde_json::Value::Object(data))
203}