1use serde::Serialize;
8use serde::de::DeserializeOwned;
9use std::io::Write;
10
11use crate::shared::error::Error;
12
13pub fn emit<W: Write, T: Serialize>(writer: &mut W, code: &str, payload: &T) -> Result<(), Error> {
20 emit_inner(writer, code, payload, RedactionMode::Default)
21}
22
23pub fn emit_unredacted<W: Write, T: Serialize>(
26 writer: &mut W,
27 code: &str,
28 payload: &T,
29) -> Result<(), Error> {
30 emit_inner(writer, code, payload, RedactionMode::None)
31}
32
33#[derive(Clone, Copy)]
34enum RedactionMode {
35 Default,
36 None,
37}
38
39fn emit_inner<W: Write, T: Serialize>(
40 writer: &mut W,
41 code: &str,
42 payload: &T,
43 redaction: RedactionMode,
44) -> Result<(), Error> {
45 let mut value = serde_json::to_value(payload).map_err(|e| {
46 Error::new(
47 crate::shared::error::ErrorCode::InternalError,
48 format!("AFDATA: failed to serialize payload: {e}"),
49 )
50 })?;
51 value = wrap_payload(code, value)?;
52
53 let options = match redaction {
54 RedactionMode::Default => agent_first_data::OutputOptions::default(),
55 RedactionMode::None => agent_first_data::OutputOptions {
56 redaction: agent_first_data::Redactor::new()
57 .policy(agent_first_data::RedactionPolicy::RedactionNone),
58 style: agent_first_data::OutputStyle::Raw,
59 },
60 };
61 let mut emitter = agent_first_data::CliEmitter::with_options(
62 writer,
63 agent_first_data::OutputFormat::Json,
64 options,
65 )
66 .with_strict_protocol();
67 emitter.emit_result(value).map_err(|err| {
68 Error::new(
69 crate::shared::error::ErrorCode::InternalError,
70 err.to_string(),
71 )
72 })?;
73 Ok(())
74}
75
76fn wrap_payload(code: &str, value: serde_json::Value) -> Result<serde_json::Value, Error> {
77 let serde_json::Value::Object(mut map) = value else {
78 return Err(Error::new(
79 crate::shared::error::ErrorCode::InternalError,
80 "AFDATA result payload must serialize to a JSON object",
81 ));
82 };
83
84 map.insert("code".into(), serde_json::Value::String(code.to_string()));
85 Ok(serde_json::Value::Object(map))
86}
87
88pub fn emit_error<W: Write>(writer: &mut W, err: &Error) -> Result<(), Error> {
90 let mut emitter =
91 agent_first_data::CliEmitter::new(writer, agent_first_data::OutputFormat::Json)
92 .with_strict_protocol();
93 let event = agent_first_data::json_error(err.error_code.as_str(), &err.detail)
94 .retryable_if(err.retryable)
95 .build()
96 .map_err(|err| {
97 Error::new(
98 crate::shared::error::ErrorCode::InternalError,
99 err.to_string(),
100 )
101 })?;
102 emitter.emit(event).map_err(|emit_err| {
103 Error::new(
104 crate::shared::error::ErrorCode::InternalError,
105 emit_err.to_string(),
106 )
107 })
108}
109
110pub fn emit_error_with<W: Write>(
112 writer: &mut W,
113 code: &str,
114 message: &str,
115 fields: serde_json::Value,
116 trace: serde_json::Value,
117) -> Result<(), Error> {
118 let mut emitter =
119 agent_first_data::CliEmitter::new(writer, agent_first_data::OutputFormat::Json)
120 .with_strict_protocol();
121 let retryable = fields
122 .get("retryable")
123 .and_then(serde_json::Value::as_bool)
124 .unwrap_or(false);
125 let fields = match fields {
126 serde_json::Value::Object(mut fields) => {
127 fields.remove("retryable");
128 serde_json::Value::Object(fields)
129 }
130 other => other,
131 };
132 let event = agent_first_data::json_error(code, message)
133 .retryable_if(retryable)
134 .fields(fields)
135 .trace(trace)
136 .build()
137 .map_err(|err| {
138 Error::new(
139 crate::shared::error::ErrorCode::InternalError,
140 err.to_string(),
141 )
142 })?;
143 emitter.emit(event).map_err(|err| {
144 Error::new(
145 crate::shared::error::ErrorCode::InternalError,
146 err.to_string(),
147 )
148 })
149}
150
151pub fn error_value(code: &str, message: &str, retryable: bool) -> serde_json::Value {
153 agent_first_data::json_error(code, message)
154 .retryable_if(retryable)
155 .build()
156 .map(Into::into)
157 .unwrap_or_else(|_| serde_json::json!({}))
158}
159
160pub fn result_value(code: &str, mut payload: serde_json::Value) -> serde_json::Value {
162 let payload = match &mut payload {
163 serde_json::Value::Object(fields) => {
164 fields
165 .entry("code".to_string())
166 .or_insert_with(|| serde_json::Value::String(code.to_string()));
167 payload
168 }
169 _ => serde_json::json!({"code": code, "value": payload}),
170 };
171 agent_first_data::json_result(payload)
172 .build()
173 .map(Into::into)
174 .unwrap_or_else(|_| serde_json::json!({}))
175}
176
177pub fn decode_result<T: DeserializeOwned>(bytes: &[u8]) -> Result<T, Error> {
178 let text = std::str::from_utf8(bytes).map_err(|error| {
179 Error::new(
180 crate::shared::error::ErrorCode::InternalError,
181 format!("decode AFDATA event: {error}"),
182 )
183 })?;
184 match agent_first_data::decode_protocol_event(text) {
185 Ok(agent_first_data::DecodedEvent::Result(result)) => serde_json::from_value(result.result)
186 .map_err(|error| {
187 Error::new(
188 crate::shared::error::ErrorCode::InternalError,
189 format!("decode AFDATA result payload: {error}"),
190 )
191 }),
192 Ok(_) => Err(Error::new(
193 crate::shared::error::ErrorCode::InternalError,
194 "expected AFDATA result event",
195 )),
196 Err(error) => Err(Error::new(
197 crate::shared::error::ErrorCode::InternalError,
198 format!("invalid AFDATA event: {error}"),
199 )),
200 }
201}
202
203pub fn decode_error(bytes: &[u8]) -> Result<Error, Error> {
204 let text = std::str::from_utf8(bytes).map_err(|error| {
205 Error::new(
206 crate::shared::error::ErrorCode::InternalError,
207 format!("decode AFDATA error event: {error}"),
208 )
209 })?;
210 match agent_first_data::decode_protocol_event(text) {
211 Ok(agent_first_data::DecodedEvent::Error(error)) => {
212 let code = serde_json::from_value(serde_json::Value::String(error.code))
213 .unwrap_or(crate::shared::error::ErrorCode::InternalError);
214 Ok(Error::new(code, error.message).with_retryable(error.retryable))
215 }
216 Ok(_) => Err(Error::new(
217 crate::shared::error::ErrorCode::InternalError,
218 "expected AFDATA error event",
219 )),
220 Err(error) => Err(Error::new(
221 crate::shared::error::ErrorCode::InternalError,
222 format!("invalid AFDATA error event: {error}"),
223 )),
224 }
225}
226
227#[cfg(test)]
228mod tests {
229 use super::*;
230
231 #[derive(Serialize)]
232 struct HealthPayload {
233 status: &'static str,
234 uptime_s: u64,
235 }
236
237 #[test]
238 fn json_result_event_is_single_line_with_code_field() {
239 let mut buf = Vec::new();
240 let payload = HealthPayload {
241 status: "ok",
242 uptime_s: 42,
243 };
244 emit(&mut buf, "health", &payload).unwrap();
245 let s = String::from_utf8(buf).unwrap_or_default();
246 assert!(s.ends_with('\n'));
247 let trimmed = s.trim_end();
248 let parsed: serde_json::Value = serde_json::from_str(trimmed).unwrap();
249 assert_eq!(parsed["kind"], "result");
250 assert_eq!(parsed["result"]["code"], "health");
251 assert_eq!(parsed["result"]["status"], "ok");
252 assert_eq!(parsed["result"]["uptime_s"], 42);
253 assert_eq!(trimmed.lines().count(), 1);
254 }
255
256 #[test]
257 fn error_event_uses_error_code_tag() {
258 let mut buf = Vec::new();
259 let err = Error::new(
260 crate::shared::error::ErrorCode::NavigationTimeout,
261 "no load",
262 );
263 emit_error(&mut buf, &err).unwrap();
264 let parsed: serde_json::Value =
265 serde_json::from_slice(&buf).unwrap_or(serde_json::Value::Null);
266 assert_eq!(parsed["kind"], "error");
267 assert_eq!(parsed["error"]["code"], "navigation_timeout");
268 assert_eq!(parsed["error"]["message"], "no load");
269 assert_eq!(parsed["error"]["retryable"], true);
270 }
271
272 #[test]
273 fn error_extension_fields_are_flattened_into_error_payload() {
274 let mut buf = Vec::new();
275 emit_error_with(
276 &mut buf,
277 "navigation_timeout",
278 "no load",
279 serde_json::json!({
280 "retryable": true,
281 "stage": "capture_text",
282 "details": "scalar detail remains an explicitly named field"
283 }),
284 serde_json::json!({"duration_ms": 10}),
285 )
286 .unwrap();
287 let parsed: serde_json::Value = serde_json::from_slice(&buf).unwrap();
288 assert_eq!(parsed["error"]["stage"], "capture_text");
289 assert_eq!(parsed["error"]["retryable"], true);
290 assert_eq!(
291 parsed["error"]["details"],
292 "scalar detail remains an explicitly named field"
293 );
294 assert!(parsed["error"].get("fields").is_none());
295 }
296
297 #[derive(Serialize)]
298 struct SecretPayload {
299 token_secret: &'static str,
300 }
301
302 #[test]
303 fn afdata_event_redacts_secret_fields() {
304 let mut buf = Vec::new();
305 emit(
306 &mut buf,
307 "container_status",
308 &SecretPayload {
309 token_secret: "supersecret",
310 },
311 )
312 .unwrap();
313 let parsed: serde_json::Value = serde_json::from_slice(&buf).unwrap();
314 assert_eq!(parsed["result"]["token_secret"], "***");
315 }
316
317 #[test]
318 fn envelope_uses_sdk_result_payload_without_nested_envelope() {
319 let mut buf = Vec::new();
320 let payload = serde_json::json!({"code": -32000, "message": "cdp error"});
321 emit(&mut buf, "cdp", &payload).unwrap();
322 let parsed: serde_json::Value = serde_json::from_slice(&buf).unwrap();
323 assert_eq!(parsed["kind"], "result");
324 assert_eq!(parsed["result"]["code"], "cdp");
325 assert_eq!(parsed["result"]["code"], "cdp");
326 assert_eq!(parsed["result"]["message"], "cdp error");
327 assert!(parsed["result"].get("result").is_none());
328 }
329}