From 41b7b289d0917fb84d3d785d51949c56a6482e81 Mon Sep 17 00:00:00 2001 From: Hare Date: Mon, 14 Sep 2026 20:32:21 +0900 Subject: [PATCH] fix: serialize Workdir lease admission with writes --- crates/workdir/src/scope.rs | 209 +++++++++++++++++++++++++++++++++++- 1 file changed, 208 insertions(+), 1 deletion(-) diff --git a/crates/workdir/src/scope.rs b/crates/workdir/src/scope.rs index 104c58ad..5e158d23 100644 --- a/crates/workdir/src/scope.rs +++ b/crates/workdir/src/scope.rs @@ -707,7 +707,7 @@ impl ScopedWorkdirSession { ActiveWriteLease { validity: Arc::downgrade(&validity), cleanup_pending: Arc::downgrade(&cleanup_pending), - rules: request.rules.clone(), + rules: write_rules, }, ); } @@ -792,6 +792,7 @@ impl WorkdirSession for ScopedWorkdirSession { } async fn write(&self, mut request: WriteRequest) -> Result { + let _scope_guard = self.scope_lock.lock().await; let path = self .resolve_operation_path(&request.path, WorkdirToolScopePermission::Write) .await?; @@ -801,6 +802,7 @@ impl WorkdirSession for ScopedWorkdirSession { } async fn edit(&self, mut request: EditRequest) -> Result { + let _scope_guard = self.scope_lock.lock().await; let path = self .resolve_operation_path(&request.path, WorkdirToolScopePermission::Write) .await?; @@ -1294,6 +1296,134 @@ mod tests { ))) } + #[derive(Debug)] + struct BlockingWriteSession { + inner: Arc, + entered: tokio::sync::watch::Sender, + release: Arc, + block_next_write: std::sync::atomic::AtomicBool, + } + + #[async_trait] + impl WorkdirSession for BlockingWriteSession { + fn workdir(&self) -> &Workdir { + self.inner.workdir() + } + + fn capabilities(&self) -> WorkdirSessionCapabilities { + self.inner.capabilities() + } + + async fn authorize_scope_path( + &self, + request: WorkdirScopeAuthorizationRequest, + ) -> Result<(), WorkdirError> { + self.inner.authorize_scope_path(request).await + } + + async fn scope_rules_overlap( + &self, + request: WorkdirScopeOverlapRequest, + ) -> Result { + self.inner.scope_rules_overlap(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 { + if self.block_next_write.swap(false, Ordering::AcqRel) { + let _ = self.entered.send(true); + self.release.notified().await; + } + WorkdirSession::write(self.inner.as_ref(), request).await + } + + async fn edit(&self, request: EditRequest) -> Result { + self.inner.edit(request).await + } + + 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 { + self.inner.start_command(request).await + } + + async fn command_status( + &self, + handle: CommandHandle, + ) -> Result { + self.inner.command_status(handle).await + } + + async fn command_output( + &self, + request: CommandOutputRequest, + ) -> Result { + self.inner.command_output(request).await + } + + async fn cancel_command(&self, handle: CommandHandle) -> Result<(), WorkdirError> { + self.inner.cancel_command(handle).await + } + + fn subscribe_command_events(&self) -> Option> { + self.inner.subscribe_command_events() + } + + fn command_snapshot(&self) -> Vec { + self.inner.command_snapshot() + } + + async fn close(&self) -> Result<(), WorkdirError> { + self.inner.close().await + } + } + + fn blocking_session( + root: &Path, + ) -> ( + WorkdirToolBroker, + tokio::sync::watch::Receiver, + Arc, + ) { + let scope = SharedScope::new(Scope::writable(root).unwrap()); + let inner = Arc::new(LocalWorkdirSession::materialized_bound( + Workdir::new("blocking-delegation-test"), + root.to_path_buf(), + root.to_path_buf(), + scope, + WorkdirSessionCapabilities::ALL, + )); + let (entered, receiver) = tokio::sync::watch::channel(false); + let release = Arc::new(tokio::sync::Notify::new()); + let source = Arc::new(BlockingWriteSession { + inner, + entered, + release: release.clone(), + block_next_write: std::sync::atomic::AtomicBool::new(true), + }); + (WorkdirToolBroker::new(source), receiver, release) + } + fn request(path: &str, permission: WorkdirToolScopePermission) -> WorkdirToolScope { WorkdirToolScope { rules: vec![WorkdirToolScopeRule { @@ -1616,6 +1746,83 @@ mod tests { )); } + #[tokio::test] + async fn write_and_overlapping_scope_admission_are_serialized() { + let root = TempDir::new().unwrap(); + fs::create_dir_all(root.path().join("shared")).unwrap(); + let (parent, mut entered, release) = blocking_session(root.path()); + let writer = { + let parent = parent.clone(); + tokio::spawn(async move { parent.write(write("shared/file", "written")).await }) + }; + entered.changed().await.unwrap(); + assert!(*entered.borrow()); + + let mut admission = { + let parent = parent.clone(); + tokio::spawn(async move { + parent + .scope(request("shared", WorkdirToolScopePermission::Write)) + .await + }) + }; + assert!( + tokio::time::timeout(std::time::Duration::from_millis(50), &mut admission) + .await + .is_err(), + "scope admission must wait for the in-flight parent write" + ); + + release.notify_waiters(); + writer.await.unwrap().unwrap(); + let lease = tokio::time::timeout(std::time::Duration::from_secs(1), admission) + .await + .expect("scope admission should resume after write completion") + .unwrap() + .unwrap(); + drop(lease); + } + + #[tokio::test] + async fn read_rules_do_not_expand_child_write_lease_conflicts() { + 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 + .scope(WorkdirToolScope { + rules: vec![ + WorkdirToolScopeRule { + target: fs_path("leased"), + permission: WorkdirToolScopePermission::Write, + recursive: true, + symlink_policy: SymlinkPolicy::Resolved, + }, + WorkdirToolScopeRule { + target: FsPath::root(), + permission: WorkdirToolScopePermission::Read, + recursive: true, + symlink_policy: SymlinkPolicy::Resolved, + }, + ], + cwd: fs_path("leased"), + command: false, + }) + .await + .unwrap(); + + parent + .write(write("other/parent", "allowed")) + .await + .unwrap(); + let sibling = parent + .scope(request("other", WorkdirToolScopePermission::Write)) + .await + .unwrap(); + sibling.write(write("sibling", "allowed")).await.unwrap(); + drop(child); + } + #[cfg(unix)] #[tokio::test] async fn sibling_write_scopes_reject_distinct_aliases_to_same_resolved_target() {