Skip to main content

agent_first_data/
protocol.rs

1use serde::Serialize;
2use serde_json::Value;
3use std::fmt;
4use std::ops::Deref;
5
6// ═══════════════════════════════════════════
7// Event Type and Build Errors (0.16 API)
8// ═══════════════════════════════════════════
9
10/// A typed, strict-valid AFDATA protocol v1 event.
11///
12/// Wraps the complete JSON envelope and provides access to the underlying value.
13/// Events produced by builders are guaranteed to pass `validate_protocol_event(_, true)`.
14#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
15pub struct Event(Value);
16
17impl Event {
18    /// Access the underlying JSON value.
19    pub fn as_value(&self) -> &Value {
20        &self.0
21    }
22
23    /// Convert into the underlying JSON value.
24    pub fn into_value(self) -> Value {
25        self.0
26    }
27}
28
29impl From<Event> for Value {
30    fn from(event: Event) -> Self {
31        event.0
32    }
33}
34
35impl Deref for Event {
36    type Target = Value;
37
38    fn deref(&self) -> &Self::Target {
39        &self.0
40    }
41}
42
43impl fmt::Display for Event {
44    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
45        write!(f, "{}", self.0)
46    }
47}
48
49/// Error type for builder failures.
50///
51/// Errors occur when:
52/// - A reserved field is overwritten (code, message, hint, retryable for error; message for progress/log; code is deleted from log vocabulary)
53/// - An object field (.fields or .extend) is not a JSON object
54/// - A required field (code/message for convenience functions) is empty
55#[derive(Clone, Debug, PartialEq, Eq)]
56pub enum BuildError {
57    ReservedField(String),
58    NonObjectField(String),
59    EmptyRequiredField(String),
60}
61
62impl fmt::Display for BuildError {
63    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
64        match self {
65            Self::ReservedField(msg) => write!(f, "reserved field: {msg}"),
66            Self::NonObjectField(msg) => write!(f, "non-object field: {msg}"),
67            Self::EmptyRequiredField(msg) => write!(f, "empty required field: {msg}"),
68        }
69    }
70}
71
72impl std::error::Error for BuildError {}
73
74/// Log level enumeration (serialized as lowercase).
75#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
76#[serde(rename_all = "lowercase")]
77pub enum LogLevel {
78    Debug,
79    Info,
80    Warn,
81    Error,
82}
83
84impl LogLevel {
85    fn as_str(&self) -> &'static str {
86        match self {
87            Self::Debug => "debug",
88            Self::Info => "info",
89            Self::Warn => "warn",
90            Self::Error => "error",
91        }
92    }
93}
94
95// ═══════════════════════════════════════════
96// Result Builder
97// ═══════════════════════════════════════════
98
99/// Builder for result events.
100pub struct ResultBuilder {
101    payload: Value,
102    trace: Option<Value>,
103}
104
105impl ResultBuilder {
106    /// Set the trace object.
107    pub fn trace(mut self, trace: Value) -> Self {
108        self.trace = Some(trace);
109        self
110    }
111
112    /// Build the event.
113    pub fn build(self) -> Result<Event, BuildError> {
114        let trace = self
115            .trace
116            .unwrap_or_else(|| Value::Object(serde_json::Map::new()));
117        let mut obj = serde_json::Map::new();
118        obj.insert("kind".to_string(), Value::String("result".to_string()));
119        obj.insert("result".to_string(), self.payload);
120        obj.insert("trace".to_string(), trace);
121        Ok(Event(Value::Object(obj)))
122    }
123}
124
125/// Fluent builder: start building a result event.
126pub fn json_result(payload: Value) -> ResultBuilder {
127    ResultBuilder {
128        payload,
129        trace: None,
130    }
131}
132
133// ═══════════════════════════════════════════
134// Error Builder
135// ═══════════════════════════════════════════
136
137/// Builder for error events.
138pub struct ErrorBuilder {
139    code: String,
140    message: String,
141    retryable: bool,
142    hint: Option<String>,
143    fields: serde_json::Map<String, Value>,
144    trace: Option<Value>,
145    build_error: Option<BuildError>,
146}
147
148impl ErrorBuilder {
149    /// Mark the error as retryable.
150    pub fn retryable(mut self) -> Self {
151        self.retryable = true;
152        self
153    }
154
155    /// Set retryable based on a boolean condition.
156    pub fn retryable_if(mut self, should_retry: bool) -> Self {
157        self.retryable = should_retry;
158        self
159    }
160
161    /// Set the hint.
162    pub fn hint(mut self, hint: &str) -> Self {
163        self.hint = Some(hint.to_string());
164        self
165    }
166
167    /// Set the hint if present.
168    pub fn hint_if_some(mut self, hint: Option<&str>) -> Self {
169        if let Some(h) = hint {
170            self.hint = Some(h.to_string());
171        }
172        self
173    }
174
175    /// Add an extension field.
176    pub fn field(mut self, name: &str, value: Value) -> Self {
177        if self.build_error.is_none() {
178            match name {
179                "code" | "message" | "hint" | "retryable" => {
180                    self.build_error = Some(BuildError::ReservedField(format!(
181                        "cannot write reserved field {name:?} to error payload"
182                    )));
183                }
184                _ => {
185                    self.fields.insert(name.to_string(), value);
186                }
187            }
188        }
189        self
190    }
191
192    /// Add multiple extension fields from a JSON object.
193    pub fn fields(mut self, fields: Value) -> Self {
194        if self.build_error.is_none() {
195            match fields {
196                Value::Object(map) => {
197                    for (k, v) in map {
198                        match k.as_str() {
199                            "code" | "message" | "hint" | "retryable" => {
200                                self.build_error = Some(BuildError::ReservedField(format!(
201                                    "cannot write reserved field {k:?} to error payload"
202                                )));
203                                return self;
204                            }
205                            _ => {
206                                self.fields.insert(k, v);
207                            }
208                        }
209                    }
210                }
211                _ => {
212                    self.build_error = Some(BuildError::NonObjectField(
213                        "fields() argument must be a JSON object".to_string(),
214                    ));
215                }
216            }
217        }
218        self
219    }
220
221    /// Extend with a serializable value (must serialize to JSON object).
222    pub fn extend<T: Serialize>(mut self, value: T) -> Self {
223        if self.build_error.is_none() {
224            match serde_json::to_value(&value) {
225                Ok(Value::Object(map)) => {
226                    for (k, v) in map {
227                        match k.as_str() {
228                            "code" | "message" | "hint" | "retryable" => {
229                                self.build_error = Some(BuildError::ReservedField(format!(
230                                    "cannot write reserved field {k:?} to error payload"
231                                )));
232                                return self;
233                            }
234                            _ => {
235                                self.fields.insert(k, v);
236                            }
237                        }
238                    }
239                }
240                Ok(_) => {
241                    self.build_error = Some(BuildError::NonObjectField(
242                        "extend() argument must serialize to a JSON object".to_string(),
243                    ));
244                }
245                Err(_) => {
246                    self.build_error = Some(BuildError::NonObjectField(
247                        "extend() argument serialization failed".to_string(),
248                    ));
249                }
250            }
251        }
252        self
253    }
254
255    /// Set the trace object.
256    pub fn trace(mut self, trace: Value) -> Self {
257        self.trace = Some(trace);
258        self
259    }
260
261    /// Build the event.
262    pub fn build(self) -> Result<Event, BuildError> {
263        if let Some(err) = self.build_error {
264            return Err(err);
265        }
266
267        let mut error_obj = self.fields;
268        error_obj.insert("code".to_string(), Value::String(self.code));
269        error_obj.insert("message".to_string(), Value::String(self.message));
270        error_obj.insert("retryable".to_string(), Value::Bool(self.retryable));
271        if let Some(h) = self.hint {
272            error_obj.insert("hint".to_string(), Value::String(h));
273        }
274
275        let trace = self
276            .trace
277            .unwrap_or_else(|| Value::Object(serde_json::Map::new()));
278        let mut obj = serde_json::Map::new();
279        obj.insert("kind".to_string(), Value::String("error".to_string()));
280        obj.insert("error".to_string(), Value::Object(error_obj));
281        obj.insert("trace".to_string(), trace);
282
283        Ok(Event(Value::Object(obj)))
284    }
285}
286
287/// Fluent builder: start building an error event.
288///
289/// Panics if `code` or `message` is empty (required by protocol contract).
290#[allow(clippy::panic)]
291pub fn json_error(code: &str, message: &str) -> ErrorBuilder {
292    if code.is_empty() {
293        panic!("json_error: code must not be empty");
294    }
295    if message.is_empty() {
296        panic!("json_error: message must not be empty");
297    }
298    ErrorBuilder {
299        code: code.to_string(),
300        message: message.to_string(),
301        retryable: false,
302        hint: None,
303        fields: serde_json::Map::new(),
304        trace: None,
305        build_error: None,
306    }
307}
308
309// ═══════════════════════════════════════════
310// Progress Builder
311// ═══════════════════════════════════════════
312
313/// Builder for progress events.
314pub struct ProgressBuilder {
315    message: String,
316    fields: serde_json::Map<String, Value>,
317    trace: Option<Value>,
318    build_error: Option<BuildError>,
319}
320
321impl ProgressBuilder {
322    /// Add an extension field.
323    pub fn field(mut self, name: &str, value: Value) -> Self {
324        if self.build_error.is_none() {
325            if name == "message" {
326                self.build_error = Some(BuildError::ReservedField(
327                    "cannot write reserved field \"message\" to progress payload".to_string(),
328                ));
329            } else {
330                self.fields.insert(name.to_string(), value);
331            }
332        }
333        self
334    }
335
336    /// Add multiple extension fields from a JSON object.
337    pub fn fields(mut self, fields: Value) -> Self {
338        if self.build_error.is_none() {
339            match fields {
340                Value::Object(map) => {
341                    for (k, v) in map {
342                        if k == "message" {
343                            self.build_error = Some(BuildError::ReservedField(
344                                "cannot write reserved field \"message\" to progress payload"
345                                    .to_string(),
346                            ));
347                            return self;
348                        }
349                        self.fields.insert(k, v);
350                    }
351                }
352                _ => {
353                    self.build_error = Some(BuildError::NonObjectField(
354                        "fields() argument must be a JSON object".to_string(),
355                    ));
356                }
357            }
358        }
359        self
360    }
361
362    /// Extend with a serializable value (must serialize to JSON object).
363    pub fn extend<T: Serialize>(mut self, value: T) -> Self {
364        if self.build_error.is_none() {
365            match serde_json::to_value(&value) {
366                Ok(Value::Object(map)) => {
367                    for (k, v) in map {
368                        if k == "message" {
369                            self.build_error = Some(BuildError::ReservedField(
370                                "cannot write reserved field \"message\" to progress payload"
371                                    .to_string(),
372                            ));
373                            return self;
374                        }
375                        self.fields.insert(k, v);
376                    }
377                }
378                Ok(_) => {
379                    self.build_error = Some(BuildError::NonObjectField(
380                        "extend() argument must serialize to a JSON object".to_string(),
381                    ));
382                }
383                Err(_) => {
384                    self.build_error = Some(BuildError::NonObjectField(
385                        "extend() argument serialization failed".to_string(),
386                    ));
387                }
388            }
389        }
390        self
391    }
392
393    /// Set the trace object.
394    pub fn trace(mut self, trace: Value) -> Self {
395        self.trace = Some(trace);
396        self
397    }
398
399    /// Build the event.
400    pub fn build(self) -> Result<Event, BuildError> {
401        if let Some(err) = self.build_error {
402            return Err(err);
403        }
404
405        let mut progress_obj = self.fields;
406        progress_obj.insert("message".to_string(), Value::String(self.message));
407
408        let trace = self
409            .trace
410            .unwrap_or_else(|| Value::Object(serde_json::Map::new()));
411        let mut obj = serde_json::Map::new();
412        obj.insert("kind".to_string(), Value::String("progress".to_string()));
413        obj.insert("progress".to_string(), Value::Object(progress_obj));
414        obj.insert("trace".to_string(), trace);
415
416        Ok(Event(Value::Object(obj)))
417    }
418}
419
420/// Fluent builder: start building a progress event.
421///
422/// Panics if `message` is empty (required by protocol contract).
423#[allow(clippy::panic)]
424pub fn json_progress(message: &str) -> ProgressBuilder {
425    if message.is_empty() {
426        panic!("json_progress: message must not be empty");
427    }
428    ProgressBuilder {
429        message: message.to_string(),
430        fields: serde_json::Map::new(),
431        trace: None,
432        build_error: None,
433    }
434}
435
436// ═══════════════════════════════════════════
437// Log Builder
438// ═══════════════════════════════════════════
439
440/// Builder for log events.
441pub struct LogBuilder {
442    level: LogLevel,
443    message: String,
444    fields: serde_json::Map<String, Value>,
445    trace: Option<Value>,
446    build_error: Option<BuildError>,
447}
448
449impl LogBuilder {
450    /// Add an extension field.
451    pub fn field(mut self, name: &str, value: Value) -> Self {
452        if self.build_error.is_none() {
453            match name {
454                "message" | "level" | "code" => {
455                    self.build_error = Some(BuildError::ReservedField(format!(
456                        "cannot write reserved field {name:?} to log payload"
457                    )));
458                }
459                _ => {
460                    self.fields.insert(name.to_string(), value);
461                }
462            }
463        }
464        self
465    }
466
467    /// Add multiple extension fields from a JSON object.
468    pub fn fields(mut self, fields: Value) -> Self {
469        if self.build_error.is_none() {
470            match fields {
471                Value::Object(map) => {
472                    for (k, v) in map {
473                        match k.as_str() {
474                            "message" | "level" | "code" => {
475                                self.build_error = Some(BuildError::ReservedField(format!(
476                                    "cannot write reserved field {k:?} to log payload"
477                                )));
478                                return self;
479                            }
480                            _ => {
481                                self.fields.insert(k, v);
482                            }
483                        }
484                    }
485                }
486                _ => {
487                    self.build_error = Some(BuildError::NonObjectField(
488                        "fields() argument must be a JSON object".to_string(),
489                    ));
490                }
491            }
492        }
493        self
494    }
495
496    /// Extend with a serializable value (must serialize to JSON object).
497    pub fn extend<T: Serialize>(mut self, value: T) -> Self {
498        if self.build_error.is_none() {
499            match serde_json::to_value(&value) {
500                Ok(Value::Object(map)) => {
501                    for (k, v) in map {
502                        match k.as_str() {
503                            "message" | "level" | "code" => {
504                                self.build_error = Some(BuildError::ReservedField(format!(
505                                    "cannot write reserved field {k:?} to log payload"
506                                )));
507                                return self;
508                            }
509                            _ => {
510                                self.fields.insert(k, v);
511                            }
512                        }
513                    }
514                }
515                Ok(_) => {
516                    self.build_error = Some(BuildError::NonObjectField(
517                        "extend() argument must serialize to a JSON object".to_string(),
518                    ));
519                }
520                Err(_) => {
521                    self.build_error = Some(BuildError::NonObjectField(
522                        "extend() argument serialization failed".to_string(),
523                    ));
524                }
525            }
526        }
527        self
528    }
529
530    /// Set the trace object.
531    pub fn trace(mut self, trace: Value) -> Self {
532        self.trace = Some(trace);
533        self
534    }
535
536    /// Build the event.
537    pub fn build(self) -> Result<Event, BuildError> {
538        if let Some(err) = self.build_error {
539            return Err(err);
540        }
541
542        let mut log_obj = self.fields;
543        log_obj.insert(
544            "level".to_string(),
545            Value::String(self.level.as_str().to_string()),
546        );
547        log_obj.insert("message".to_string(), Value::String(self.message));
548
549        let trace = self
550            .trace
551            .unwrap_or_else(|| Value::Object(serde_json::Map::new()));
552        let mut obj = serde_json::Map::new();
553        obj.insert("kind".to_string(), Value::String("log".to_string()));
554        obj.insert("log".to_string(), Value::Object(log_obj));
555        obj.insert("trace".to_string(), trace);
556
557        Ok(Event(Value::Object(obj)))
558    }
559}
560
561/// Fluent builder: start building a log event.
562///
563/// Panics if `message` is empty (required by protocol contract).
564#[allow(clippy::panic)]
565pub fn json_log(level: LogLevel, message: &str) -> LogBuilder {
566    if message.is_empty() {
567        panic!("json_log: message must not be empty");
568    }
569    LogBuilder {
570        level,
571        message: message.to_string(),
572        fields: serde_json::Map::new(),
573        trace: None,
574        build_error: None,
575    }
576}
577
578// ═══════════════════════════════════════════
579// CLI Helper
580// ═══════════════════════════════════════════
581
582/// Build a CLI error event with optional hint.
583///
584/// Equivalent to: `json_error("cli_error", message).hint_if_some(hint).build()`
585///
586/// Always returns an event with:
587/// - code: "cli_error"
588/// - retryable: false
589/// - trace: {}
590///
591/// Panics if `message` is empty (required by protocol contract).
592#[allow(clippy::panic, clippy::expect_used)]
593pub fn build_cli_error(message: &str, hint: Option<&str>) -> Event {
594    if message.is_empty() {
595        panic!("build_cli_error: message must not be empty");
596    }
597    json_error("cli_error", message)
598        .hint_if_some(hint)
599        .build()
600        .expect("build_cli_error: builder returned error unexpectedly")
601}
602
603// ═══════════════════════════════════════════
604// Validation
605// ═══════════════════════════════════════════
606
607/// Validate one protocol v1 event envelope.
608///
609/// `strict` additionally enforces the recommended strict profile: `trace` is
610/// required, and kind-specific payload shapes are checked (see
611/// [`validate_protocol_event_strict_payload`]).
612pub fn validate_protocol_event(event: &Value, strict: bool) -> Result<(), String> {
613    validate_protocol_event_base(event)?;
614    if strict {
615        validate_protocol_event_strict_payload(event)?;
616    }
617    Ok(())
618}
619
620fn validate_protocol_event_base(event: &Value) -> Result<(), String> {
621    let Some(obj) = event.as_object() else {
622        return Err("event must be a JSON object".to_string());
623    };
624    let Some(kind) = obj.get("kind").and_then(Value::as_str) else {
625        return Err("event.kind must be one of result, error, progress, log".to_string());
626    };
627    if !matches!(kind, "result" | "error" | "progress" | "log") {
628        return Err(format!("unsupported event kind {kind:?}"));
629    }
630    if !obj.contains_key(kind) {
631        return Err(format!("event payload field {kind:?} is required"));
632    }
633    for key in obj.keys() {
634        if key != "kind" && key != kind && key != "trace" {
635            return Err(format!("unexpected top-level field {key:?}"));
636        }
637    }
638    if let Some(trace) = obj.get("trace")
639        && !trace.is_object()
640    {
641        return Err("event.trace must be a JSON object when present".to_string());
642    }
643    if kind == "error" {
644        validate_error_payload(obj.get("error"))?;
645    }
646    Ok(())
647}
648
649fn validate_error_payload(error: Option<&Value>) -> Result<(), String> {
650    let Some(error) = error.and_then(Value::as_object) else {
651        return Err("event.error must be a JSON object".to_string());
652    };
653    match error.get("code").and_then(Value::as_str) {
654        Some(code) if !code.is_empty() => {}
655        _ => return Err("event.error.code must be a non-empty string".to_string()),
656    }
657    match error.get("message").and_then(Value::as_str) {
658        Some(message) if !message.is_empty() => {}
659        _ => return Err("event.error.message must be a non-empty string".to_string()),
660    }
661    if error.get("hint").is_some_and(|hint| !hint.is_string()) {
662        return Err("event.error.hint must be a string when present".to_string());
663    }
664    Ok(())
665}
666
667/// Validate a finite structured CLI event stream:
668/// `(log | progress)* -> exactly one (result | error) -> end`.
669///
670/// `strict` is forwarded to [`validate_protocol_event`] for every event.
671pub fn validate_protocol_stream(events: &[Value], strict: bool) -> Result<(), String> {
672    let mut terminal_kind: Option<&str> = None;
673    for (idx, event) in events.iter().enumerate() {
674        validate_protocol_event(event, strict).map_err(|err| format!("event {idx}: {err}"))?;
675        let Some(kind) = event.get("kind").and_then(Value::as_str) else {
676            return Err(format!("event {idx}: missing kind"));
677        };
678        match kind {
679            "log" | "progress" => {
680                if terminal_kind.is_some() {
681                    return Err(format!("event {idx}: non-terminal event after terminal"));
682                }
683            }
684            "result" | "error" => {
685                if terminal_kind.is_some() {
686                    return Err(format!("event {idx}: duplicate terminal event"));
687                }
688                terminal_kind = Some(kind);
689            }
690            _ => return Err(format!("event {idx}: unsupported event kind {kind:?}")),
691        }
692    }
693    if terminal_kind.is_none() {
694        return Err("event stream must contain exactly one terminal result or error".to_string());
695    }
696    Ok(())
697}
698
699/// Validate one protocol v1 event's payload against the recommended strict profile.
700///
701/// Assumes the base envelope shape (see [`validate_protocol_event_base`]) already passed.
702fn validate_protocol_event_strict_payload(event: &Value) -> Result<(), String> {
703    if !event.get("trace").is_some_and(Value::is_object) {
704        return Err("event.trace is required by the strict profile".to_string());
705    }
706    match event.get("kind").and_then(Value::as_str) {
707        Some("error") => validate_strict_error_payload(event.get("error")),
708        Some("log") => validate_strict_log_payload(event.get("log")),
709        Some("progress") => validate_strict_progress_payload(event.get("progress")),
710        _ => Ok(()),
711    }
712}
713
714fn validate_strict_error_payload(error: Option<&Value>) -> Result<(), String> {
715    let Some(error) = error.and_then(Value::as_object) else {
716        return Err("event.error must be a JSON object in the strict profile".to_string());
717    };
718    require_non_empty_string(error, "code", "event.error")?;
719    require_non_empty_string(error, "message", "event.error")?;
720    if error.get("retryable").and_then(Value::as_bool).is_none() {
721        return Err("event.error.retryable must be a boolean in the strict profile".to_string());
722    }
723    if error.contains_key("hint") && !error.get("hint").is_some_and(Value::is_string) {
724        return Err("event.error.hint must be a string when present".to_string());
725    }
726    Ok(())
727}
728
729fn validate_strict_log_payload(log: Option<&Value>) -> Result<(), String> {
730    let Some(log) = log.and_then(Value::as_object) else {
731        return Err("event.log must be a JSON object in the strict profile".to_string());
732    };
733
734    // 0.16 change: log.code is deleted from protocol vocabulary, must not be present
735    if log.contains_key("code") {
736        return Err(
737            "event.log.code must not be present in the strict profile (deleted in 0.16)"
738                .to_string(),
739        );
740    }
741
742    require_non_empty_string(log, "message", "event.log")?;
743    let valid_level = log
744        .get("level")
745        .and_then(Value::as_str)
746        .is_some_and(|level| matches!(level, "debug" | "info" | "warn" | "error"));
747    if !valid_level {
748        return Err(
749            "event.log.level must be one of debug, info, warn, error in the strict profile"
750                .to_string(),
751        );
752    }
753    Ok(())
754}
755
756fn validate_strict_progress_payload(progress: Option<&Value>) -> Result<(), String> {
757    let Some(progress) = progress.and_then(Value::as_object) else {
758        return Err("event.progress must be a JSON object in the strict profile".to_string());
759    };
760    require_non_empty_string(progress, "message", "event.progress")
761}
762
763fn require_non_empty_string(
764    payload: &serde_json::Map<String, Value>,
765    field: &str,
766    path: &str,
767) -> Result<(), String> {
768    if payload
769        .get(field)
770        .and_then(Value::as_str)
771        .is_some_and(|value| !value.is_empty())
772    {
773        return Ok(());
774    }
775    Err(format!(
776        "{path}.{field} must be a non-empty string in the strict profile"
777    ))
778}
779
780// ═══════════════════════════════════════════
781// Reader API: decode_protocol_event
782// ═══════════════════════════════════════════
783
784/// A decoded, strict-valid AFDATA protocol v1 event, typed by kind.
785#[derive(Clone, Debug, PartialEq)]
786pub enum DecodedEvent {
787    Result(DecodedResult),
788    Error(DecodedError),
789    Progress(DecodedProgress),
790    Log(DecodedLog),
791}
792
793/// Decoded `kind:"result"` event.
794#[derive(Clone, Debug, PartialEq)]
795pub struct DecodedResult {
796    /// The raw `result` payload value.
797    pub result: Value,
798    /// The raw `trace` object, when present.
799    pub trace: Option<Value>,
800}
801
802/// Decoded `kind:"error"` event.
803#[derive(Clone, Debug, PartialEq)]
804pub struct DecodedError {
805    pub code: String,
806    pub message: String,
807    pub retryable: bool,
808    pub hint: Option<String>,
809    /// Payload keys beyond `code`, `message`, `retryable`, `hint`.
810    pub fields: serde_json::Map<String, Value>,
811    /// The raw `trace` object, when present.
812    pub trace: Option<Value>,
813}
814
815/// Decoded `kind:"progress"` event.
816#[derive(Clone, Debug, PartialEq)]
817pub struct DecodedProgress {
818    pub message: String,
819    /// Payload keys beyond `message`.
820    pub fields: serde_json::Map<String, Value>,
821    /// The raw `trace` object, when present.
822    pub trace: Option<Value>,
823}
824
825/// Decoded `kind:"log"` event.
826#[derive(Clone, Debug, PartialEq)]
827pub struct DecodedLog {
828    pub level: LogLevel,
829    pub message: String,
830    /// Payload keys beyond `level`, `message`.
831    pub fields: serde_json::Map<String, Value>,
832    /// The raw `trace` object, when present.
833    pub trace: Option<Value>,
834}
835
836/// Error returned by [`decode_protocol_event`].
837#[derive(Clone, Debug, PartialEq, Eq)]
838pub enum EventDecodeError {
839    /// `text` is not valid JSON.
840    InvalidJson(String),
841    /// The parsed JSON value failed strict protocol validation.
842    InvalidEvent(String),
843}
844
845impl fmt::Display for EventDecodeError {
846    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
847        match self {
848            Self::InvalidJson(err) => write!(f, "invalid JSON: {err}"),
849            Self::InvalidEvent(err) => write!(f, "invalid protocol event: {err}"),
850        }
851    }
852}
853
854impl std::error::Error for EventDecodeError {}
855
856/// Parse one protocol v1 line, strict-validate it, and return a typed decoded event.
857///
858/// `text` is a single JSON text value (one protocol line), not a JSONL stream.
859pub fn decode_protocol_event(text: &str) -> Result<DecodedEvent, EventDecodeError> {
860    let value: Value =
861        serde_json::from_str(text).map_err(|err| EventDecodeError::InvalidJson(err.to_string()))?;
862    validate_protocol_event(&value, true).map_err(EventDecodeError::InvalidEvent)?;
863
864    // Strict validation above guarantees the envelope is an object with a
865    // recognized `kind` and a matching, object-shaped payload for
866    // error/progress/log; the `ok_or_else` fallbacks below are defensive only.
867    let malformed = || {
868        EventDecodeError::InvalidEvent(
869            "event passed strict validation but has an unexpected shape".to_string(),
870        )
871    };
872    let obj = value.as_object().ok_or_else(malformed)?;
873    let trace = obj.get("trace").cloned();
874    match obj.get("kind").and_then(Value::as_str) {
875        Some("result") => Ok(DecodedEvent::Result(DecodedResult {
876            result: obj.get("result").cloned().unwrap_or(Value::Null),
877            trace,
878        })),
879        Some("error") => {
880            let mut fields = obj
881                .get("error")
882                .and_then(Value::as_object)
883                .ok_or_else(malformed)?
884                .clone();
885            let code = fields
886                .remove("code")
887                .and_then(|v| v.as_str().map(str::to_string))
888                .ok_or_else(malformed)?;
889            let message = fields
890                .remove("message")
891                .and_then(|v| v.as_str().map(str::to_string))
892                .ok_or_else(malformed)?;
893            let retryable = fields
894                .remove("retryable")
895                .and_then(|v| v.as_bool())
896                .ok_or_else(malformed)?;
897            let hint = fields
898                .remove("hint")
899                .and_then(|v| v.as_str().map(str::to_string));
900            Ok(DecodedEvent::Error(DecodedError {
901                code,
902                message,
903                retryable,
904                hint,
905                fields,
906                trace,
907            }))
908        }
909        Some("progress") => {
910            let mut fields = obj
911                .get("progress")
912                .and_then(Value::as_object)
913                .ok_or_else(malformed)?
914                .clone();
915            let message = fields
916                .remove("message")
917                .and_then(|v| v.as_str().map(str::to_string))
918                .ok_or_else(malformed)?;
919            Ok(DecodedEvent::Progress(DecodedProgress {
920                message,
921                fields,
922                trace,
923            }))
924        }
925        Some("log") => {
926            let mut fields = obj
927                .get("log")
928                .and_then(Value::as_object)
929                .ok_or_else(malformed)?
930                .clone();
931            let message = fields
932                .remove("message")
933                .and_then(|v| v.as_str().map(str::to_string))
934                .ok_or_else(malformed)?;
935            let level = match fields.remove("level").as_ref().and_then(Value::as_str) {
936                Some("debug") => LogLevel::Debug,
937                Some("info") => LogLevel::Info,
938                Some("warn") => LogLevel::Warn,
939                Some("error") => LogLevel::Error,
940                _ => return Err(malformed()),
941            };
942            Ok(DecodedEvent::Log(DecodedLog {
943                level,
944                message,
945                fields,
946                trace,
947            }))
948        }
949        _ => Err(malformed()),
950    }
951}