use std::collections::VecDeque;
use std::io;
use crate::canonical::{CanonicalError, ErrorKind, Event, ExitClass};
use crate::pipeline::Sink;
use crate::protocol::{DecodeState, Decoder, Frame, Protocol};
use crate::transport::{Bytes, TransportResponse};
pub(super) fn response_events(
proto: &'static dyn Protocol,
resp: TransportResponse,
streamed: bool,
hint: Option<String>,
) -> Box<dyn Iterator<Item = Event>> {
let status = resp.status;
let mut state = DecodeState::default();
if !is_2xx(status) {
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(events.into_iter())
} else if !streamed {
let events = match super::drain(resp.body) {
Ok(data) => 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(StreamEvents::new(proto, resp.body))
}
}
struct StreamEvents {
proto: &'static dyn Protocol,
body: Box<dyn Iterator<Item = io::Result<Bytes>>>,
decoder: Box<dyn Decoder>,
state: DecodeState,
pending: VecDeque<Event>,
body_done: bool,
finished: bool,
}
impl StreamEvents {
fn new(
proto: &'static dyn Protocol,
body: Box<dyn Iterator<Item = io::Result<Bytes>>>,
) -> Self {
StreamEvents {
decoder: proto.framing().decoder(),
proto,
body,
state: DecodeState::default(),
pending: VecDeque::new(),
body_done: false,
finished: false,
}
}
fn decode_into_pending(&mut self, frame: Frame) {
match self.proto.decode(frame, &mut self.state) {
Ok(events) => self.pending.extend(events),
Err(e) => self.pending.push_back(Event::Error(e)),
}
}
}
impl Iterator for StreamEvents {
type Item = Event;
fn next(&mut self) -> Option<Event> {
loop {
if let Some(ev) = self.pending.pop_front() {
return Some(ev);
}
if self.finished {
return None;
}
if self.body_done {
let tail = self.decoder.finish().unwrap_or_default();
for frame in tail {
self.decode_into_pending(frame);
}
if !self.state.terminated {
self.pending
.push_back(Event::Error(transport_err("premature upstream EOF")));
}
self.finished = true;
continue;
}
match self.body.next() {
Some(Ok(chunk)) => {
for frame in self.decoder.push(chunk).unwrap_or_default() {
self.decode_into_pending(frame);
}
}
Some(Err(_)) => {
self.pending
.push_back(Event::Error(transport_err("transport stream dropped")));
self.finished = true;
}
None => self.body_done = true,
}
}
}
}
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 flatten(result: Result<Vec<Event>, CanonicalError>) -> Vec<Event> {
result.unwrap_or_else(|e| vec![Event::Error(e)])
}
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,
}
}
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, 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,
}
}
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"]);
}
}