mod stream;
use crate::canonical::{parse_retry_after, CanonicalError, ErrorKind, Event, ExitClass};
use crate::pipeline::Sink;
use crate::protocol::{DecodeState, Frame, Protocol};
use crate::transport::TransportResponse;
pub(super) fn response_events(
proto: &'static dyn Protocol,
resp: TransportResponse,
streamed: bool,
hint: Option<String>,
now: u64,
) -> Box<dyn Iterator<Item = Event>> {
let status = resp.status;
let mut state = DecodeState::default();
if !is_2xx(status) {
let retry_after = resp
.retry_after
.as_deref()
.and_then(|h| parse_retry_after(h, now));
let events = match super::drain(resp.body) {
Ok(data) => {
let frame = Frame {
event: None,
data,
status: Some(status),
};
with_hint(proto.decode(frame, &mut state), hint.as_deref())
}
Err(_) => vec![Event::Error(transport_err(
"failed to read error response body",
))],
};
Box::new(stamp_retry_after(events, retry_after).into_iter())
} else if !streamed {
let events = match super::drain(resp.body) {
Ok(data) => ensure_terminal(flatten(proto.decode_full(&data, &mut state))),
Err(_) => vec![Event::Error(transport_err("failed to read response body"))],
};
Box::new(events.into_iter())
} else {
Box::new(stream::StreamEvents::new(proto, resp.body))
}
}
fn with_hint(result: Result<Vec<Event>, CanonicalError>, hint: Option<&str>) -> Vec<Event> {
let events = flatten(result);
match hint {
Some(h) => events.into_iter().map(|ev| append_hint(ev, h)).collect(),
None => events,
}
}
fn append_hint(ev: Event, hint: &str) -> Event {
match ev {
Event::Error(mut e) => {
e.message = format!("{}; {hint}", e.message);
Event::Error(e)
}
other => other,
}
}
fn stamp_retry_after(events: Vec<Event>, secs: Option<u32>) -> Vec<Event> {
match secs {
None => events,
Some(_) => events
.into_iter()
.map(|ev| match ev {
Event::Error(mut e) => {
e.retry_after_seconds = secs;
Event::Error(e)
}
other => other,
})
.collect(),
}
}
fn flatten(result: Result<Vec<Event>, CanonicalError>) -> Vec<Event> {
result.unwrap_or_else(|e| vec![Event::Error(e)])
}
fn ensure_terminal(mut events: Vec<Event>) -> Vec<Event> {
let has_verdict = events
.iter()
.any(|e| matches!(e, Event::Finish { .. } | Event::Error(_)));
if !has_verdict {
events.push(Event::Error(transport_err(
"non-stream response carried no completion (empty or finish-less aggregate)",
)));
}
events
}
pub(super) fn write_event(sink: &mut dyn Sink, ev: Event, exit: &mut u8) -> Result<(), u8> {
if let Event::Error(e) = &ev {
*exit = e.exit_code();
}
sink.write(&ev).map_err(|io| ExitClass::from_io(&io).code())
}
pub(super) fn fail_inband(sink: &mut dyn Sink, err: CanonicalError) -> u8 {
let mut exit = err.exit_code();
match write_event(sink, Event::Error(err), &mut exit)
.and_then(|()| write_event(sink, Event::End, &mut exit))
{
Ok(()) => exit,
Err(code) => code,
}
}
pub(super) fn transport_err(message: &str) -> CanonicalError {
CanonicalError {
kind: ErrorKind::Transport,
message: message.to_owned(),
provider_detail: None,
retry_after_seconds: None,
}
}
pub(super) fn premature_eof_with_body(head: &[u8]) -> CanonicalError {
let mut err = transport_err("premature upstream EOF");
err.provider_detail = body_detail(head);
err
}
fn body_detail(data: &[u8]) -> Option<serde_json::Value> {
if data.is_empty() {
return None;
}
match serde_json::from_slice::<serde_json::Value>(data) {
Ok(v) => Some(v),
Err(_) => Some(serde_json::Value::String(
String::from_utf8_lossy(data).trim().to_owned(),
)),
}
}
pub(super) fn is_2xx(status: u16) -> bool {
(200..300).contains(&status)
}
pub(super) fn exit_from_status(status: u16) -> u8 {
if is_2xx(status) {
ExitClass::Ok.code()
} else {
ExitClass::from_kind(ErrorKind::from_http_status(status)).code()
}
}
#[cfg(test)]
mod tests {
use super::{append_hint, stamp_retry_after, with_hint};
use crate::canonical::{CanonicalError, ErrorKind, Event};
fn err(msg: &str) -> CanonicalError {
CanonicalError {
kind: ErrorKind::Provider { status: 404 },
message: msg.to_owned(),
provider_detail: None,
retry_after_seconds: None,
}
}
fn msg(ev: Event) -> String {
match ev {
Event::Error(e) => e.message,
_ => "<non-error>".to_owned(),
}
}
#[test]
fn the_hint_enriches_only_the_error_and_no_ops_otherwise() {
assert_eq!(msg(append_hint(Event::Error(err("boom")), "x")), "boom; x");
assert_eq!(msg(append_hint(Event::End, "x")), "<non-error>");
let kept = with_hint(Ok(vec![Event::Error(err("a"))]), None);
assert_eq!(kept.into_iter().map(msg).collect::<Vec<_>>(), ["a"]);
let hinted = with_hint(Ok(vec![Event::Error(err("a"))]), Some("x"));
assert_eq!(hinted.into_iter().map(msg).collect::<Vec<_>>(), ["a; x"]);
let folded = with_hint(Err(err("y")), Some("x"));
assert_eq!(folded.into_iter().map(msg).collect::<Vec<_>>(), ["y; x"]);
}
#[test]
fn the_retry_after_stamps_only_errors_and_no_ops_on_none() {
let retry = |ev: &Event| match ev {
Event::Error(e) => e.retry_after_seconds,
_ => None,
};
let stamped = stamp_retry_after(vec![Event::Error(err("boom")), Event::End], Some(30));
assert_eq!(retry(&stamped[0]), Some(30));
assert_eq!(retry(&stamped[1]), None); let untouched = stamp_retry_after(vec![Event::Error(err("boom"))], None);
assert_eq!(retry(&untouched[0]), None);
}
}