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::drive::Sent;
use super::events::{exit_from_status, transport_err, write_event};
use super::Host;
pub(super) fn send_raw(
reader: &mut dyn Read,
merged: PartialConfig,
host: &Host,
) -> Result<Sent, CanonicalError> {
let bytes = read_to_vec(reader)?;
let streamed = peek_stream(&bytes);
let cfg = merged.into_resolved(None, Some(host.cache))?;
let registry = Registry::builtin();
let proto = registry.protocol(cfg.provider.protocol);
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 mut wire = WireRequest::new(format!("{}{}", ctx.base_url, proto.path(&ctx)), bytes);
wire.exec = proto.exec_spec(&ctx);
let resp = super::request::send(wire, &cfg, proto, &ctx, host)?;
Ok(Sent {
proto,
resp,
streamed,
hint: None,
})
}
fn peek_stream(bytes: &[u8]) -> bool {
#[derive(serde::Deserialize)]
struct StreamPeek {
stream: Option<bool>,
}
serde_json::from_slice::<StreamPeek>(bytes)
.ok()
.and_then(|p| p.stream)
.unwrap_or(true)
}
pub(super) 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,
}
}
pub(super) 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,
retry_after_seconds: None,
})?;
Ok(buf)
}