use std::time::Instant;
use crate::{
batch::{observe_phase, TransformObserver, TransformPhase, TransformSignal},
decode::{looks_like_json, DecodeError, InputFormat},
Error, Result,
};
pub(super) fn decode_json_observed<R>(
bytes: &[u8],
format: InputFormat,
signal: TransformSignal,
observer: &mut Option<&mut dyn TransformObserver>,
decode_json: impl FnOnce(&[u8]) -> std::result::Result<R, DecodeError>,
decode_jsonl: impl FnOnce(&[u8]) -> std::result::Result<R, DecodeError>,
signal_label: &str,
) -> Result<R> {
let start = Instant::now();
let (request, phase) = match format {
InputFormat::Json => (decode_json(bytes)?, TransformPhase::JsonDecode),
InputFormat::Jsonl => (decode_jsonl(bytes)?, TransformPhase::JsonlDecode),
_ => {
return Err(Error::Decode(DecodeError::Unsupported(format!(
"expected JSON or JSONL {signal_label} input"
))));
}
};
observe_phase(observer, signal, phase, start.elapsed());
Ok(request)
}
pub(super) fn auto_dispatch<T, O>(
bytes: &[u8],
observer: &mut O,
json: impl FnOnce(&[u8], &mut O) -> Result<T>,
jsonl: impl FnOnce(&[u8], &mut O) -> Result<T>,
protobuf: impl FnOnce(&[u8], &mut O) -> Result<T>,
) -> Result<T> {
if looks_like_json(bytes) {
match json(bytes, observer) {
Ok(value) => Ok(value),
Err(json_err) => match jsonl(bytes, observer) {
Ok(value) => Ok(value),
Err(_) => protobuf(bytes, observer).map_err(|proto_err| {
Error::Decode(DecodeError::Unsupported(format!(
"json decode failed: {json_err}; protobuf fallback failed: {proto_err}"
)))
}),
},
}
} else {
match protobuf(bytes, observer) {
Ok(value) => Ok(value),
Err(proto_err) => json(bytes, observer).map_err(|json_err| {
Error::Decode(DecodeError::Unsupported(format!(
"protobuf decode failed: {proto_err}; json fallback failed: {json_err}"
)))
}),
}
}
}