1use serde::Serialize;
2use serde_json::Value;
3use std::fmt;
4use std::ops::Deref;
5
6#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
15pub struct Event(Value);
16
17impl Event {
18 pub fn as_value(&self) -> &Value {
20 &self.0
21 }
22
23 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#[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#[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
95pub struct ResultBuilder {
101 payload: Value,
102 trace: Option<Value>,
103}
104
105impl ResultBuilder {
106 pub fn trace(mut self, trace: Value) -> Self {
108 self.trace = Some(trace);
109 self
110 }
111
112 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
125pub fn json_result(payload: Value) -> ResultBuilder {
127 ResultBuilder {
128 payload,
129 trace: None,
130 }
131}
132
133pub 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 pub fn retryable(mut self) -> Self {
151 self.retryable = true;
152 self
153 }
154
155 pub fn retryable_if(mut self, should_retry: bool) -> Self {
157 self.retryable = should_retry;
158 self
159 }
160
161 pub fn hint(mut self, hint: &str) -> Self {
163 self.hint = Some(hint.to_string());
164 self
165 }
166
167 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 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 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 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 pub fn trace(mut self, trace: Value) -> Self {
257 self.trace = Some(trace);
258 self
259 }
260
261 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#[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
309pub 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 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 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 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 pub fn trace(mut self, trace: Value) -> Self {
395 self.trace = Some(trace);
396 self
397 }
398
399 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#[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
436pub 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 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 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 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 pub fn trace(mut self, trace: Value) -> Self {
532 self.trace = Some(trace);
533 self
534 }
535
536 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#[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#[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
603pub 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
667pub 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
699fn 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 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#[derive(Clone, Debug, PartialEq)]
786pub enum DecodedEvent {
787 Result(DecodedResult),
788 Error(DecodedError),
789 Progress(DecodedProgress),
790 Log(DecodedLog),
791}
792
793#[derive(Clone, Debug, PartialEq)]
795pub struct DecodedResult {
796 pub result: Value,
798 pub trace: Option<Value>,
800}
801
802#[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 pub fields: serde_json::Map<String, Value>,
811 pub trace: Option<Value>,
813}
814
815#[derive(Clone, Debug, PartialEq)]
817pub struct DecodedProgress {
818 pub message: String,
819 pub fields: serde_json::Map<String, Value>,
821 pub trace: Option<Value>,
823}
824
825#[derive(Clone, Debug, PartialEq)]
827pub struct DecodedLog {
828 pub level: LogLevel,
829 pub message: String,
830 pub fields: serde_json::Map<String, Value>,
832 pub trace: Option<Value>,
834}
835
836#[derive(Clone, Debug, PartialEq, Eq)]
838pub enum EventDecodeError {
839 InvalidJson(String),
841 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
856pub 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 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}