diff --git a/crates/workdir/src/delegation.rs b/crates/workdir/src/delegation.rs new file mode 100644 index 00000000..19c3beea --- /dev/null +++ b/crates/workdir/src/delegation.rs @@ -0,0 +1,781 @@ +use std::collections::HashMap; +use std::path::Path; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, Mutex, Weak}; + +use async_trait::async_trait; +use fs_operation::{ + EditRequest, EditResult, FsPath, GlobRequest, GlobResult, GrepRequest, GrepResult, ListRequest, + ListResult, ReadRequest, ReadResult, StatRequest, StatResult, WriteRequest, WriteResult, +}; + +use crate::{ + CommandHandle, CommandOutput, CommandOutputRequest, CommandRequest, CommandStatus, Workdir, + WorkdirError, WorkdirSession, WorkdirSessionCapabilities, WorkdirSessionCapability, + WorkdirSessionHandle, +}; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WorkdirDelegationPermission { + Read, + Write, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct WorkdirDelegationRule { + pub target: FsPath, + pub permission: WorkdirDelegationPermission, + pub recursive: bool, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct WorkdirDelegationRequest { + pub rules: Vec, + pub cwd: FsPath, +} + +pub struct WorkdirDelegation { + pub scoped_session: WorkdirSessionHandle, + pub capabilities: WorkdirSessionCapabilities, + validity: Arc, +} + +impl std::fmt::Debug for WorkdirDelegation { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("WorkdirDelegation") + .field("workdir", &self.scoped_session.workdir()) + .field("capabilities", &self.capabilities) + .field("active", &self.is_active()) + .finish() + } +} + +impl WorkdirDelegation { + pub fn is_active(&self) -> bool { + self.validity.is_active() + } + + pub fn release(&self) { + self.validity.active.store(false, Ordering::Release); + } +} + +impl Drop for WorkdirDelegation { + fn drop(&mut self) { + self.release(); + } +} + +#[derive(Debug)] +struct SessionValidity { + active: AtomicBool, + parent: Option>, +} + +impl SessionValidity { + fn root() -> Arc { + Arc::new(Self { + active: AtomicBool::new(true), + parent: None, + }) + } + + fn child(parent: Arc) -> Arc { + Arc::new(Self { + active: AtomicBool::new(true), + parent: Some(parent), + }) + } + + fn is_active(&self) -> bool { + self.active.load(Ordering::Acquire) + && self.parent.as_ref().is_none_or(|parent| parent.is_active()) + } +} + +#[derive(Clone, Debug)] +struct ActiveWriteLease { + validity: Weak, + rules: Vec, +} + +struct DelegatingWorkdirSession { + source: WorkdirSessionHandle, + scope: Option>, + capabilities: WorkdirSessionCapabilities, + validity: Arc, + child_write_leases: Mutex>, + next_lease_id: AtomicU64, + closes_source: bool, +} + +impl std::fmt::Debug for DelegatingWorkdirSession { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("DelegatingWorkdirSession") + .field("workdir", &self.source.workdir()) + .field("scope", &self.scope) + .field("capabilities", &self.capabilities) + .field("active", &self.validity.is_active()) + .finish_non_exhaustive() + } +} + +/// Wrap a provider session with logical-path delegation and parent write gates. +pub fn delegation_capable_session(source: WorkdirSessionHandle) -> WorkdirSessionHandle { + let capabilities = source.capabilities(); + Arc::new(DelegatingWorkdirSession { + source, + scope: None, + capabilities, + validity: SessionValidity::root(), + child_write_leases: Mutex::new(HashMap::new()), + next_lease_id: AtomicU64::new(1), + closes_source: true, + }) +} + +impl DelegatingWorkdirSession { + fn ensure_active(&self) -> Result<(), WorkdirError> { + if self.validity.is_active() { + Ok(()) + } else { + Err(WorkdirError::SessionClosed) + } + } + + fn ensure_capability( + &self, + required: WorkdirSessionCapability, + operation: &'static str, + ) -> Result<(), WorkdirError> { + self.ensure_active()?; + if self.capabilities.supports(required) { + Ok(()) + } else { + Err(WorkdirError::Denied(format!( + "delegated workdir session does not permit {operation}" + ))) + } + } + + fn ensure_path( + &self, + path: &FsPath, + permission: WorkdirDelegationPermission, + ) -> Result<(), WorkdirError> { + self.ensure_active()?; + if let Some(scope) = &self.scope { + if !scope + .iter() + .any(|rule| rule_allows_path(rule, path, permission)) + { + return Err(WorkdirError::Denied(format!( + "logical workdir path `{path}` is outside the delegated {permission:?} scope" + ))); + } + } + if permission == WorkdirDelegationPermission::Write { + self.ensure_parent_write_available(path)?; + } + Ok(()) + } + + fn ensure_read( + &self, + path: &FsPath, + capability: WorkdirSessionCapability, + ) -> Result<(), WorkdirError> { + self.ensure_capability(capability, "read operations")?; + self.ensure_path(path, WorkdirDelegationPermission::Read) + } + + fn ensure_write( + &self, + path: &FsPath, + capability: WorkdirSessionCapability, + ) -> Result<(), WorkdirError> { + self.ensure_capability(capability, "write operations")?; + self.ensure_path(path, WorkdirDelegationPermission::Write) + } + + fn ensure_command(&self, starting: bool) -> Result<(), WorkdirError> { + self.ensure_capability(WorkdirSessionCapability::Command, "command execution")?; + if starting && self.has_active_write_lease() { + return Err(WorkdirError::Denied( + "command execution is denied while a child holds a write delegation".into(), + )); + } + Ok(()) + } + + fn ensure_parent_write_available(&self, path: &FsPath) -> Result<(), WorkdirError> { + let mut leases = self + .child_write_leases + .lock() + .expect("workdir delegation lease mutex poisoned"); + leases.retain(|_, lease| lease.validity.upgrade().is_some_and(|v| v.is_active())); + if leases.values().any(|lease| { + lease.rules.iter().any(|rule| { + rule.permission == WorkdirDelegationPermission::Write + && rule_allows_path(rule, path, WorkdirDelegationPermission::Write) + }) + }) { + Err(WorkdirError::Denied(format!( + "logical workdir path `{path}` is leased to a child session" + ))) + } else { + Ok(()) + } + } + + fn has_active_write_lease(&self) -> bool { + let mut leases = self + .child_write_leases + .lock() + .expect("workdir delegation lease mutex poisoned"); + leases.retain(|_, lease| lease.validity.upgrade().is_some_and(|v| v.is_active())); + leases.values().any(|lease| { + lease + .rules + .iter() + .any(|rule| rule.permission == WorkdirDelegationPermission::Write) + }) + } + + fn validate_delegation_rules( + &self, + rules: &[WorkdirDelegationRule], + ) -> Result { + self.ensure_active()?; + if rules.is_empty() { + return Err(WorkdirError::Denied( + "workdir delegation requires at least one logical scope rule".into(), + )); + } + let writable = rules + .iter() + .any(|rule| rule.permission == WorkdirDelegationPermission::Write); + if !self.capabilities.supports(WorkdirSessionCapability::Read) + || (writable + && (!self.capabilities.supports(WorkdirSessionCapability::Write) + || !self.capabilities.supports(WorkdirSessionCapability::Edit))) + { + return Err(WorkdirError::Denied( + "parent workdir session cannot delegate the requested capabilities".into(), + )); + } + for requested in rules { + if let Some(scope) = &self.scope { + if !scope + .iter() + .any(|parent| rule_contains_rule(parent, requested)) + { + return Err(WorkdirError::Denied(format!( + "logical workdir scope `{}` exceeds the parent delegation", + requested.target + ))); + } + } + } + let mut delegated = vec![WorkdirSessionCapability::Read]; + for capability in [ + WorkdirSessionCapability::Glob, + WorkdirSessionCapability::Grep, + ] { + if self.capabilities.supports(capability) { + delegated.push(capability); + } + } + if writable { + delegated.push(WorkdirSessionCapability::Write); + delegated.push(WorkdirSessionCapability::Edit); + } + Ok(WorkdirSessionCapabilities::from_capabilities(delegated)) + } +} + +#[async_trait] +impl WorkdirSession for DelegatingWorkdirSession { + fn workdir(&self) -> &Workdir { + self.source.workdir() + } + + fn capabilities(&self) -> WorkdirSessionCapabilities { + self.capabilities + } + + fn is_delegation_capable(&self) -> bool { + true + } + + async fn capture_delegation_source(&self) -> Result { + self.ensure_active()?; + if self.scope.is_some() { + return Err(WorkdirError::Denied( + "scoped Workdir sessions cannot expose their provider source".into(), + )); + } + self.source.capture_delegation_source().await + } + + async fn delegate( + &self, + request: WorkdirDelegationRequest, + ) -> Result { + let capabilities = self.validate_delegation_rules(&request.rules)?; + if !request + .rules + .iter() + .any(|rule| rule_allows_path(rule, &request.cwd, WorkdirDelegationPermission::Read)) + { + return Err(WorkdirError::Denied(format!( + "delegated cwd `{}` is outside the delegated readable scope", + request.cwd + ))); + } + let source = self.source.capture_delegation_source().await?; + let validity = SessionValidity::child(self.validity.clone()); + let id = self.next_lease_id.fetch_add(1, Ordering::Relaxed); + if request + .rules + .iter() + .any(|rule| rule.permission == WorkdirDelegationPermission::Write) + { + self.child_write_leases + .lock() + .expect("workdir delegation lease mutex poisoned") + .insert( + id, + ActiveWriteLease { + validity: Arc::downgrade(&validity), + rules: request.rules.clone(), + }, + ); + } + let child: WorkdirSessionHandle = Arc::new(DelegatingWorkdirSession { + source, + scope: Some(request.rules), + capabilities, + validity: validity.clone(), + child_write_leases: Mutex::new(HashMap::new()), + next_lease_id: AtomicU64::new(1), + closes_source: false, + }); + let scoped_session: WorkdirSessionHandle = + if capabilities == WorkdirSessionCapabilities::READ_ONLY { + Arc::new(ReadOnlyWorkdirSession::new(child)) + } else { + child + }; + Ok(WorkdirDelegation { + scoped_session, + capabilities, + validity, + }) + } + + async fn stat(&self, request: StatRequest) -> Result { + self.ensure_read(&request.path, WorkdirSessionCapability::Read)?; + self.source.stat(request).await + } + + async fn read(&self, request: ReadRequest) -> Result { + self.ensure_read(&request.path, WorkdirSessionCapability::Read)?; + self.source.read(request).await + } + + async fn write(&self, request: WriteRequest) -> Result { + self.ensure_write(&request.path, WorkdirSessionCapability::Write)?; + self.source.write(request).await + } + + async fn edit(&self, request: EditRequest) -> Result { + self.ensure_write(&request.path, WorkdirSessionCapability::Edit)?; + self.source.edit(request).await + } + + async fn list(&self, request: ListRequest) -> Result { + self.ensure_read(&request.path, WorkdirSessionCapability::Read)?; + self.source.list(request).await + } + + async fn glob(&self, request: GlobRequest) -> Result { + self.ensure_read(&request.path, WorkdirSessionCapability::Glob)?; + self.source.glob(request).await + } + + async fn grep(&self, request: GrepRequest) -> Result { + self.ensure_read(&request.path, WorkdirSessionCapability::Grep)?; + self.source.grep(request).await + } + + async fn start_command(&self, request: CommandRequest) -> Result { + self.ensure_command(true)?; + self.source.start_command(request).await + } + + async fn command_status(&self, handle: CommandHandle) -> Result { + self.ensure_command(false)?; + self.source.command_status(handle).await + } + + async fn command_output( + &self, + request: CommandOutputRequest, + ) -> Result { + self.ensure_command(false)?; + self.source.command_output(request).await + } + + async fn cancel_command(&self, handle: CommandHandle) -> Result<(), WorkdirError> { + self.ensure_command(false)?; + self.source.cancel_command(handle).await + } + + async fn close(&self) -> Result<(), WorkdirError> { + self.validity.active.store(false, Ordering::Release); + if self.closes_source { + self.source.close().await + } else { + Ok(()) + } + } +} + +/// A fail-closed read-only view over an already scoped delegated session. +#[derive(Debug)] +pub struct ReadOnlyWorkdirSession { + inner: WorkdirSessionHandle, +} + +impl ReadOnlyWorkdirSession { + pub fn new(inner: WorkdirSessionHandle) -> Self { + Self { inner } + } +} + +#[async_trait] +impl WorkdirSession for ReadOnlyWorkdirSession { + fn workdir(&self) -> &Workdir { + self.inner.workdir() + } + + fn capabilities(&self) -> WorkdirSessionCapabilities { + WorkdirSessionCapabilities::READ_ONLY + } + + fn is_delegation_capable(&self) -> bool { + true + } + + async fn delegate( + &self, + request: WorkdirDelegationRequest, + ) -> Result { + if request + .rules + .iter() + .any(|rule| rule.permission == WorkdirDelegationPermission::Write) + { + return Err(WorkdirError::Denied( + "read-only workdir session cannot delegate write access".into(), + )); + } + self.inner.delegate(request).await + } + + async fn stat(&self, request: StatRequest) -> Result { + self.inner.stat(request).await + } + + async fn read(&self, request: ReadRequest) -> Result { + self.inner.read(request).await + } + + async fn write(&self, _request: WriteRequest) -> Result { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn edit(&self, _request: EditRequest) -> Result { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn list(&self, request: ListRequest) -> Result { + self.inner.list(request).await + } + + async fn glob(&self, request: GlobRequest) -> Result { + self.inner.glob(request).await + } + + async fn grep(&self, request: GrepRequest) -> Result { + self.inner.grep(request).await + } + + async fn start_command(&self, _request: CommandRequest) -> Result { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn command_status(&self, _handle: CommandHandle) -> Result { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn command_output( + &self, + _request: CommandOutputRequest, + ) -> Result { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn cancel_command(&self, _handle: CommandHandle) -> Result<(), WorkdirError> { + Err(WorkdirError::Denied("read-only workdir session".into())) + } + + async fn close(&self) -> Result<(), WorkdirError> { + self.inner.close().await + } +} + +fn rule_allows_path( + rule: &WorkdirDelegationRule, + path: &FsPath, + required: WorkdirDelegationPermission, +) -> bool { + if required == WorkdirDelegationPermission::Write + && rule.permission != WorkdirDelegationPermission::Write + { + return false; + } + path_in_rule(rule, path) +} + +fn path_in_rule(rule: &WorkdirDelegationRule, path: &FsPath) -> bool { + let target = Path::new(rule.target.as_str()); + let path = Path::new(path.as_str()); + if path == target { + return true; + } + let Ok(suffix) = path.strip_prefix(target) else { + return false; + }; + let depth = suffix.components().count(); + rule.recursive || depth <= 1 +} + +fn rule_contains_rule(parent: &WorkdirDelegationRule, child: &WorkdirDelegationRule) -> bool { + if child.permission == WorkdirDelegationPermission::Write + && parent.permission != WorkdirDelegationPermission::Write + { + return false; + } + if !path_in_rule(parent, &child.target) { + return false; + } + parent.recursive || !child.recursive +} + +#[cfg(test)] +mod tests { + use std::fs; + + use manifest::{Permission, Scope, ScopeConfig, ScopeRule, SharedScope}; + use tempfile::TempDir; + + use super::*; + use crate::LocalWorkdirSession; + + fn fs_path(path: &str) -> FsPath { + FsPath::new(path).unwrap() + } + + fn session(root: &Path) -> WorkdirSessionHandle { + let scope = SharedScope::new( + Scope::from_config(&ScopeConfig { + allow: vec![ScopeRule { + target: root.to_path_buf(), + permission: Permission::Write, + recursive: true, + }], + deny: Vec::new(), + }) + .unwrap(), + ); + delegation_capable_session(Arc::new(LocalWorkdirSession::materialized_bound( + Workdir::new("delegation-test"), + root.to_path_buf(), + root.to_path_buf(), + scope, + WorkdirSessionCapabilities::ALL, + ))) + } + + fn request(path: &str, permission: WorkdirDelegationPermission) -> WorkdirDelegationRequest { + WorkdirDelegationRequest { + rules: vec![WorkdirDelegationRule { + target: fs_path(path), + permission, + recursive: true, + }], + cwd: fs_path(path), + } + } + + fn read(path: &str) -> ReadRequest { + ReadRequest { + path: fs_path(path), + offset: 0, + limit: 20, + max_bytes: 1024, + } + } + + fn write(path: &str, content: &str) -> WriteRequest { + WriteRequest { + path: fs_path(path), + content: content.as_bytes().to_vec(), + expected_hash: None, + } + } + + #[test] + fn non_recursive_rule_covers_target_and_direct_children_only() { + let rule = WorkdirDelegationRule { + target: fs_path("docs"), + permission: WorkdirDelegationPermission::Read, + recursive: false, + }; + assert!(path_in_rule(&rule, &fs_path("docs"))); + assert!(path_in_rule(&rule, &fs_path("docs/readme.md"))); + assert!(!path_in_rule(&rule, &fs_path("docs/guides/start.md"))); + } + + #[tokio::test] + async fn read_only_delegation_allows_prefix_and_denies_siblings_and_mutation() { + let root = TempDir::new().unwrap(); + fs::create_dir_all(root.path().join("docs")).unwrap(); + fs::create_dir_all(root.path().join("secret")).unwrap(); + fs::write(root.path().join("docs/readme.md"), "visible").unwrap(); + fs::write(root.path().join("secret/key"), "hidden").unwrap(); + let parent = session(root.path()); + + let child = parent + .delegate(request("docs", WorkdirDelegationPermission::Read)) + .await + .unwrap(); + assert_eq!(child.capabilities, WorkdirSessionCapabilities::READ_ONLY); + assert_eq!( + child + .scoped_session + .read(read("docs/readme.md")) + .await + .unwrap() + .bytes, + b"visible" + ); + assert!(matches!( + child.scoped_session.read(read("secret/key")).await, + Err(WorkdirError::Denied(_)) + )); + assert!(matches!( + child.scoped_session.write(write("docs/new.md", "no")).await, + Err(WorkdirError::Denied(_)) + )); + assert!( + !child + .capabilities + .supports(WorkdirSessionCapability::Command) + ); + } + + #[tokio::test] + async fn write_lease_blocks_parent_region_until_release() { + let root = TempDir::new().unwrap(); + fs::create_dir_all(root.path().join("leased")).unwrap(); + fs::create_dir_all(root.path().join("other")).unwrap(); + let parent = session(root.path()); + let child = parent + .delegate(request("leased", WorkdirDelegationPermission::Write)) + .await + .unwrap(); + + assert!(matches!( + parent.write(write("leased/file", "parent")).await, + Err(WorkdirError::Denied(_)) + )); + parent.write(write("other/file", "parent")).await.unwrap(); + child + .scoped_session + .write(write("leased/file", "child")) + .await + .unwrap(); + child.release(); + parent + .write(write("leased/parent", "parent")) + .await + .unwrap(); + assert!(matches!( + child.scoped_session.read(read("leased/file")).await, + Err(WorkdirError::SessionClosed) + )); + } + + #[tokio::test] + async fn nested_delegation_is_attenuated_and_parent_revocation_cascades() { + let root = TempDir::new().unwrap(); + fs::create_dir_all(root.path().join("docs/sub")).unwrap(); + fs::create_dir_all(root.path().join("docs/peer")).unwrap(); + fs::write(root.path().join("docs/sub/a"), "a").unwrap(); + fs::write(root.path().join("docs/peer/b"), "b").unwrap(); + let root_session = session(root.path()); + let child = root_session + .delegate(request("docs", WorkdirDelegationPermission::Read)) + .await + .unwrap(); + let nested = child + .scoped_session + .delegate(request("docs/sub", WorkdirDelegationPermission::Read)) + .await + .unwrap(); + + nested + .scoped_session + .read(read("docs/sub/a")) + .await + .unwrap(); + assert!(matches!( + nested.scoped_session.read(read("docs/peer/b")).await, + Err(WorkdirError::Denied(_)) + )); + assert!( + child + .scoped_session + .delegate(request("docs/sub", WorkdirDelegationPermission::Write)) + .await + .is_err() + ); + + child.release(); + assert!(matches!( + nested.scoped_session.read(read("docs/sub/a")).await, + Err(WorkdirError::SessionClosed) + )); + } + + #[tokio::test] + async fn closing_parent_invalidates_delegated_sessions() { + let root = TempDir::new().unwrap(); + fs::create_dir_all(root.path().join("docs")).unwrap(); + fs::write(root.path().join("docs/a"), "a").unwrap(); + let parent = session(root.path()); + let child = parent + .delegate(request("docs", WorkdirDelegationPermission::Read)) + .await + .unwrap(); + + parent.close().await.unwrap(); + assert!(matches!( + child.scoped_session.read(read("docs/a")).await, + Err(WorkdirError::SessionClosed) + )); + } +} diff --git a/crates/workdir/src/http.rs b/crates/workdir/src/http.rs index 8dd73c17..c6855751 100644 --- a/crates/workdir/src/http.rs +++ b/crates/workdir/src/http.rs @@ -120,7 +120,10 @@ impl WorkdirTransportError { WorkdirError::UnknownCommand(_) => { (Code::UnknownCommand, "Workdir command was not found") } - WorkdirError::Unavailable(_) => (Code::Unavailable, "Workdir session is unavailable"), + WorkdirError::Unavailable(_) | WorkdirError::SessionClosed => { + (Code::Unavailable, "Workdir session is unavailable") + } + WorkdirError::Denied(_) => (Code::InvalidRequest, "Workdir operation was denied"), WorkdirError::Transport(_) => (Code::Internal, "Workdir transport failed"), WorkdirError::InvalidPath(_) | WorkdirError::RelativePath(_) @@ -169,7 +172,7 @@ mod client { use reqwest::{Client, StatusCode, Url}; use super::*; - use crate::{Workdir, WorkdirSession}; + use crate::{Workdir, WorkdirSession, WorkdirSessionHandle}; /// Provides a fresh bearer token for each Runtime request. Backend /// implementations can mint short-lived capability tokens without making a @@ -310,6 +313,21 @@ mod client { self.capabilities } + async fn capture_delegation_source(&self) -> Result { + if self.closed.load(Ordering::Acquire) { + return Err(WorkdirError::SessionClosed); + } + Ok(Arc::new(Self { + client: self.client.clone(), + base_url: self.base_url.clone(), + authorization: self.authorization.clone(), + workdir: self.workdir.clone(), + session_id: self.session_id.clone(), + capabilities: self.capabilities, + closed: AtomicBool::new(false), + })) + } + async fn stat(&self, request: StatRequest) -> Result { match self.operate(WorkdirSessionOperation::Stat(request)).await? { WorkdirSessionOperationResult::Stat(result) => Ok(result), diff --git a/crates/workdir/src/lib.rs b/crates/workdir/src/lib.rs index 957c7bbd..d165e67f 100644 --- a/crates/workdir/src/lib.rs +++ b/crates/workdir/src/lib.rs @@ -5,6 +5,7 @@ //! bound to one Worker. Tools consume sessions; they do not own Workdir //! materialization or cleanup. +mod delegation; pub mod http; mod local; mod operation; @@ -12,11 +13,14 @@ pub mod workspace; use std::path::{Path, PathBuf}; use std::sync::Arc; -use std::sync::atomic::{AtomicBool, Ordering}; use async_trait::async_trait; use serde::{Deserialize, Serialize}; +pub use delegation::{ + ReadOnlyWorkdirSession, WorkdirDelegation, WorkdirDelegationPermission, + WorkdirDelegationRequest, WorkdirDelegationRule, delegation_capable_session, +}; pub use fs_operation::{ ContentHash, EditRequest, EditResult, EntryKind, FsPath as WorkdirPath, GlobRequest, GlobResult, GrepOutputMode, GrepRequest, GrepResult, ListEntry, ListRequest, ListResult, @@ -140,6 +144,30 @@ pub trait WorkdirSession: std::fmt::Debug + Send + Sync { fn workdir(&self) -> &Workdir; fn capabilities(&self) -> WorkdirSessionCapabilities; + fn is_delegation_capable(&self) -> bool { + false + } + + /// Capture a provider-specific source for a delegated child session. + /// Remote providers use this boundary to pin attachment identity without + /// exposing transport handles or host paths. + async fn capture_delegation_source(&self) -> Result { + Err(WorkdirError::Denied( + "workdir provider does not support delegated sessions".into(), + )) + } + + /// Attenuate this session into a revocable child lease. Only sessions + /// created with [`delegation_capable_session`] implement this operation. + async fn delegate( + &self, + _request: WorkdirDelegationRequest, + ) -> Result { + Err(WorkdirError::Denied( + "workdir session is not delegation-capable".into(), + )) + } + async fn stat(&self, request: StatRequest) -> Result; async fn read(&self, request: ReadRequest) -> Result; async fn write(&self, request: WriteRequest) -> Result; @@ -160,107 +188,14 @@ pub trait WorkdirSession: std::fmt::Debug + Send + Sync { pub type WorkdirSessionHandle = Arc; -/// Ephemeral least-authority view over an existing Workdir session. -/// -/// The wrapper exposes only stat/read/list/glob/grep and never forwards write, -/// edit, command, or close authority to the underlying Worker session. Closing -/// the wrapper is terminal for the view but deliberately leaves the owner's -/// source session open. -#[derive(Debug)] -pub struct ReadOnlyWorkdirSession { - source: WorkdirSessionHandle, - closed: AtomicBool, -} - -impl ReadOnlyWorkdirSession { - pub fn new(source: WorkdirSessionHandle) -> Self { - Self { - source, - closed: AtomicBool::new(false), - } - } - - fn ensure_open(&self) -> Result<(), WorkdirError> { - if self.closed.load(Ordering::Acquire) { - Err(WorkdirError::Unavailable( - "read-only Workdir session is closed".to_string(), - )) - } else { - Ok(()) - } - } -} - -#[async_trait] -impl WorkdirSession for ReadOnlyWorkdirSession { - fn workdir(&self) -> &Workdir { - self.source.workdir() - } - - fn capabilities(&self) -> WorkdirSessionCapabilities { - WorkdirSessionCapabilities::READ_ONLY - } - - async fn stat(&self, request: StatRequest) -> Result { - self.ensure_open()?; - self.source.stat(request).await - } - - async fn read(&self, request: ReadRequest) -> Result { - self.ensure_open()?; - self.source.read(request).await - } - - async fn write(&self, _request: WriteRequest) -> Result { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Write)) - } - - async fn edit(&self, _request: EditRequest) -> Result { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Edit)) - } - - async fn list(&self, request: ListRequest) -> Result { - self.ensure_open()?; - self.source.list(request).await - } - - async fn glob(&self, request: GlobRequest) -> Result { - self.ensure_open()?; - self.source.glob(request).await - } - - async fn grep(&self, request: GrepRequest) -> Result { - self.ensure_open()?; - self.source.grep(request).await - } - - async fn start_command(&self, _request: CommandRequest) -> Result { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Command)) - } - - async fn command_status(&self, _handle: CommandHandle) -> Result { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Command)) - } - - async fn command_output( - &self, - _request: CommandOutputRequest, - ) -> Result { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Command)) - } - - async fn cancel_command(&self, _handle: CommandHandle) -> Result<(), WorkdirError> { - Err(WorkdirError::Unsupported(WorkdirSessionCapability::Command)) - } - - async fn close(&self) -> Result<(), WorkdirError> { - self.closed.store(true, Ordering::Release); - Ok(()) - } -} - #[derive(Debug, thiserror::Error)] pub enum WorkdirError { + #[error("Workdir operation denied: {0}")] + Denied(String), + + #[error("Workdir session is closed")] + SessionClosed, + #[error("Workdir session does not support {0:?}")] Unsupported(WorkdirSessionCapability), diff --git a/crates/workdir/src/local.rs b/crates/workdir/src/local.rs index dd36b997..6fdfc99c 100644 --- a/crates/workdir/src/local.rs +++ b/crates/workdir/src/local.rs @@ -29,8 +29,8 @@ 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, + WorkdirSession, WorkdirSessionCapabilities, WorkdirSessionCapability, WorkdirSessionHandle, + WriteRequest, WriteResult, }; #[cfg(test)] use crate::{EntryKind, WriteOutcome}; @@ -371,6 +371,10 @@ impl WorkdirSession for LocalWorkdirSession { self.inner.capabilities } + async fn capture_delegation_source(&self) -> Result { + Ok(Arc::new(self.clone())) + } + async fn stat(&self, request: StatRequest) -> Result { self.ensure_capability(WorkdirSessionCapability::Read)?; let logical = request.path.clone(); diff --git a/crates/workdir/src/workspace.rs b/crates/workdir/src/workspace.rs index f55a7589..ce2ed61c 100644 --- a/crates/workdir/src/workspace.rs +++ b/crates/workdir/src/workspace.rs @@ -314,3 +314,17 @@ mod tests { assert_eq!(decoded, detail); } } + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct WorkspaceWorkdirSessionOperationRequest { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub expected_session_fence: Option, + pub operation: crate::http::WorkdirSessionOperation, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct WorkspaceWorkdirSessionFence { + pub value: String, +}