//! Scope-aware filesystem primitive. //! //! `LocalWorkdirSession` is the write/read gate layered on top of a [`manifest::Scope`] //! and a Worker's working directory. The scope decides which paths are //! readable and writable; the cwd is carried alongside for convenience //! (Glob/Grep default their search base to it). //! //! `LocalWorkdirSession` is cheap to clone (`Arc` inside). Tool-specific session //! state, such as read-before-edit tracking, remains owned by the tool layer. use std::collections::HashMap; #[cfg(test)] use std::io::Write as _; use std::io::{Read as _, Seek as _, SeekFrom}; use std::path::{Path, PathBuf}; use std::process::Stdio; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::Duration; use async_trait::async_trait; use manifest::{Scope, SharedScope}; use sha2::{Digest, Sha256}; use tokio::process::Command; use tokio::sync::{Mutex, Notify}; use tokio::task::JoinHandle; use crate::{ CommandHandle, CommandOutput, CommandOutputRequest, CommandRequest, CommandStatus, EditRequest, EditResult, GlobRequest, GlobResult, GrepRequest, GrepResult, ListRequest, ListResult, ReadRequest, ReadResult, StatRequest, StatResult, Workdir, WorkdirError, WorkdirPath, WorkdirSession, WorkdirSessionCapabilities, WorkdirSessionCapability, WriteRequest, WriteResult, }; #[cfg(test)] use crate::{EntryKind, WriteOutcome}; #[derive(Debug)] enum LocalCommand { Running { task: JoinHandle>, completion: Arc, }, Completed(CommandOutput), } #[derive(Debug)] struct ScopeAccess(Arc); impl fs_operation::FsAccessPolicy for ScopeAccess { fn is_readable(&self, path: &Path) -> bool { self.0.is_readable(path) } fn is_writable(&self, path: &Path) -> bool { self.0.is_writable(path) } } #[derive(Debug)] struct LocalWorkdirSessionInner { workdir: Workdir, root: PathBuf, scope: SharedScope, cwd: PathBuf, capabilities: WorkdirSessionCapabilities, closed: AtomicBool, close_lock: Mutex<()>, next_command_id: AtomicU64, commands: Mutex>, } impl Drop for LocalWorkdirSessionInner { fn drop(&mut self) { if let Ok(mut commands) = self.commands.try_lock() { for (_, command) in commands.drain() { if let LocalCommand::Running { task, .. } = command { task.abort(); } } } } } /// Scope-aware filesystem handle. Clone-cheap (`Arc` inside). /// /// The wrapped [`SharedScope`] is shared with every clone of this /// `LocalWorkdirSession` and with whoever else holds the same `SharedScope` /// handle (typically the owning Worker). Mutations to that `SharedScope` /// propagate atomically; the next permission check inside any /// `LocalWorkdirSession` reads the new view. #[derive(Debug, Clone)] pub struct LocalWorkdirSession { inner: Arc, } /// First symlink encountered while resolving a path. #[derive(Debug, Clone, PartialEq, Eq)] pub struct SymlinkInfo { /// The symlink path as it appears in the original path chain. pub link_path: PathBuf, /// The symlink target resolved relative to the symlink's parent when the /// link stores a relative target. pub target_path: PathBuf, /// Best-effort resolved form of the full requested path after replacing /// the symlink component with its target and rejoining any remaining tail. /// Existing targets are canonicalized; broken targets are left absolute. pub resolved_path: PathBuf, /// Whether the symlink target itself exists. A missing target is a broken /// symlink even when the symlink lives inside an allowed scope. pub target_exists: bool, } fn local_workdir_identity(root: &Path) -> Workdir { let canonical = root.canonicalize().unwrap_or_else(|_| root.to_path_buf()); let digest = Sha256::digest(canonical.to_string_lossy().as_bytes()); let digest = digest .iter() .map(|byte| format!("{byte:02x}")) .collect::(); Workdir::new(format!("local-{digest}")) } impl LocalWorkdirSession { /// Create a new [`LocalWorkdirSession`] wrapping `scope` and `cwd` in a fresh /// [`SharedScope`]. Use [`LocalWorkdirSession::with_shared_scope`] when you /// need the resulting `LocalWorkdirSession` to share scope state with another /// holder of the `SharedScope` (typically the Worker). pub fn new(scope: Scope, cwd: PathBuf) -> Self { Self::materialized( cwd.clone(), cwd, SharedScope::new(scope), WorkdirSessionCapabilities::ALL, ) } pub fn with_shared_scope(scope: SharedScope, cwd: PathBuf) -> Self { Self::materialized(cwd.clone(), cwd, scope, WorkdirSessionCapabilities::ALL) } /// Construct a standalone local session with a deterministic identity /// derived from the canonical materialization root. pub fn materialized( root: PathBuf, cwd: PathBuf, scope: SharedScope, capabilities: WorkdirSessionCapabilities, ) -> Self { let workdir = local_workdir_identity(&root); Self::materialized_bound(workdir, root, cwd, scope, capabilities) } /// Open a local session for an authority-assigned persistent Workdir. pub fn materialized_bound( workdir: Workdir, root: PathBuf, cwd: PathBuf, scope: SharedScope, capabilities: WorkdirSessionCapabilities, ) -> Self { Self { inner: Arc::new(LocalWorkdirSessionInner { workdir, root, scope, cwd, capabilities, closed: AtomicBool::new(false), close_lock: Mutex::new(()), next_command_id: AtomicU64::new(1), commands: Mutex::new(HashMap::new()), }), } } pub fn root(&self) -> &Path { &self.inner.root } /// Snapshot the current scope. Cheap; the returned `Arc` is /// a coherent point-in-time view that subsequent mutations do not /// affect. pub fn scope(&self) -> Arc { self.inner.scope.snapshot() } /// Shared scope handle backing this `LocalWorkdirSession`. Cloning it lets a /// caller (usually the Worker) hold the same view and push updates /// that are immediately reflected in subsequent permission checks. pub fn shared_scope(&self) -> &SharedScope { &self.inner.scope } /// The Worker's working directory. Glob/Grep default their search base /// to this path when callers omit an explicit `path` parameter. pub fn cwd(&self) -> &Path { &self.inner.cwd } // ========================================================================= // Read — scope-checked against readability // ========================================================================= /// Read the full contents of `path` as raw bytes. /// /// Follows symlinks. Rejects directories, relative paths, paths not /// readable by the scope, and missing files. #[cfg(test)] pub(crate) fn read_bytes(&self, path: &Path) -> Result, WorkdirError> { if !path.is_absolute() { return Err(WorkdirError::RelativePath(path.to_path_buf())); } let symlink = first_symlink(path); let scope = self.inner.scope.load(); if !scope.is_readable(path) { return Err(symlink_out_of_scope_or_plain( path, symlink.as_ref(), "read", &scope, )); } if let Some(info) = symlink.as_ref() { if !info.target_exists { return Err(broken_symlink_error(path, info)); } } let meta = std::fs::metadata(path).map_err(|e| match e.kind() { std::io::ErrorKind::NotFound => WorkdirError::NotFound(path.to_path_buf()), _ => WorkdirError::io(path, e), })?; if meta.is_dir() { return Err(if let Some(info) = symlink.as_ref() { WorkdirError::SymlinkTargetIsDirectory { path: path.to_path_buf(), target: info.resolved_path.clone(), } } else { WorkdirError::IsDirectory(path.to_path_buf()) }); } std::fs::read(path).map_err(|e| WorkdirError::io(path, e)) } // ========================================================================= // Write — scope-checked, atomic // ========================================================================= /// Atomically write `content` to `path`, creating or overwriting it. /// /// - `path` must be absolute and writable under the scope. /// - Paths that are readable but not writable return [`WorkdirError::ReadOnly`]. /// - Paths outside the scope entirely return [`WorkdirError::OutOfScope`]. /// - Missing parent directories are created. /// - The actual write uses a sibling tempfile + `persist`, so the /// target file transitions atomically between states. /// /// This method does **not** consult tool-specific read history. #[cfg(test)] pub(crate) fn write(&self, path: &Path, content: &[u8]) -> Result { if !path.is_absolute() { return Err(WorkdirError::RelativePath(path.to_path_buf())); } let symlink = first_symlink(path); let scope = self.inner.scope.load(); if !scope.is_writable(path) { return Err(if scope.is_readable(path) { WorkdirError::ReadOnly(path.to_path_buf()) } else { symlink_out_of_scope_or_plain(path, symlink.as_ref(), "write", &scope) }); } drop(scope); if let Some(info) = symlink.as_ref() { if !info.target_exists { return Err(broken_symlink_error(path, info)); } } // Reject existing directory targets. match std::fs::metadata(path) { Ok(meta) if meta.is_dir() => { return Err(if let Some(info) = symlink.as_ref() { WorkdirError::SymlinkTargetIsDirectory { path: path.to_path_buf(), target: info.resolved_path.clone(), } } else { WorkdirError::IsDirectory(path.to_path_buf()) }); } _ => {} } let existed = path.exists(); let write_target = if existed { path.canonicalize().unwrap_or_else(|_| path.to_path_buf()) } else { path.to_path_buf() }; let parent = write_target.parent().ok_or_else(|| { WorkdirError::InvalidArgument(format!( "path has no parent directory: {}", write_target.display() )) })?; if !parent.as_os_str().is_empty() && !parent.exists() { std::fs::create_dir_all(parent).map_err(|e| WorkdirError::io(parent, e))?; } let tmp_parent: &Path = if parent.as_os_str().is_empty() { Path::new(".") } else { parent }; let mut tmp = tempfile::NamedTempFile::new_in(tmp_parent) .map_err(|e| WorkdirError::io(tmp_parent, e))?; tmp.write_all(content) .map_err(|e| WorkdirError::io(&write_target, e))?; tmp.as_file() .sync_all() .map_err(|e| WorkdirError::io(&write_target, e))?; tmp.persist(&write_target) .map_err(|e| WorkdirError::io(&write_target, e.error))?; Ok(WriteOutcome { bytes_written: content.len(), created: !existed, }) } fn ensure_open(&self) -> Result<(), WorkdirError> { if self.inner.closed.load(Ordering::Acquire) { Err(WorkdirError::Unavailable(format!( "Workdir session for {} is closed", self.inner.workdir.id() ))) } else { Ok(()) } } fn ensure_capability(&self, capability: WorkdirSessionCapability) -> Result<(), WorkdirError> { self.ensure_open()?; if self.inner.capabilities.supports(capability) { Ok(()) } else { Err(WorkdirError::Unsupported(capability)) } } fn resolve(&self, path: &WorkdirPath) -> PathBuf { if path.is_root() { self.inner.root.clone() } else { self.inner.root.join(path.as_str()) } } } #[async_trait] impl WorkdirSession for LocalWorkdirSession { fn workdir(&self) -> &Workdir { &self.inner.workdir } fn capabilities(&self) -> WorkdirSessionCapabilities { self.inner.capabilities } async fn stat(&self, request: StatRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Read)?; let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_stat(&self.inner.root, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn read(&self, request: ReadRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Read)?; let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_read(&self.inner.root, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn write(&self, request: WriteRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Write)?; let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_write(&self.inner.root, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn edit(&self, request: EditRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Edit)?; let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_edit(&self.inner.root, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn list(&self, request: ListRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Read)?; let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_list(&self.inner.root, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn glob(&self, request: GlobRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Glob)?; let logical = request.path.clone(); let base = self.resolve(&request.path); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_glob(&self.inner.root, &base, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn grep(&self, request: GrepRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Grep)?; let base = self.resolve(&request.path); let logical = request.path.clone(); let access = ScopeAccess(self.inner.scope.snapshot()); fs_operation::run_grep(&self.inner.root, base, request, &access) .map_err(WorkdirError::from) .map_err(|error| sanitize_error(error, &logical)) } async fn start_command(&self, request: CommandRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Command)?; let id = self.inner.next_command_id.fetch_add(1, Ordering::Relaxed); let handle = CommandHandle(format!("command-{id}")); let cwd = self.inner.cwd.clone(); let completion = Arc::new(Notify::new()); let task_completion = Arc::clone(&completion); let task = tokio::spawn(async move { let output = run_command(cwd, request).await; task_completion.notify_one(); output }); let mut commands = self.inner.commands.lock().await; if let Err(error) = self.ensure_open() { task.abort(); completion.notify_one(); return Err(error); } commands.insert(handle.0.clone(), LocalCommand::Running { task, completion }); Ok(handle) } async fn command_status(&self, handle: CommandHandle) -> Result { self.ensure_capability(WorkdirSessionCapability::Command)?; let commands = self.inner.commands.lock().await; let command = commands .get(&handle.0) .ok_or_else(|| WorkdirError::UnknownCommand(handle.0.clone()))?; Ok(match command { LocalCommand::Running { task, .. } if !task.is_finished() => CommandStatus::Running, LocalCommand::Running { .. } => CommandStatus::Completed, LocalCommand::Completed(output) => output.status, }) } async fn command_output( &self, request: CommandOutputRequest, ) -> Result { self.ensure_capability(WorkdirSessionCapability::Command)?; let command = loop { self.ensure_open()?; let mut commands = self.inner.commands.lock().await; let Some(command) = commands.get(&request.handle.0) else { return Err(WorkdirError::UnknownCommand(request.handle.0.clone())); }; let completion = match command { LocalCommand::Running { task, completion, .. } if !task.is_finished() => Some(Arc::clone(completion)), _ => None, }; if let Some(completion) = completion { if !request.wait { return Ok(CommandOutput { status: CommandStatus::Running, exit_code: None, timed_out: false, content: String::new(), next_cursor: None, truncated: false, }); } drop(commands); completion.notified().await; continue; } break commands .remove(&request.handle.0) .expect("command checked above"); }; let output = match command { LocalCommand::Running { task, .. } => task .await .map_err(|error| WorkdirError::Unavailable(error.to_string()))??, LocalCommand::Completed(output) => output, }; let page = command_output_page(&output, request.cursor, request.limit); if page.next_cursor.is_some() { let mut commands = self.inner.commands.lock().await; if !self.inner.closed.load(Ordering::Acquire) { commands.insert(request.handle.0, LocalCommand::Completed(output)); } } Ok(page) } async fn cancel_command(&self, handle: CommandHandle) -> Result<(), WorkdirError> { self.ensure_capability(WorkdirSessionCapability::Command)?; let command = self .inner .commands .lock() .await .remove(&handle.0) .ok_or_else(|| WorkdirError::UnknownCommand(handle.0))?; if let LocalCommand::Running { task, completion } = command { task.abort(); completion.notify_one(); } Ok(()) } async fn close(&self) -> Result<(), WorkdirError> { let _close_guard = self.inner.close_lock.lock().await; if self.inner.closed.swap(true, Ordering::AcqRel) { return Ok(()); } let mut commands = self.inner.commands.lock().await; for (_, command) in commands.drain() { if let LocalCommand::Running { task, completion } = command { task.abort(); completion.notify_one(); } } Ok(()) } } fn command_output_page(output: &CommandOutput, cursor: usize, limit: usize) -> CommandOutput { let total_chars = output.content.chars().count(); let start = cursor.min(total_chars); let content = output .content .chars() .skip(start) .take(limit.max(1)) .collect::(); let end = start + content.chars().count(); CommandOutput { status: output.status, exit_code: output.exit_code, timed_out: output.timed_out, content, next_cursor: (end < total_chars).then_some(end), truncated: output.truncated || end < total_chars, } } fn sanitize_error(error: WorkdirError, logical: &WorkdirPath) -> WorkdirError { let path = PathBuf::from(logical.as_str()); match error { WorkdirError::RelativePath(_) => WorkdirError::InvalidPath(logical.to_string()), WorkdirError::OutOfScope(_) => WorkdirError::OutOfScope(path), WorkdirError::SymlinkOutOfScope { required_permission, .. } => WorkdirError::SymlinkOutOfScope { path, target: PathBuf::from(""), required_permission, }, WorkdirError::BrokenSymlink { .. } => WorkdirError::BrokenSymlink { path: path.clone(), link: path, target: PathBuf::from(""), }, WorkdirError::SymlinkTargetIsDirectory { .. } => WorkdirError::SymlinkTargetIsDirectory { path, target: PathBuf::from(""), }, WorkdirError::SymlinkDirectoryNotTraversed { tool, .. } => { WorkdirError::SymlinkDirectoryNotTraversed { tool, path, target: PathBuf::from(""), } } WorkdirError::ReadOnly(_) => WorkdirError::ReadOnly(path), WorkdirError::IsDirectory(_) => WorkdirError::IsDirectory(path), WorkdirError::NotFound(_) => WorkdirError::NotFound(path), WorkdirError::Io { source, .. } => WorkdirError::Unavailable(format!( "I/O operation failed for {logical}: {}", source.kind() )), other => other, } } async fn run_command(cwd: PathBuf, request: CommandRequest) -> Result { let stdout = tempfile::NamedTempFile::new().map_err(|error| WorkdirError::io(&cwd, error))?; let stderr = tempfile::NamedTempFile::new().map_err(|error| WorkdirError::io(&cwd, error))?; let stdout_path = stdout.into_temp_path(); let stderr_path = stderr.into_temp_path(); let stdout_file = std::fs::File::create(&stdout_path) .map_err(|error| WorkdirError::io(&stdout_path, error))?; let stderr_file = std::fs::File::create(&stderr_path) .map_err(|error| WorkdirError::io(&stderr_path, error))?; let mut child = Command::new("bash") .arg("-c") .arg(&request.command) .current_dir(&cwd) .stdin(Stdio::null()) .stdout(Stdio::from(stdout_file)) .stderr(Stdio::from(stderr_file)) .kill_on_drop(true) .spawn() .map_err(|error| WorkdirError::io(&cwd, error))?; let timed_out = match tokio::time::timeout( Duration::from_secs(request.timeout_secs.max(1)), child.wait(), ) .await { Ok(result) => { let status = result.map_err(|error| WorkdirError::io(&cwd, error))?; let (content, truncated) = read_command_output_files(&stdout_path, &stderr_path, request.output_limit.max(1))?; return Ok(CommandOutput { status: CommandStatus::Completed, exit_code: status.code(), timed_out: false, content, next_cursor: None, truncated, }); } Err(_) => { let _ = child.kill().await; true } }; let (content, truncated) = read_command_output_files(&stdout_path, &stderr_path, request.output_limit.max(1))?; Ok(CommandOutput { status: CommandStatus::Failed, exit_code: None, timed_out, content, next_cursor: None, truncated, }) } fn read_command_output_files( stdout_path: &Path, stderr_path: &Path, limit: usize, ) -> Result<(String, bool), WorkdirError> { let stdout_size = std::fs::metadata(stdout_path) .map_err(|error| WorkdirError::io(stdout_path, error))? .len() as usize; let stderr_size = std::fs::metadata(stderr_path) .map_err(|error| WorkdirError::io(stderr_path, error))? .len() as usize; let total = stdout_size.saturating_add(stderr_size); let stderr_budget = stderr_size.min(limit / 2); let stdout_budget = stdout_size.min(limit.saturating_sub(stderr_budget)); let remaining = limit.saturating_sub(stdout_budget + stderr_budget); let stdout_budget = (stdout_budget + remaining.min(stdout_size - stdout_budget)).min(limit); let stderr_budget = limit.saturating_sub(stdout_budget).min(stderr_size); let stdout = read_tail(stdout_path, stdout_budget)?; let stderr = read_tail(stderr_path, stderr_budget)?; let mut content = String::new(); if !stdout.is_empty() { content.push_str(&String::from_utf8_lossy(&stdout)); } if !stderr.is_empty() { if !content.is_empty() && !content.ends_with('\n') { content.push('\n'); } content.push_str(&String::from_utf8_lossy(&stderr)); } Ok((content, total > limit)) } fn read_tail(path: &Path, limit: usize) -> Result, WorkdirError> { if limit == 0 { return Ok(Vec::new()); } let mut file = std::fs::File::open(path).map_err(|error| WorkdirError::io(path, error))?; let len = file .metadata() .map_err(|error| WorkdirError::io(path, error))? .len(); let start = len.saturating_sub(limit as u64); file.seek(SeekFrom::Start(start)) .map_err(|error| WorkdirError::io(path, error))?; let mut bytes = Vec::with_capacity((len - start) as usize); file.read_to_end(&mut bytes) .map_err(|error| WorkdirError::io(path, error))?; Ok(bytes) } /// Return the first symlink component in `path`, if one exists. /// /// The function only inspects existing path components. It intentionally uses /// `symlink_metadata` so the symlink itself can be diagnosed before any later /// `metadata` call follows it and collapses the reason into `NotFound` or /// `OutOfScope`. pub fn first_symlink(path: &Path) -> Option { if !path.is_absolute() { return None; } let mut cur = PathBuf::new(); let mut components = path.components().peekable(); while let Some(component) = components.next() { cur.push(component.as_os_str()); let meta = std::fs::symlink_metadata(&cur).ok()?; if !meta.file_type().is_symlink() { continue; } let raw_target = std::fs::read_link(&cur).ok()?; let target_path = if raw_target.is_absolute() { raw_target } else { cur.parent() .unwrap_or_else(|| Path::new("/")) .join(raw_target) }; let target_exists = target_path.exists(); let mut resolved_path = target_path .canonicalize() .unwrap_or_else(|_| target_path.clone()); for remaining in components { resolved_path.push(remaining.as_os_str()); } return Some(SymlinkInfo { link_path: cur, target_path, resolved_path, target_exists, }); } None } pub fn direct_symlink(path: &Path) -> Option { let meta = std::fs::symlink_metadata(path).ok()?; if meta.file_type().is_symlink() { first_symlink(path) } else { None } } #[cfg(test)] fn symlink_out_of_scope_or_plain( path: &Path, symlink: Option<&SymlinkInfo>, required_permission: &'static str, scope: &Scope, ) -> WorkdirError { if let Some(info) = symlink { let link_parent_readable = info .link_path .parent() .map(|parent| scope.is_readable(parent)) .unwrap_or(false); if info.target_exists && link_parent_readable { return WorkdirError::SymlinkOutOfScope { path: path.to_path_buf(), target: info.resolved_path.clone(), required_permission, }; } } WorkdirError::OutOfScope(path.to_path_buf()) } #[cfg(test)] fn broken_symlink_error(path: &Path, info: &SymlinkInfo) -> WorkdirError { WorkdirError::BrokenSymlink { path: path.to_path_buf(), link: info.link_path.clone(), target: info.target_path.clone(), } } // ============================================================================= // Tests // ============================================================================= #[cfg(test)] mod tests { use super::*; use manifest::{Permission, ScopeConfig, ScopeRule}; use std::fs; use tempfile::TempDir; fn make_fs(dir: &TempDir) -> LocalWorkdirSession { LocalWorkdirSession::new( Scope::writable(dir.path()).unwrap(), dir.path().to_path_buf(), ) } #[tokio::test] async fn logical_provider_operations_cover_read_write_edit_stat_and_list() { let dir = TempDir::new().unwrap(); let workdir = make_fs(&dir); let path = WorkdirPath::new("notes/item.txt").unwrap(); let written = WorkdirSession::write( &workdir, WriteRequest { path: path.clone(), content: b"alpha\nbeta\n".to_vec(), expected_hash: None, }, ) .await .unwrap(); assert!(written.created); let read = WorkdirSession::read( &workdir, ReadRequest { path: path.clone(), offset: 0, limit: 20, max_bytes: 1024, }, ) .await .unwrap(); assert_eq!(read.bytes, b"alpha\nbeta\n"); assert!(!read.truncated); let bounded = WorkdirSession::read( &workdir, ReadRequest { path: path.clone(), offset: 0, limit: 20, max_bytes: 6, }, ) .await .unwrap(); assert_eq!(bounded.bytes, b"alpha\n"); assert!(bounded.truncated); assert_eq!(bounded.content_hash, read.content_hash); let edited = WorkdirSession::edit( &workdir, EditRequest { path: path.clone(), old_string: "beta".to_owned(), new_string: "gamma".to_owned(), replace_all: false, expected_hash: read.content_hash, }, ) .await .unwrap(); assert_eq!(edited.replacements, 1); let stat = WorkdirSession::stat(&workdir, StatRequest { path: path.clone() }) .await .unwrap(); assert_eq!(stat.path, path); assert_eq!(stat.kind, EntryKind::File); let listed = WorkdirSession::list( &workdir, ListRequest { path: WorkdirPath::new("notes").unwrap(), limit: 10, }, ) .await .unwrap(); assert_eq!(listed.total_entries, 1); assert_eq!(listed.entries[0].path.as_str(), "notes/item.txt"); let error = WorkdirSession::edit( &workdir, EditRequest { path: path.clone(), old_string: "gamma".to_owned(), new_string: "delta".to_owned(), replace_all: false, expected_hash: read.content_hash, }, ) .await .unwrap_err(); assert!(matches!(error, WorkdirError::Conflict(_))); std::fs::remove_file(dir.path().join("notes/item.txt")).unwrap(); let error = WorkdirSession::write( &workdir, WriteRequest { path, content: b"replacement".to_vec(), expected_hash: Some(edited.content_hash), }, ) .await .unwrap_err(); assert!(matches!(error, WorkdirError::Conflict(_))); } #[tokio::test] async fn close_is_terminal_for_one_session_without_deleting_workdir_identity() { let dir = TempDir::new().unwrap(); std::fs::write(dir.path().join("item.txt"), "persisted").unwrap(); let workdir = Workdir::new("working-directory-42"); let scope = SharedScope::new(Scope::writable(dir.path()).unwrap()); let session = LocalWorkdirSession::materialized_bound( workdir.clone(), dir.path().to_path_buf(), dir.path().to_path_buf(), scope.clone(), WorkdirSessionCapabilities::ALL, ); assert_eq!(session.workdir(), &workdir); let command = WorkdirSession::start_command( &session, CommandRequest { command: "sleep 30".to_owned(), timeout_secs: 60, output_limit: 1024, }, ) .await .unwrap(); let waiting_session = session.clone(); let waiting_command = command.clone(); let waiter = tokio::spawn(async move { WorkdirSession::command_output( &waiting_session, CommandOutputRequest { handle: waiting_command, cursor: 0, limit: 1024, wait: true, }, ) .await }); tokio::task::yield_now().await; WorkdirSession::close(&session).await.unwrap(); WorkdirSession::close(&session).await.unwrap(); let waiter_error = tokio::time::timeout(Duration::from_secs(1), waiter) .await .expect("close should wake command output waiters") .unwrap() .unwrap_err(); assert!(matches!(waiter_error, WorkdirError::Unavailable(_))); assert!(matches!( WorkdirSession::command_status(&session, command).await, Err(WorkdirError::Unavailable(_)) )); assert!(matches!( WorkdirSession::read( &session, ReadRequest { path: WorkdirPath::new("item.txt").unwrap(), offset: 0, limit: 10, max_bytes: 1024, }, ) .await, Err(WorkdirError::Unavailable(_)) )); let restored = LocalWorkdirSession::materialized_bound( workdir.clone(), dir.path().to_path_buf(), dir.path().to_path_buf(), scope, WorkdirSessionCapabilities::ALL, ); assert_eq!(restored.workdir(), &workdir); let read = WorkdirSession::read( &restored, ReadRequest { path: WorkdirPath::new("item.txt").unwrap(), offset: 0, limit: 10, max_bytes: 1024, }, ) .await .unwrap(); assert_eq!(read.bytes, b"persisted"); } #[tokio::test] async fn capability_boundary_rejects_direct_unsupported_operation() { let dir = TempDir::new().unwrap(); let workdir = LocalWorkdirSession::materialized( dir.path().to_path_buf(), dir.path().to_path_buf(), SharedScope::new(Scope::writable(dir.path()).unwrap()), WorkdirSessionCapabilities::READ_ONLY, ); assert_eq!(workdir.root(), dir.path()); assert_eq!(workdir.cwd(), dir.path()); let error = WorkdirSession::write( &workdir, WriteRequest { path: WorkdirPath::new("blocked.txt").unwrap(), content: b"blocked".to_vec(), expected_hash: None, }, ) .await .unwrap_err(); assert!(matches!( error, WorkdirError::Unsupported(WorkdirSessionCapability::Write) )); assert!(!dir.path().join("blocked.txt").exists()); } // ------------------------------------------------------------------------- // read_bytes // ------------------------------------------------------------------------- #[test] fn read_bytes_returns_content() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let file = dir.path().join("a.txt"); fs::write(&file, b"abc").unwrap(); assert_eq!(fs.read_bytes(&file).unwrap(), b"abc"); } #[test] fn read_bytes_rejects_relative() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.read_bytes(Path::new("rel.txt")).unwrap_err(); assert!(matches!(err, WorkdirError::RelativePath(_))); } #[test] fn read_bytes_rejects_directory() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.read_bytes(dir.path()).unwrap_err(); assert!(matches!(err, WorkdirError::IsDirectory(_))); } #[test] fn read_bytes_rejects_missing() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.read_bytes(&dir.path().join("nope.txt")).unwrap_err(); assert!(matches!(err, WorkdirError::NotFound(_))); } #[test] fn read_bytes_rejects_paths_outside_scope() { let dir = TempDir::new().unwrap(); let outside = TempDir::new().unwrap(); let outside_file = outside.path().join("x.txt"); fs::write(&outside_file, b"hi").unwrap(); let scoped = make_fs(&dir); let err = scoped.read_bytes(&outside_file).unwrap_err(); assert!(matches!(err, WorkdirError::OutOfScope(_))); } #[cfg(unix)] #[test] fn read_bytes_reports_broken_symlink_target() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let link = dir.path().join("external-project"); let target = dir.path().join("missing-target"); symlink(&target, &link).unwrap(); let err = fs.read_bytes(&link).unwrap_err(); assert!( matches!( err, WorkdirError::BrokenSymlink { ref path, link: ref err_link, target: ref err_target } if path == &link && err_link == &link && err_target == &target ), "expected broken symlink diagnostic, got {err:?}" ); } #[cfg(unix)] #[test] fn read_bytes_reports_symlink_target_outside_scope() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let outside = TempDir::new().unwrap(); let target = outside.path().join("target.txt"); fs::write(&target, b"secret").unwrap(); let link = dir.path().join("outside-repo.txt"); symlink(&target, &link).unwrap(); let fs = make_fs(&dir); let err = fs.read_bytes(&link).unwrap_err(); assert!( matches!( err, WorkdirError::SymlinkOutOfScope { ref path, target: ref err_target, required_permission: "read" } if path == &link && err_target == &target.canonicalize().unwrap() ), "expected symlink out-of-scope diagnostic, got {err:?}" ); } #[cfg(unix)] #[test] fn read_bytes_allows_symlink_file_when_target_is_inside_scope() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let target = dir.path().join("target.txt"); fs::write(&target, b"visible").unwrap(); let link = dir.path().join("link.txt"); symlink(&target, &link).unwrap(); let fs = make_fs(&dir); assert_eq!(fs.read_bytes(&link).unwrap(), b"visible"); } #[cfg(unix)] #[test] fn read_bytes_reports_symlink_to_directory_as_wrong_file_type() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let target_dir = dir.path().join("target-dir"); fs::create_dir(&target_dir).unwrap(); let link = dir.path().join("dir-link"); symlink(&target_dir, &link).unwrap(); let fs = make_fs(&dir); let err = fs.read_bytes(&link).unwrap_err(); assert!( matches!( err, WorkdirError::SymlinkTargetIsDirectory { ref path, ref target } if path == &link && target == &target_dir.canonicalize().unwrap() ), "expected symlink directory type diagnostic, got {err:?}" ); } // ------------------------------------------------------------------------- // write // ------------------------------------------------------------------------- #[test] fn write_creates_new_file() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let file = dir.path().join("new.txt"); let out = fs.write(&file, b"hello").unwrap(); assert!(out.created); assert_eq!(out.bytes_written, 5); assert_eq!(fs::read(&file).unwrap(), b"hello"); } #[test] fn write_overwrites_existing() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let file = dir.path().join("a.txt"); fs::write(&file, b"old").unwrap(); let out = fs.write(&file, b"new").unwrap(); assert!(!out.created); assert_eq!(fs::read(&file).unwrap(), b"new"); } #[cfg(unix)] #[test] fn write_existing_symlink_file_updates_in_scope_target() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let target = dir.path().join("target.txt"); fs::write(&target, b"old").unwrap(); let link = dir.path().join("link.txt"); symlink(&target, &link).unwrap(); let out = fs.write(&link, b"new").unwrap(); assert!(!out.created); assert_eq!(fs::read(&target).unwrap(), b"new"); assert!( fs::symlink_metadata(&link) .unwrap() .file_type() .is_symlink() ); } #[cfg(unix)] #[test] fn write_reports_symlink_target_outside_scope() { use std::os::unix::fs::symlink; let dir = TempDir::new().unwrap(); let outside = TempDir::new().unwrap(); let target = outside.path().join("target.txt"); fs::write(&target, b"secret").unwrap(); let link = dir.path().join("outside-repo.txt"); symlink(&target, &link).unwrap(); let fs = make_fs(&dir); let err = fs.write(&link, b"new").unwrap_err(); assert!( matches!( err, WorkdirError::SymlinkOutOfScope { ref path, target: ref err_target, required_permission: "write" } if path == &link && err_target == &target.canonicalize().unwrap() ), "expected write symlink out-of-scope diagnostic, got {err:?}" ); } #[test] fn write_rejects_out_of_scope() { let dir = TempDir::new().unwrap(); let outside = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.write(&outside.path().join("x"), b"x").unwrap_err(); assert!(matches!(err, WorkdirError::OutOfScope(_))); } #[test] fn write_rejects_readonly_path() { let dir = TempDir::new().unwrap(); let sub = dir.path().join("sub"); fs::create_dir(&sub).unwrap(); let cfg = ScopeConfig { allow: vec![ScopeRule { target: dir.path().to_path_buf(), permission: Permission::Write, recursive: true, }], deny: vec![ScopeRule { target: sub.clone(), permission: Permission::Write, recursive: true, }], }; let scope = Scope::from_config(&cfg).unwrap(); let scoped = LocalWorkdirSession::new(scope, dir.path().to_path_buf()); let err = scoped.write(&sub.join("locked.txt"), b"x").unwrap_err(); assert!( matches!(err, WorkdirError::ReadOnly(_)), "expected ReadOnly, got {err:?}" ); } #[test] fn write_rejects_relative_path() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.write(Path::new("rel.txt"), b"x").unwrap_err(); assert!(matches!(err, WorkdirError::RelativePath(_))); } #[test] fn write_creates_missing_parents_inside_scope() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let nested = dir.path().join("a/b/c/deep.txt"); fs.write(&nested, b"x").unwrap(); assert_eq!(fs::read(&nested).unwrap(), b"x"); } #[test] fn write_rejects_directory_target() { let dir = TempDir::new().unwrap(); let fs = make_fs(&dir); let err = fs.write(dir.path(), b"x").unwrap_err(); assert!(matches!(err, WorkdirError::IsDirectory(_))); } // ------------------------------------------------------------------------- // Dynamic scope: SharedScope mutations propagate into LocalWorkdirSession decisions // ------------------------------------------------------------------------- #[test] fn add_allow_rule_through_shared_scope_grows_readable_set() { use manifest::SharedScope; let dir = TempDir::new().unwrap(); let extra = TempDir::new().unwrap(); let extra_file = extra.path().join("x.txt"); fs::write(&extra_file, b"hi").unwrap(); let shared = SharedScope::new(Scope::writable(dir.path()).unwrap()); let fs = LocalWorkdirSession::with_shared_scope(shared.clone(), dir.path().to_path_buf()); // Before: extra is out of scope. let err = fs.read_bytes(&extra_file).unwrap_err(); assert!(matches!(err, WorkdirError::OutOfScope(_))); // Push an allow(Read) rule. shared .update(|cur| { cur.with_added_allow_rules([ScopeRule { target: extra.path().to_path_buf(), permission: Permission::Read, recursive: true, }]) }) .unwrap(); // After: read goes through. assert_eq!(fs.read_bytes(&extra_file).unwrap(), b"hi"); // But write still fails — allow only granted Read. let err = fs.write(&extra.path().join("y.txt"), b"x").unwrap_err(); assert!( matches!(err, WorkdirError::ReadOnly(_)), "expected ReadOnly, got {err:?}" ); } #[test] fn revoke_write_through_shared_scope_blocks_subsequent_writes() { use manifest::SharedScope; let dir = TempDir::new().unwrap(); let sub = dir.path().join("sub"); fs::create_dir(&sub).unwrap(); let target = sub.join("a.txt"); let shared = SharedScope::new(Scope::writable(dir.path()).unwrap()); let fs = LocalWorkdirSession::with_shared_scope(shared.clone(), dir.path().to_path_buf()); // Write succeeds initially. fs.write(&target, b"first").unwrap(); // Revoke Write on `sub` (push a deny(Write) rule). shared .update(|cur| { cur.with_added_deny_rules([ScopeRule { target: sub.clone(), permission: Permission::Write, recursive: true, }]) }) .unwrap(); // Subsequent write fails with ReadOnly — Read is preserved. let err = fs.write(&target, b"second").unwrap_err(); assert!( matches!(err, WorkdirError::ReadOnly(_)), "expected ReadOnly after revoke, got {err:?}" ); // Read still works. assert_eq!(fs.read_bytes(&target).unwrap(), b"first"); } #[test] fn shared_scope_changes_propagate_across_clones() { use manifest::SharedScope; let dir = TempDir::new().unwrap(); let target = dir.path().join("a.txt"); let shared = SharedScope::new(Scope::writable(dir.path()).unwrap()); let fs1 = LocalWorkdirSession::with_shared_scope(shared.clone(), dir.path().to_path_buf()); let fs2 = fs1.clone(); // fs1 writes; both clones see the file. fs1.write(&target, b"hi").unwrap(); assert_eq!(fs2.read_bytes(&target).unwrap(), b"hi"); // Revoke write through the original handle. shared .update(|cur| { cur.with_added_deny_rules([ScopeRule { target: dir.path().to_path_buf(), permission: Permission::Write, recursive: true, }]) }) .unwrap(); // Both clones reject writes now — they share the same SharedScope. assert!(matches!( fs1.write(&target, b"x").unwrap_err(), WorkdirError::ReadOnly(_) )); assert!(matches!( fs2.write(&target, b"x").unwrap_err(), WorkdirError::ReadOnly(_) )); } #[tokio::test] async fn provider_executes_glob_grep_and_command_at_the_materialization() { let dir = TempDir::new().unwrap(); std::fs::create_dir(dir.path().join("src")).unwrap(); std::fs::write( dir.path().join("src/main.rs"), "fn main() { /* NEEDLE */ }\n", ) .unwrap(); let workdir = make_fs(&dir); let glob = WorkdirSession::glob( &workdir, GlobRequest { pattern: "**/*.rs".into(), path: WorkdirPath::root(), limit: 10, }, ) .await .unwrap(); assert_eq!(glob.paths, [WorkdirPath::new("src/main.rs").unwrap()]); let grep = WorkdirSession::grep( &workdir, GrepRequest { pattern: "NEEDLE".into(), path: WorkdirPath::root(), glob: None, file_type: None, case_insensitive: false, before_context: 0, after_context: 0, multiline: false, output_mode: crate::GrepOutputMode::Content, limit: 10, offset: 0, }, ) .await .unwrap(); assert_eq!(grep.match_count, 1); assert!(grep.output.contains("src/main.rs")); assert!(!grep.output.contains(dir.path().to_string_lossy().as_ref())); let handle = WorkdirSession::start_command( &workdir, CommandRequest { command: "pwd && printf provider-command".into(), timeout_secs: 5, output_limit: 4096, }, ) .await .unwrap(); let output = WorkdirSession::command_output( &workdir, CommandOutputRequest { handle, cursor: 0, limit: 4096, wait: true, }, ) .await .unwrap(); assert_eq!(output.exit_code, Some(0)); assert!(output.content.contains("provider-command")); assert!( output .content .contains(dir.path().to_string_lossy().as_ref()) ); } #[tokio::test] async fn completed_command_output_can_be_read_in_bounded_unicode_pages() { let dir = TempDir::new().unwrap(); let workdir = make_fs(&dir); let handle = WorkdirSession::start_command( &workdir, CommandRequest { command: "printf 'aéz'".into(), timeout_secs: 5, output_limit: 1024, }, ) .await .unwrap(); let first = WorkdirSession::command_output( &workdir, CommandOutputRequest { handle: handle.clone(), cursor: 0, limit: 2, wait: true, }, ) .await .unwrap(); assert_eq!(first.content, "aé"); assert_eq!(first.next_cursor, Some(2)); let second = WorkdirSession::command_output( &workdir, CommandOutputRequest { handle: handle.clone(), cursor: first.next_cursor.unwrap(), limit: 2, wait: false, }, ) .await .unwrap(); assert_eq!(second.content, "z"); assert_eq!(second.next_cursor, None); assert!(matches!( WorkdirSession::command_status(&workdir, handle).await, Err(WorkdirError::UnknownCommand(_)) )); } #[tokio::test] async fn provider_cancels_active_command() { let dir = TempDir::new().unwrap(); let workdir = make_fs(&dir); let handle = WorkdirSession::start_command( &workdir, CommandRequest { command: "sleep 30".into(), timeout_secs: 60, output_limit: 1024, }, ) .await .unwrap(); assert_eq!( WorkdirSession::command_status(&workdir, handle.clone()) .await .unwrap(), CommandStatus::Running ); let waiting_workdir = workdir.clone(); let waiting_handle = handle.clone(); let waiter = tokio::spawn(async move { WorkdirSession::command_output( &waiting_workdir, CommandOutputRequest { handle: waiting_handle, cursor: 0, limit: 1024, wait: true, }, ) .await }); tokio::task::yield_now().await; WorkdirSession::cancel_command(&workdir, handle.clone()) .await .unwrap(); let waiter_error = tokio::time::timeout(Duration::from_secs(1), waiter) .await .expect("cancel should wake command output waiters") .unwrap() .unwrap_err(); assert!(matches!(waiter_error, WorkdirError::UnknownCommand(_))); assert!(matches!( WorkdirSession::command_status(&workdir, handle).await, Err(WorkdirError::UnknownCommand(_)) )); } }