use std::io::Read;
use crate::canonical::{CanonicalError, ErrorKind, Event};
use crate::config::PartialConfig;
use crate::pipeline::Sink;
use crate::protocol::{Framing, WireRequest};
use crate::registry::Registry;
use crate::transport::TransportResponse;
use super::events::{exit_from_status, fail_inband, transport_err, write_event};
use super::Host;
pub(super) fn serve_raw(
reader: &mut dyn Read,
merged: PartialConfig,
sink: &mut dyn Sink,
host: &Host,
) -> u8 {
let bytes = match read_to_vec(reader) {
Ok(b) => b,
Err(e) => return fail_inband(sink, e),
};
let cfg = match merged.into_resolved(None) {
Ok(c) => c,
Err(e) => return fail_inband(sink, e.into()),
};
let registry = Registry::builtin();
let proto = registry.protocol(cfg.provider.protocol);
let auth = registry.auth(cfg.provider.auth);
let beta: Vec<(&str, &str)> = cfg
.provider
.beta_headers
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
let ctx = cfg.provider_ctx(&beta);
let authc = cfg.auth_ctx();
let mut wire = WireRequest::new(format!("{}{}", ctx.base_url, proto.path(&ctx)), bytes);
wire.set_header("content-type", proto.content_type());
for (k, v) in ctx.beta_headers {
wire.set_header(k, v);
}
wire.timeouts = cfg.timeouts();
if let Err(e) = auth.apply(
&mut wire,
&ctx,
&authc,
host.store,
host.clock,
host.transport,
) {
return fail_inband(sink, e);
}
let resp = match host.transport.send(wire) {
Ok(r) => r,
Err(e) => return fail_inband(sink, e),
};
stream_raw(sink, resp)
}
fn stream_raw(sink: &mut dyn Sink, resp: TransportResponse) -> u8 {
let mut exit = exit_from_status(resp.status);
let mut framer = Framing::Identity.decoder();
let outcome = (|| {
for chunk in resp.body {
match chunk {
Ok(c) => {
for frame in framer.push(c).unwrap_or_default() {
write_event(sink, Event::Raw(frame.into_bytes()), &mut exit)?;
}
}
Err(_) => {
let err = transport_err("transport stream dropped");
return write_event(sink, Event::Error(err), &mut exit);
}
}
}
Ok(())
})();
match outcome.and_then(|()| write_event(sink, Event::End, &mut exit)) {
Ok(()) => exit,
Err(code) => code,
}
}
fn read_to_vec(reader: &mut dyn Read) -> Result<Vec<u8>, CanonicalError> {
let mut buf = Vec::new();
reader.read_to_end(&mut buf).map_err(|e| CanonicalError {
kind: ErrorKind::ParseInput,
message: format!("failed to read stdin: {e}"),
provider_detail: None,
})?;
Ok(buf)
}