use std::io::{self, Write};
use std::path::PathBuf;
use bytes::{Bytes, BytesMut};
use futures_util::Stream;
use hyper::body::Frame;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use super::context::{stream_build_context, stream_build_context_with_inline};
use crate::error::Result;
type BodyItem = io::Result<Frame<Bytes>>;
const CHANNEL_CAP: usize = 8;
const CHUNK_BYTES: usize = 64 * 1024;
pub(super) enum ContextSource {
Inline(String),
Dockerfile(String),
}
struct ChannelWriter {
tx: mpsc::Sender<BodyItem>,
buf: BytesMut,
}
impl ChannelWriter {
fn send_pending(&mut self) -> io::Result<()> {
if self.buf.is_empty() {
return Ok(());
}
let chunk = self.buf.split().freeze();
self.tx.blocking_send(Ok(Frame::data(chunk))).map_err(|_| {
io::Error::new(
io::ErrorKind::BrokenPipe,
"build-context receiver dropped before the tar finished",
)
})
}
}
impl Write for ChannelWriter {
fn write(&mut self, data: &[u8]) -> io::Result<usize> {
self.buf.extend_from_slice(data);
if self.buf.len() >= CHUNK_BYTES {
self.send_pending()?;
}
Ok(data.len())
}
fn flush(&mut self) -> io::Result<()> {
self.send_pending()
}
}
pub(super) fn context_body(
context: PathBuf,
source: ContextSource,
secret_files: Vec<(String, Vec<u8>)>,
) -> (
JoinHandle<Result<()>>,
impl Stream<Item = BodyItem> + Send + 'static,
) {
let (tx, rx) = mpsc::channel::<BodyItem>(CHANNEL_CAP);
let producer = tokio::task::spawn_blocking(move || -> Result<()> {
let mut writer = ChannelWriter {
tx,
buf: BytesMut::with_capacity(CHUNK_BYTES),
};
match source {
ContextSource::Inline(inline) => {
stream_build_context_with_inline(&mut writer, &context, &inline, &secret_files)
}
ContextSource::Dockerfile(dockerfile) => {
stream_build_context(&mut writer, &context, &dockerfile, &secret_files)
}
}
});
let body = futures_util::stream::unfold(rx, |mut rx| async move {
rx.recv().await.map(|item| (item, rx))
});
(producer, body)
}
#[cfg(test)]
mod tests {
use super::super::context::build_context_tar;
use super::*;
use futures_util::StreamExt;
#[tokio::test]
async fn streamed_body_matches_buffered_tar() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("Dockerfile"), "FROM scratch\n").unwrap();
std::fs::write(dir.path().join("app.txt"), "hello world").unwrap();
let (producer, body) = context_body(
dir.path().to_path_buf(),
ContextSource::Dockerfile("Dockerfile".to_string()),
Vec::new(),
);
futures_util::pin_mut!(body);
let mut streamed = Vec::new();
while let Some(item) = body.next().await {
let frame = item.expect("no stream error");
if let Ok(data) = frame.into_data() {
streamed.extend_from_slice(&data);
}
}
producer.await.expect("join").expect("producer succeeds");
let buffered = build_context_tar(dir.path(), "Dockerfile", &[]).unwrap();
assert_eq!(
streamed, buffered,
"streamed tar must be byte-identical to the buffered tar"
);
}
#[tokio::test]
async fn missing_context_errors_via_producer() {
let (producer, body) = context_body(
std::path::PathBuf::from("/nonexistent/podup/context"),
ContextSource::Dockerfile("Dockerfile".to_string()),
Vec::new(),
);
futures_util::pin_mut!(body);
while body.next().await.is_some() {}
let produced = producer.await.expect("join");
assert!(produced.is_err(), "walking a missing context must error");
}
}