pub mod downloader;
pub mod finder;
pub mod platform;
use std::io::{BufReader, Read, Write};
use std::path::PathBuf;
use std::process::{Child, Command, ExitStatus, Stdio};
use std::thread;
use std::time::Duration;
pub struct P4Cli {
bin_path: PathBuf,
_cache: Option<tempfile::TempDir>,
}
impl P4Cli {
pub fn new() -> std::io::Result<Self> {
if let Some(path) = finder::find_system_p4() {
return Ok(Self {
bin_path: path,
_cache: None,
});
}
let (bin_path, cache) = downloader::download_p4()?;
Ok(Self {
bin_path,
_cache: Some(cache),
})
}
pub fn run<S: AsRef<std::ffi::OsStr>>(&self, args: &[S]) -> std::io::Result<P4Output> {
self.command().args(args).run()
}
pub fn stream<S: AsRef<std::ffi::OsStr>>(&self, args: &[S]) -> std::io::Result<P4Stream> {
self.command().args(args).stream()
}
pub fn command(&self) -> P4Command<'_> {
P4Command::new(self)
}
}
pub struct P4Output {
exit_code: i32,
timed_out: bool,
stdout: Vec<u8>,
stderr: Vec<u8>,
}
impl P4Output {
pub fn exit_code(&self) -> i32 {
self.exit_code
}
pub fn timed_out(&self) -> bool {
self.timed_out
}
pub fn success(&self) -> bool {
!self.timed_out && self.exit_code == 0
}
pub fn stdout(&self) -> &[u8] {
&self.stdout
}
pub fn stderr(&self) -> &[u8] {
&self.stderr
}
pub fn stdout_str(&self) -> std::io::Result<&str> {
std::str::from_utf8(&self.stdout).map_err(std::io::Error::other)
}
pub fn stderr_str(&self) -> std::io::Result<&str> {
std::str::from_utf8(&self.stderr).map_err(std::io::Error::other)
}
pub fn stdout_lines(&self) -> std::io::Result<Vec<&str>> {
let s = self.stdout_str()?;
if s.is_empty() {
Ok(Vec::new())
} else {
Ok(s.lines().collect())
}
}
pub fn stderr_lines(&self) -> std::io::Result<Vec<&str>> {
let s = self.stderr_str()?;
if s.is_empty() {
Ok(Vec::new())
} else {
Ok(s.lines().collect())
}
}
}
impl std::fmt::Debug for P4Output {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("P4Output")
.field("exit_code", &self.exit_code)
.field("timed_out", &self.timed_out)
.field("stdout_len", &self.stdout.len())
.field("stderr_len", &self.stderr.len())
.finish()
}
}
pub enum P4StreamEvent {
Stdout(Vec<u8>),
Stderr(Vec<u8>),
Exit(i32),
}
impl P4StreamEvent {
pub fn as_utf8(&self) -> Option<&str> {
match self {
P4StreamEvent::Stdout(data) | P4StreamEvent::Stderr(data) => {
std::str::from_utf8(data).ok()
}
P4StreamEvent::Exit(_) => None,
}
}
}
impl std::fmt::Display for P4StreamEvent {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
P4StreamEvent::Stdout(data) | P4StreamEvent::Stderr(data) => {
if let Ok(text) = std::str::from_utf8(data) {
write!(f, "{text}")
} else {
write!(f, "<{} bytes>", data.len())
}
}
P4StreamEvent::Exit(code) => write!(f, "(exit {code})"),
}
}
}
pub struct P4Stream {
rx: std::sync::mpsc::Receiver<std::io::Result<P4StreamEvent>>,
child: Option<Child>,
exhausted: bool,
}
impl std::fmt::Debug for P4Stream {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("P4Stream")
.field("exhausted", &self.exhausted)
.finish()
}
}
impl Iterator for P4Stream {
type Item = std::io::Result<P4StreamEvent>;
fn next(&mut self) -> Option<Self::Item> {
if self.exhausted {
return None;
}
match self.rx.recv() {
Ok(item) => Some(item),
Err(_) => {
self.exhausted = true;
match self.child.take() {
Some(mut c) => match c.wait() {
Ok(status) => Some(Ok(P4StreamEvent::Exit(status.code().unwrap_or(-1)))),
Err(e) => Some(Err(e)),
},
None => None,
}
}
}
}
}
impl Drop for P4Stream {
fn drop(&mut self) {
if let Some(ref mut child) = self.child {
let _ = child.kill();
let _ = child.wait();
}
}
}
pub struct P4Command<'a> {
cli: &'a P4Cli,
args: Vec<std::ffi::OsString>,
timeout: Option<Duration>,
cwd: Option<PathBuf>,
envs: Vec<(std::ffi::OsString, std::ffi::OsString)>,
stdin_data: Option<Vec<u8>>,
}
impl<'a> P4Command<'a> {
fn new(cli: &'a P4Cli) -> Self {
Self {
cli,
args: Vec::new(),
timeout: None,
cwd: None,
envs: Vec::new(),
stdin_data: None,
}
}
pub fn arg(&mut self, arg: impl AsRef<std::ffi::OsStr>) -> &mut Self {
self.args.push(arg.as_ref().to_os_string());
self
}
pub fn args(&mut self, args: &[impl AsRef<std::ffi::OsStr>]) -> &mut Self {
self.args
.extend(args.iter().map(|a| a.as_ref().to_os_string()));
self
}
pub fn timeout(&mut self, timeout: Duration) -> &mut Self {
self.timeout = Some(timeout);
self
}
pub fn cwd(&mut self, path: impl Into<PathBuf>) -> &mut Self {
self.cwd = Some(path.into());
self
}
pub fn env(
&mut self,
key: impl Into<std::ffi::OsString>,
val: impl Into<std::ffi::OsString>,
) -> &mut Self {
self.envs.push((key.into(), val.into()));
self
}
pub fn stdin(&mut self, data: impl Into<Vec<u8>>) -> &mut Self {
self.stdin_data = Some(data.into());
self
}
pub fn run(&mut self) -> std::io::Result<P4Output> {
let mut cmd = Command::new(&self.cli.bin_path);
cmd.args(&self.args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
if self.stdin_data.is_some() {
cmd.stdin(Stdio::piped());
} else {
cmd.stdin(Stdio::null());
}
if let Some(ref cwd) = self.cwd {
cmd.current_dir(cwd);
}
for (k, v) in &self.envs {
cmd.env(k, v);
}
let mut child = cmd.spawn()?;
let stdin_handle = self.stdin_data.take().and_then(|data| {
child.stdin.take().map(|mut stdin| {
thread::spawn(move || match stdin.write_all(&data) {
Err(e) => Err(e),
Ok(_) => Ok(()),
})
})
});
let stdout = child
.stdout
.take()
.ok_or_else(|| std::io::Error::other("stdout was not captured"))?;
let stderr = child
.stderr
.take()
.ok_or_else(|| std::io::Error::other("stderr was not captured"))?;
let stdout_handle = thread::spawn(move || {
let mut buf = Vec::new();
BufReader::new(stdout).read_to_end(&mut buf)?;
Ok::<_, std::io::Error>(buf)
});
let stderr_handle = thread::spawn(move || {
let mut buf = Vec::new();
BufReader::new(stderr).read_to_end(&mut buf)?;
Ok::<_, std::io::Error>(buf)
});
let (exit_status, timed_out) = wait_process(&mut child, self.timeout)?;
if let Some(handle) = stdin_handle {
handle
.join()
.map_err(|_| std::io::Error::other("stdin thread panicked"))?
.map_err(|e| std::io::Error::other(format!("stdin write failed: {e}")))?;
}
let stdout_buf = stdout_handle
.join()
.map_err(|_| std::io::Error::other("stdout reader thread panicked"))?
.map_err(|e| std::io::Error::other(format!("stdout read failed: {e}")))?;
let stderr_buf = stderr_handle
.join()
.map_err(|_| std::io::Error::other("stderr reader thread panicked"))?
.map_err(|e| std::io::Error::other(format!("stderr read failed: {e}")))?;
Ok(P4Output {
exit_code: exit_status.code().unwrap_or(-1),
timed_out,
stdout: stdout_buf,
stderr: stderr_buf,
})
}
pub fn stream(&mut self) -> std::io::Result<P4Stream> {
if self.timeout.is_some() {
return Err(std::io::Error::other(
"timeout is not supported on stream(); use run() instead",
));
}
let mut cmd = Command::new(&self.cli.bin_path);
cmd.args(&self.args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
if self.stdin_data.is_some() {
cmd.stdin(Stdio::piped());
} else {
cmd.stdin(Stdio::null());
}
if let Some(ref cwd) = self.cwd {
cmd.current_dir(cwd);
}
for (k, v) in &self.envs {
cmd.env(k, v);
}
let mut child = cmd.spawn()?;
if let Some(data) = self.stdin_data.take()
&& let Some(mut stdin) = child.stdin.take()
{
thread::spawn(move || {
let _ = stdin.write_all(&data);
});
}
let mut stdout = child
.stdout
.take()
.ok_or_else(|| std::io::Error::other("stdout was not captured"))?;
let mut stderr = child
.stderr
.take()
.ok_or_else(|| std::io::Error::other("stderr was not captured"))?;
let (tx, rx) = std::sync::mpsc::sync_channel(64);
let tx_out = tx.clone();
thread::spawn(move || {
let mut buf = vec![0u8; 65536];
loop {
let n = match stdout.read(&mut buf) {
Ok(0) => break,
Ok(n) => n,
Err(e) => {
let _ = tx_out.send(Err(e));
break;
}
};
if tx_out
.send(Ok(P4StreamEvent::Stdout(buf[..n].to_vec())))
.is_err()
{
break;
}
}
});
let tx_err = tx.clone();
thread::spawn(move || {
let mut buf = vec![0u8; 65536];
loop {
let n = match stderr.read(&mut buf) {
Ok(0) => break,
Ok(n) => n,
Err(e) => {
let _ = tx_err.send(Err(e));
break;
}
};
if tx_err
.send(Ok(P4StreamEvent::Stderr(buf[..n].to_vec())))
.is_err()
{
break;
}
}
});
Ok(P4Stream {
rx,
child: Some(child),
exhausted: false,
})
}
}
fn wait_process(
child: &mut Child,
timeout: Option<Duration>,
) -> std::io::Result<(ExitStatus, bool)> {
match timeout {
None => Ok((child.wait()?, false)),
Some(t) => wait_with_timeout(child, t),
}
}
fn wait_with_timeout(child: &mut Child, timeout: Duration) -> std::io::Result<(ExitStatus, bool)> {
let start = std::time::Instant::now();
loop {
if let Some(status) = child.try_wait()? {
return Ok((status, false));
}
if start.elapsed() >= timeout {
child.kill()?;
return Ok((child.wait()?, true));
}
thread::sleep(Duration::from_millis(50));
}
}