use async_trait::async_trait;
use clap::{CommandFactory, Parser};
use std::path::Path;
use crate::interpreter::ExecResult;
use crate::operation::KernelOperation;
use crate::tools::{schema_from_clap, ExecContext, ToolCtx, GlobalFlags, Tool, ToolArgs, ToolSchema};
pub struct Tee;
#[derive(Parser, Debug)]
#[command(name = "tee", about = "Read from stdin and write to stdout and files")]
struct TeeArgs {
#[arg(id = "append", short = 'a', long = "append")]
_append: bool,
#[command(flatten)]
global: GlobalFlags,
paths: Vec<String>,
}
#[async_trait]
impl Tool for Tee {
fn name(&self) -> &str {
"tee"
}
fn schema(&self) -> ToolSchema {
schema_from_clap(
&TeeArgs::command(),
"tee",
"Read from stdin and write to stdout and files",
[
("Save and display", "echo hello | tee output.txt"),
("Append to log", "echo entry | tee -a log.txt"),
],
)
.with_operations([KernelOperation::FsOverwrite.as_str()])
}
async fn execute(&self, args: ToolArgs, ctx: &mut dyn ToolCtx) -> ExecResult {
let Some(ctx) = ctx.as_any_mut().downcast_mut::<ExecContext>() else {
return ExecResult::failure(1, "internal error: kernel builtin requires ExecContext");
};
let argv = match args.to_argv() {
Ok(v) => v,
Err(e) => return ExecResult::failure(2, format!("tee: {e}")),
};
let parsed = match TeeArgs::try_parse_from(
std::iter::once("tee".to_string()).chain(argv),
) {
Ok(p) => p,
Err(e) => return ExecResult::failure(2, format!("tee: {e}")),
};
parsed.global.apply(ctx);
if args.positional.is_empty() {
return ExecResult::failure(1, "tee: missing file argument");
}
let append = args.has_flag("append") || args.has_flag("a");
let paths = match crate::interpreter::values_to_text_sink_named(&args.positional, "a path") {
Ok(p) => p,
Err(e) => return ExecResult::failure(1, format!("tee: {e}")),
};
let targets: Vec<(String, bool)> = paths.iter().map(|p| (p.clone(), append)).collect();
let snapshots = match ctx
.snapshot_overwrites("tee",
&targets)
.await
{
Ok(s) => s,
Err(blocked) => return blocked,
};
let input = match ctx.read_stdin_to_bytes().await {
Ok(i) => i.unwrap_or_default(),
Err(e) => return ExecResult::failure(1, format!("tee: {e}")),
};
let mut errors: Vec<String> = Vec::new();
for path_str in &paths {
let resolved = ctx.resolve_path(path_str);
let path = Path::new(&resolved);
let write_result = if append {
ctx.backend
.append(path, &input)
.await
.map_err(|e| e.to_string())
} else {
let expected = snapshots.get(&resolved);
ctx.overwrite_checked(path, &input, expected).await
};
if let Err(e) = write_result {
errors.push(format!("tee: {}: {}", path_str, e));
}
}
let mut result = ExecResult::success_text_or_bytes(input);
if !errors.is_empty() {
result.err = ExecResult::terminate_diagnostic(errors.join("\n"));
result = result.with_code(1);
}
result
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::backend::WriteMode;
use crate::ast::Value;
use crate::vfs::{Filesystem, MemoryFs, VfsRouter};
use std::sync::Arc;
async fn make_ctx() -> ExecContext {
let mut vfs = VfsRouter::new();
let mem = MemoryFs::new();
mem.write(Path::new("existing.txt"), b"original content\n")
.await
.unwrap();
vfs.mount("/", mem);
ExecContext::new(Arc::new(vfs))
}
#[tokio::test]
async fn test_tee_new_file() {
let mut ctx = make_ctx().await;
ctx.set_stdin("hello world\n".to_string());
let mut args = ToolArgs::new();
args.positional.push(Value::String("/output.txt".into()));
let result = Tee.execute(args, &mut ctx).await;
assert!(result.ok());
assert_eq!(&*result.text_out(), "hello world\n");
let written = ctx
.backend
.read(Path::new("/output.txt"), None)
.await
.unwrap();
assert_eq!(written, b"hello world\n");
}
#[tokio::test]
async fn test_tee_overwrite() {
let mut ctx = make_ctx().await;
ctx.set_stdin("new content\n".to_string());
let mut args = ToolArgs::new();
args.positional.push(Value::String("/existing.txt".into()));
let result = Tee.execute(args, &mut ctx).await;
assert!(result.ok());
let written = ctx
.backend
.read(Path::new("/existing.txt"), None)
.await
.unwrap();
assert_eq!(written, b"new content\n");
}
#[tokio::test]
async fn test_tee_append() {
let mut ctx = make_ctx().await;
ctx.set_stdin("appended\n".to_string());
let mut args = ToolArgs::new();
args.positional.push(Value::String("/existing.txt".into()));
args.flags.insert("a".to_string());
let result = Tee.execute(args, &mut ctx).await;
assert!(result.ok());
let written = ctx
.backend
.read(Path::new("/existing.txt"), None)
.await
.unwrap();
assert_eq!(written, b"original content\nappended\n");
}
struct FailingReadBackend {
inner: Arc<dyn crate::backend::KernelBackend>,
fail_path: std::path::PathBuf,
}
#[async_trait::async_trait]
impl crate::backend::KernelBackend for FailingReadBackend {
async fn read(
&self,
path: &Path,
range: Option<crate::backend::ReadRange>,
) -> crate::backend::BackendResult<Vec<u8>> {
if path == self.fail_path {
return Err(crate::backend::BackendError::PermissionDenied(
path.display().to_string(),
));
}
self.inner.read(path, range).await
}
async fn write(
&self,
path: &Path,
content: &[u8],
mode: WriteMode,
) -> crate::backend::BackendResult<()> {
self.inner.write(path, content, mode).await
}
async fn append(&self, path: &Path, content: &[u8]) -> crate::backend::BackendResult<()> {
if path == self.fail_path {
return Err(crate::backend::BackendError::PermissionDenied(
path.display().to_string(),
));
}
self.inner.append(path, content).await
}
async fn patch(
&self,
path: &Path,
ops: &[crate::backend::PatchOp],
) -> crate::backend::BackendResult<()> {
self.inner.patch(path, ops).await
}
async fn list(&self, path: &Path) -> crate::backend::BackendResult<Vec<crate::vfs::DirEntry>> {
self.inner.list(path).await
}
async fn stat(&self, path: &Path) -> crate::backend::BackendResult<crate::vfs::DirEntry> {
self.inner.stat(path).await
}
async fn mkdir(&self, path: &Path) -> crate::backend::BackendResult<()> {
self.inner.mkdir(path).await
}
async fn set_mtime(
&self,
path: &Path,
mtime: std::time::SystemTime,
) -> crate::backend::BackendResult<()> {
self.inner.set_mtime(path, mtime).await
}
async fn remove(&self, path: &Path, recursive: bool) -> crate::backend::BackendResult<()> {
self.inner.remove(path, recursive).await
}
async fn rename(&self, from: &Path, to: &Path) -> crate::backend::BackendResult<()> {
self.inner.rename(from, to).await
}
async fn exists(&self, path: &Path) -> bool {
self.inner.exists(path).await
}
async fn lstat(&self, path: &Path) -> crate::backend::BackendResult<crate::vfs::DirEntry> {
self.inner.lstat(path).await
}
async fn read_link(&self, path: &Path) -> crate::backend::BackendResult<std::path::PathBuf> {
self.inner.read_link(path).await
}
async fn symlink(&self, target: &Path, link: &Path) -> crate::backend::BackendResult<()> {
self.inner.symlink(target, link).await
}
async fn call_tool(
&self,
name: &str,
args: ToolArgs,
ctx: &mut dyn ToolCtx,
) -> crate::backend::BackendResult<crate::backend::ToolResult> {
self.inner.call_tool(name, args, ctx).await
}
async fn list_tools(&self) -> crate::backend::BackendResult<Vec<crate::backend::ToolInfo>> {
self.inner.list_tools().await
}
async fn get_tool(
&self,
name: &str,
) -> crate::backend::BackendResult<Option<crate::backend::ToolInfo>> {
self.inner.get_tool(name).await
}
fn read_only(&self) -> bool {
self.inner.read_only()
}
fn backend_type(&self) -> &str {
self.inner.backend_type()
}
fn mounts(&self) -> Vec<crate::backend::MountInfo> {
self.inner.mounts()
}
fn resolve_real_path(&self, path: &Path) -> Option<std::path::PathBuf> {
self.inner.resolve_real_path(path)
}
}
#[tokio::test]
async fn test_tee_append_read_failure_does_not_truncate() {
use crate::backend::KernelBackend as _;
let mut vfs = VfsRouter::new();
let mem = MemoryFs::new();
mem.write(Path::new("existing.txt"), b"original content\n")
.await
.unwrap();
vfs.mount("/", mem);
let vfs = Arc::new(vfs);
let inner: Arc<dyn crate::backend::KernelBackend> =
Arc::new(crate::backend::LocalBackend::new(vfs.clone()));
let verify_backend = crate::backend::LocalBackend::new(vfs);
let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(FailingReadBackend {
inner,
fail_path: std::path::PathBuf::from("/existing.txt"),
});
let mut ctx = ExecContext::with_backend(backend);
ctx.set_stdin("appended\n".to_string());
let mut args = ToolArgs::new();
args.positional.push(Value::String("/existing.txt".into()));
args.flags.insert("a".to_string());
let result = Tee.execute(args, &mut ctx).await;
assert!(!result.ok(), "tee -a must fail when the pre-append read fails");
assert!(
!result.err.is_empty(),
"the read failure must be reported, not swallowed"
);
let written = verify_backend
.read(Path::new("/existing.txt"), None)
.await
.unwrap();
assert_eq!(
written, b"original content\n",
"a failed pre-append read must not truncate the file; got {written:?}"
);
}
#[tokio::test]
async fn test_tee_empty_stdin() {
let mut ctx = make_ctx().await;
ctx.set_stdin("".to_string());
let mut args = ToolArgs::new();
args.positional.push(Value::String("/empty.txt".into()));
let result = Tee.execute(args, &mut ctx).await;
assert!(result.ok());
assert_eq!(&*result.text_out(), "");
let written = ctx
.backend
.read(Path::new("/empty.txt"), None)
.await
.unwrap();
assert!(written.is_empty());
}
#[tokio::test]
async fn test_tee_missing_file() {
let mut ctx = make_ctx().await;
ctx.set_stdin("data\n".to_string());
let result = Tee.execute(ToolArgs::new(), &mut ctx).await;
assert!(!result.ok());
}
#[tokio::test]
async fn overwrite_checked_rejects_concurrent_change() {
let ctx = make_ctx().await; let path = Path::new("/existing.txt");
let snapshot = crate::tools::OverwriteExpectation::Bytes(b"original content\n".to_vec());
ctx.backend
.write(path, b"changed elsewhere\n", WriteMode::Overwrite)
.await
.unwrap();
let result = ctx.overwrite_checked(path, b"my content\n", Some(&snapshot)).await;
assert!(result.is_err(), "expected a conflict, got {result:?}");
let now = ctx.backend.read(path, None).await.unwrap();
assert_eq!(now, b"changed elsewhere\n");
}
#[tokio::test]
async fn overwrite_checked_writes_when_snapshot_matches() {
let ctx = make_ctx().await;
let path = Path::new("/existing.txt");
let snapshot = crate::tools::OverwriteExpectation::Bytes(b"original content\n".to_vec());
ctx.overwrite_checked(path, b"new content\n", Some(&snapshot))
.await
.unwrap();
assert_eq!(
ctx.backend.read(path, None).await.unwrap(),
b"new content\n"
);
}
#[tokio::test]
async fn overwrite_checked_skips_cas_without_expectation() {
let ctx = make_ctx().await;
let path = Path::new("/existing.txt");
ctx.overwrite_checked(path, b"forced\n", None).await.unwrap();
assert_eq!(ctx.backend.read(path, None).await.unwrap(), b"forced\n");
}
#[tokio::test]
async fn overwrite_checked_errors_when_reread_fails_even_for_empty_snapshot() {
let ctx = make_ctx().await;
let path = Path::new("/empty.txt");
ctx.backend.write(path, b"", WriteMode::Overwrite).await.unwrap();
let empty_snapshot = crate::tools::OverwriteExpectation::Bytes(Vec::new());
ctx.backend.remove(path, false).await.unwrap();
let result = ctx
.overwrite_checked(path, b"new\n", Some(&empty_snapshot))
.await;
assert!(result.is_err(), "vanished target must error, got {result:?}");
}
}