use crate::error::*;
use crate::opts::DockerEngineClientOption;
use bytes::Bytes;
use futures::ready;
use futures::task::{Context, Poll};
use hyper::body::HttpBody;
use hyper::client::connect::Connect;
use hyper::{Body, Client, Response, Uri};
use hyperlocal::UnixConnector;
use lazy_static::lazy_static;
use percent_encoding::{utf8_percent_encode, NON_ALPHANUMERIC};
use regex::Regex;
use snafu::ResultExt;
use std::collections::HashMap;
use std::io;
use std::sync::Mutex;
use std::time::Duration;
use tokio::io::AsyncRead;
use tokio::macros::support::Pin;
use tokio::time::timeout;
pub mod container;
pub mod error;
pub mod network;
pub mod opts;
pub mod types;
pub mod version;
pub mod volume;
lazy_static! {
static ref HEADER_REGEXP: Regex = Regex::new(r"\ADocker/.+\s\((.+)\)\z").unwrap();
}
pub type LocalDockerEngineClient = DockerEngineClient<UnixConnector>;
#[derive(Clone, Debug)]
pub struct DockerEngineClient<C: Connect + Clone + Send + Sync + 'static> {
pub(crate) scheme: Option<String>,
pub(crate) host: Option<String>,
pub(crate) proto: Option<String>,
pub(crate) addr: Option<String>,
pub(crate) base_path: Option<String>,
pub(crate) timeout: Duration,
pub(crate) client: Option<Client<C, Body>>,
pub(crate) version: String,
pub(crate) custom_http_headers: Option<HashMap<String, String>>,
pub(crate) manual_override: bool,
pub(crate) negotiate_version: bool,
pub(crate) negotiated: bool,
}
impl<C: Connect + Clone + Send + Sync + 'static> Default for DockerEngineClient<C> {
fn default() -> Self {
Self {
scheme: Some("http".into()),
host: Some("unix:///var/run/docker.sock".into()),
proto: Some("unix".into()),
addr: Some("/var/run/docker.sock".into()),
base_path: None,
timeout: Duration::from_millis(5_000),
client: None,
version: "1.40".into(),
custom_http_headers: None,
manual_override: false,
negotiate_version: false,
negotiated: false,
}
}
}
impl<C: Connect + Clone + Send + Sync + 'static> DockerEngineClient<C> {
pub fn new_client_with_opts(
options: Option<Vec<DockerEngineClientOption<C>>>,
) -> Result<DockerEngineClient<C>, Error> {
let mut client: DockerEngineClient<C> = Default::default();
if let Some(client_options) = options {
for client_option in &client_options {
client_option(&mut client)?;
}
}
Ok(client)
}
fn request_uri(
&self,
path: &str,
query_params: Option<HashMap<String, String>>,
) -> Result<Uri, Error> {
let encoded_query_params = query_params
.map(|query_params| {
let mut encoded_query_params = String::from("?");
for (key, value) in query_params {
encoded_query_params.push_str(&key);
encoded_query_params.push('=');
encoded_query_params
.push_str(&utf8_percent_encode(&value, NON_ALPHANUMERIC).to_string());
encoded_query_params.push('&');
}
encoded_query_params.trim_end_matches('&').to_string()
})
.unwrap_or_else(String::new);
let path_and_query_params = format!(
"{}{}{}",
self.base_path
.as_ref()
.map(String::from)
.unwrap_or_else(|| format!("/v{}", self.version)),
path,
encoded_query_params
);
Ok(match self.proto.as_ref().unwrap().as_str() {
"tcp" => {
let scheme = self
.scheme
.as_ref()
.map(String::from)
.unwrap_or_else(|| "http".to_string());
let address = self
.addr
.as_ref()
.map(String::from)
.unwrap_or_else(|| "localhost".to_string());
format!("{}://{}{}", scheme, address, path_and_query_params)
.parse()
.unwrap()
}
"unix" => {
hyperlocal::Uri::new(self.addr.as_ref().unwrap(), &path_and_query_params).into()
}
_ => unimplemented!(),
})
}
}
pub(crate) async fn read_response_body(
mut response: Response<Body>,
read_timeout: Duration,
) -> Result<String, Error> {
let mut response_body = String::new();
while let Some(chunk) = timeout(read_timeout, response.body_mut().data())
.await
.context(HttpClientTimeoutError {})?
{
response_body.push_str(&String::from_utf8_lossy(
&chunk.context(HttpClientError {})?,
));
}
Ok(response_body)
}
pub(crate) async fn read_response_body_raw(
mut response: Response<Body>,
read_timeout: Duration,
) -> Result<Vec<u8>, Error> {
let mut response_body = Vec::new();
while let Some(chunk) = timeout(read_timeout, response.body_mut().data())
.await
.context(HttpClientTimeoutError {})?
{
response_body.append(&mut chunk.context(HttpClientError {})?.to_vec());
}
Ok(response_body)
}
pub(crate) fn parse_host_url(host: &str) -> Result<Uri, Error> {
host.parse::<Uri>().context(HttpUriError {})
}
pub(crate) fn get_docker_os(server_header: &str) -> String {
let captures = HEADER_REGEXP.captures(server_header).unwrap();
captures
.get(1)
.map(|m| String::from(m.as_str()))
.unwrap_or_else(String::new)
}
pub(crate) fn get_filters_query(
filters: types::filters::Args,
) -> Result<HashMap<String, String>, Error> {
let mut query_params: HashMap<String, String> = HashMap::new();
if !filters.fields.is_empty() {
query_params.insert(
"filters".into(),
serde_json::to_string(&filters.fields).context(JsonSerializationError {})?,
);
}
Ok(query_params)
}
pub(crate) struct AsyncHttpBodyReader {
body: Mutex<Pin<Box<dyn HttpBody<Data = Bytes, Error = hyper::Error>>>>,
}
impl AsyncHttpBodyReader {
pub(crate) fn new(response: Response<Body>) -> Self {
Self {
body: Mutex::new(Box::pin(response.into_body())),
}
}
}
impl AsyncRead for AsyncHttpBodyReader {
fn poll_read(
self: Pin<&mut Self>,
cx: &mut Context<'_>,
buf: &mut [u8],
) -> Poll<Result<usize, std::io::Error>> {
let mut body_reader = self.body.lock().unwrap();
if let Some(data) = ready!(body_reader.as_mut().poll_data(cx)) {
match data {
Ok(data) => Poll::Ready(io::Read::read(
&mut Vec::from(data.as_ref()).as_mut_slice().as_ref(),
buf,
)),
Err(err) => Poll::Ready(Err(io::Error::new(
io::ErrorKind::Other,
format!("{}", err),
))),
}
} else {
Poll::Ready(Ok(0))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
pub fn test_get_docker_os() {
let os = get_docker_os("Docker/19.03.5 (linux)");
assert_eq!(&os, "linux");
}
}