fix: close scoped SubWorker command authority
This commit is contained in:
+190
-15
@@ -106,6 +106,7 @@ pub struct WorkdirScopeLease {
|
|||||||
broker: WorkdirToolBroker,
|
broker: WorkdirToolBroker,
|
||||||
pub capabilities: WorkdirSessionCapabilities,
|
pub capabilities: WorkdirSessionCapabilities,
|
||||||
validity: Arc<SessionValidity>,
|
validity: Arc<SessionValidity>,
|
||||||
|
cleanup_pending: Arc<AtomicBool>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl std::fmt::Debug for WorkdirScopeLease {
|
impl std::fmt::Debug for WorkdirScopeLease {
|
||||||
@@ -134,12 +135,74 @@ impl WorkdirScopeLease {
|
|||||||
self.broker.scope(request).await
|
self.broker.scope(request).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn close(&self) -> Result<(), WorkdirError> {
|
||||||
|
self.validity.active.store(false, Ordering::Release);
|
||||||
|
let command_ids = self
|
||||||
|
.broker
|
||||||
|
.authority
|
||||||
|
.owned_commands
|
||||||
|
.lock()
|
||||||
|
.expect("scoped command set mutex poisoned")
|
||||||
|
.iter()
|
||||||
|
.cloned()
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
let mut first_error = None;
|
||||||
|
for command_id in command_ids {
|
||||||
|
let handle = CommandHandle(command_id.clone());
|
||||||
|
let cancel = self
|
||||||
|
.broker
|
||||||
|
.authority
|
||||||
|
.source
|
||||||
|
.cancel_command(handle.clone())
|
||||||
|
.await;
|
||||||
|
let terminal = self
|
||||||
|
.broker
|
||||||
|
.authority
|
||||||
|
.source
|
||||||
|
.command_output(CommandOutputRequest {
|
||||||
|
handle,
|
||||||
|
cursor: 0,
|
||||||
|
limit: 1,
|
||||||
|
wait: true,
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
match (cancel, terminal) {
|
||||||
|
(_, Ok(_))
|
||||||
|
| (Ok(()), Err(WorkdirError::UnknownCommand(_)))
|
||||||
|
| (Err(WorkdirError::UnknownCommand(_)), Err(WorkdirError::UnknownCommand(_))) => {
|
||||||
|
self.broker
|
||||||
|
.authority
|
||||||
|
.owned_commands
|
||||||
|
.lock()
|
||||||
|
.expect("scoped command set mutex poisoned")
|
||||||
|
.remove(&command_id);
|
||||||
|
}
|
||||||
|
(Err(error), _) | (_, Err(error)) => {
|
||||||
|
first_error.get_or_insert(error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if let Some(error) = first_error {
|
||||||
|
return Err(error);
|
||||||
|
}
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
self.finish_release();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
pub fn is_active(&self) -> bool {
|
pub fn is_active(&self) -> bool {
|
||||||
self.validity.is_active()
|
self.validity.is_active()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn release(&self) {
|
/// Revoke a scope whose owner has already terminalized every tool call.
|
||||||
|
/// Use [`Self::close`] when commands may still be live.
|
||||||
|
pub fn revoke(&self) {
|
||||||
|
self.finish_release();
|
||||||
|
}
|
||||||
|
|
||||||
|
fn finish_release(&self) {
|
||||||
self.validity.active.store(false, Ordering::Release);
|
self.validity.active.store(false, Ordering::Release);
|
||||||
|
self.cleanup_pending.store(false, Ordering::Release);
|
||||||
if let Some(forwarder) = &self.broker.event_forwarder
|
if let Some(forwarder) = &self.broker.event_forwarder
|
||||||
&& let Some(handle) = forwarder
|
&& let Some(handle) = forwarder
|
||||||
.lock()
|
.lock()
|
||||||
@@ -161,7 +224,7 @@ impl std::ops::Deref for WorkdirScopeLease {
|
|||||||
|
|
||||||
impl Drop for WorkdirScopeLease {
|
impl Drop for WorkdirScopeLease {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
self.release();
|
self.finish_release();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -195,6 +258,7 @@ impl SessionValidity {
|
|||||||
#[derive(Clone, Debug)]
|
#[derive(Clone, Debug)]
|
||||||
struct ActiveWriteLease {
|
struct ActiveWriteLease {
|
||||||
validity: Weak<SessionValidity>,
|
validity: Weak<SessionValidity>,
|
||||||
|
cleanup_pending: Weak<AtomicBool>,
|
||||||
rules: Vec<WorkdirToolScopeRule>,
|
rules: Vec<WorkdirToolScopeRule>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -471,22 +535,52 @@ impl ScopedWorkdirSession {
|
|||||||
self.ensure_scope_targets_do_not_traverse_symlinks(&request.rules)
|
self.ensure_scope_targets_do_not_traverse_symlinks(&request.rules)
|
||||||
.await?;
|
.await?;
|
||||||
let validity = SessionValidity::child(self.validity.clone());
|
let validity = SessionValidity::child(self.validity.clone());
|
||||||
|
let cleanup_pending = Arc::new(AtomicBool::new(true));
|
||||||
let id = self.next_lease_id.fetch_add(1, Ordering::Relaxed);
|
let id = self.next_lease_id.fetch_add(1, Ordering::Relaxed);
|
||||||
if request
|
if request
|
||||||
.rules
|
.rules
|
||||||
.iter()
|
.iter()
|
||||||
.any(|rule| rule.permission == WorkdirToolScopePermission::Write)
|
.any(|rule| rule.permission == WorkdirToolScopePermission::Write)
|
||||||
{
|
{
|
||||||
self.child_write_leases
|
let mut leases = self
|
||||||
|
.child_write_leases
|
||||||
.lock()
|
.lock()
|
||||||
.expect("Workdir tool scope lease mutex poisoned")
|
.expect("Workdir tool scope lease mutex poisoned");
|
||||||
.insert(
|
leases.retain(|_, lease| {
|
||||||
id,
|
lease
|
||||||
ActiveWriteLease {
|
.validity
|
||||||
validity: Arc::downgrade(&validity),
|
.upgrade()
|
||||||
rules: request.rules.clone(),
|
.is_some_and(|validity| validity.is_active())
|
||||||
},
|
|| lease
|
||||||
);
|
.cleanup_pending
|
||||||
|
.upgrade()
|
||||||
|
.is_some_and(|pending| pending.load(Ordering::Acquire))
|
||||||
|
});
|
||||||
|
let requested_write_rules = request
|
||||||
|
.rules
|
||||||
|
.iter()
|
||||||
|
.filter(|rule| rule.permission == WorkdirToolScopePermission::Write);
|
||||||
|
for requested in requested_write_rules {
|
||||||
|
if leases.values().any(|lease| {
|
||||||
|
lease
|
||||||
|
.rules
|
||||||
|
.iter()
|
||||||
|
.any(|active| rules_overlap(active, requested))
|
||||||
|
}) {
|
||||||
|
return Err(WorkdirError::Denied(format!(
|
||||||
|
"scoped write path `{}` overlaps an active child scope",
|
||||||
|
requested.target
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
leases.insert(
|
||||||
|
id,
|
||||||
|
ActiveWriteLease {
|
||||||
|
validity: Arc::downgrade(&validity),
|
||||||
|
cleanup_pending: Arc::downgrade(&cleanup_pending),
|
||||||
|
rules: request.rules.clone(),
|
||||||
|
},
|
||||||
|
);
|
||||||
}
|
}
|
||||||
let owned_commands = Arc::new(Mutex::new(HashSet::new()));
|
let owned_commands = Arc::new(Mutex::new(HashSet::new()));
|
||||||
let (command_events, _) = broadcast::channel(64);
|
let (command_events, _) = broadcast::channel(64);
|
||||||
@@ -517,6 +611,7 @@ impl ScopedWorkdirSession {
|
|||||||
broker,
|
broker,
|
||||||
capabilities,
|
capabilities,
|
||||||
validity,
|
validity,
|
||||||
|
cleanup_pending,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -787,6 +882,13 @@ fn unix_timestamp_ms() -> u64 {
|
|||||||
.min(u128::from(u64::MAX)) as u64
|
.min(u128::from(u64::MAX)) as u64
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn rules_overlap(left: &WorkdirToolScopeRule, right: &WorkdirToolScopeRule) -> bool {
|
||||||
|
left.permission == WorkdirToolScopePermission::Write
|
||||||
|
&& right.permission == WorkdirToolScopePermission::Write
|
||||||
|
&& (rule_allows_path(left, &right.target, WorkdirToolScopePermission::Write)
|
||||||
|
|| rule_allows_path(right, &left.target, WorkdirToolScopePermission::Write))
|
||||||
|
}
|
||||||
|
|
||||||
fn rule_allows_path(
|
fn rule_allows_path(
|
||||||
rule: &WorkdirToolScopeRule,
|
rule: &WorkdirToolScopeRule,
|
||||||
path: &FsPath,
|
path: &FsPath,
|
||||||
@@ -1231,7 +1333,7 @@ mod tests {
|
|||||||
));
|
));
|
||||||
parent.write(write("other/file", "parent")).await.unwrap();
|
parent.write(write("other/file", "parent")).await.unwrap();
|
||||||
child.write(write("file", "child")).await.unwrap();
|
child.write(write("file", "child")).await.unwrap();
|
||||||
child.release();
|
child.close().await.unwrap();
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
child
|
child
|
||||||
.start_command(CommandRequest {
|
.start_command(CommandRequest {
|
||||||
@@ -1255,6 +1357,79 @@ mod tests {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn sibling_write_scopes_must_not_overlap() {
|
||||||
|
let root = TempDir::new().unwrap();
|
||||||
|
fs::create_dir_all(root.path().join("shared/one")).unwrap();
|
||||||
|
fs::create_dir_all(root.path().join("other")).unwrap();
|
||||||
|
let parent = session(root.path());
|
||||||
|
let first = parent
|
||||||
|
.scope(request("shared", WorkdirToolScopePermission::Write))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert!(matches!(
|
||||||
|
parent
|
||||||
|
.scope(request("shared/one", WorkdirToolScopePermission::Write))
|
||||||
|
.await,
|
||||||
|
Err(WorkdirError::Denied(_))
|
||||||
|
));
|
||||||
|
let other = parent
|
||||||
|
.scope(request("other", WorkdirToolScopePermission::Write))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
other.close().await.unwrap();
|
||||||
|
first.close().await.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn closing_scope_cancels_and_terminalizes_owned_commands() {
|
||||||
|
let root = TempDir::new().unwrap();
|
||||||
|
fs::create_dir_all(root.path().join("work")).unwrap();
|
||||||
|
let parent = session(root.path());
|
||||||
|
let child = parent
|
||||||
|
.scope(request("work", WorkdirToolScopePermission::Write))
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let mut events = child.subscribe_command_events().unwrap();
|
||||||
|
let handle = child
|
||||||
|
.start_command(CommandRequest {
|
||||||
|
command: "sleep 30; printf leaked > marker".into(),
|
||||||
|
timeout_secs: 60,
|
||||||
|
output_limit: 1024,
|
||||||
|
cwd: None,
|
||||||
|
spill_dir: None,
|
||||||
|
tool_call_id: Some("owned-command".into()),
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(matches!(
|
||||||
|
events.recv().await.unwrap(),
|
||||||
|
CommandEvent::Started { .. }
|
||||||
|
));
|
||||||
|
|
||||||
|
child.close().await.unwrap();
|
||||||
|
|
||||||
|
assert!(matches!(
|
||||||
|
parent.command_status(handle).await,
|
||||||
|
Ok(CommandStatus::Cancelled | CommandStatus::Completed | CommandStatus::Failed)
|
||||||
|
| Err(WorkdirError::UnknownCommand(_))
|
||||||
|
));
|
||||||
|
assert!(!root.path().join("work/marker").exists());
|
||||||
|
let terminal = tokio::time::timeout(std::time::Duration::from_secs(1), async {
|
||||||
|
loop {
|
||||||
|
if let CommandEvent::Terminal { .. } = events.recv().await.unwrap() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
assert!(
|
||||||
|
terminal.is_ok(),
|
||||||
|
"scope close must publish terminal command telemetry"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn nested_delegation_is_attenuated_and_parent_revocation_cascades() {
|
async fn nested_delegation_is_attenuated_and_parent_revocation_cascades() {
|
||||||
let root = TempDir::new().unwrap();
|
let root = TempDir::new().unwrap();
|
||||||
@@ -1286,7 +1461,7 @@ mod tests {
|
|||||||
.is_err()
|
.is_err()
|
||||||
);
|
);
|
||||||
|
|
||||||
child.release();
|
child.close().await.unwrap();
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
nested.read(read("a")).await,
|
nested.read(read("a")).await,
|
||||||
Err(WorkdirError::SessionClosed)
|
Err(WorkdirError::SessionClosed)
|
||||||
@@ -1332,8 +1507,8 @@ mod tests {
|
|||||||
));
|
));
|
||||||
nested.write(write("nested", "allowed")).await.unwrap();
|
nested.write(write("nested", "allowed")).await.unwrap();
|
||||||
|
|
||||||
nested.release();
|
nested.close().await.unwrap();
|
||||||
child.release();
|
child.close().await.unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
@@ -690,7 +690,7 @@ impl SpawnedWorkerRegistry {
|
|||||||
if !record.claim_scope_reclaim() {
|
if !record.claim_scope_reclaim() {
|
||||||
return Ok(false);
|
return Ok(false);
|
||||||
}
|
}
|
||||||
record.workdir_tool_scope.release();
|
record.workdir_tool_scope.revoke();
|
||||||
let result = if let Some(parent_scope) = &self.parent_scope {
|
let result = if let Some(parent_scope) = &self.parent_scope {
|
||||||
parent_scope
|
parent_scope
|
||||||
.update(|current| current.with_removed_deny_rules(delegated_write_rules(record)))
|
.update(|current| current.with_removed_deny_rules(delegated_write_rules(record)))
|
||||||
@@ -731,6 +731,11 @@ impl SpawnedWorkerRegistry {
|
|||||||
.stop()
|
.stop()
|
||||||
.await
|
.await
|
||||||
.map_err(|error| io::Error::other(error.to_string()))?;
|
.map_err(|error| io::Error::other(error.to_string()))?;
|
||||||
|
record
|
||||||
|
.workdir_tool_scope
|
||||||
|
.close()
|
||||||
|
.await
|
||||||
|
.map_err(|error| io::Error::other(error.to_string()))?;
|
||||||
let summary = record.stop_summary();
|
let summary = record.stop_summary();
|
||||||
self.reclaim_record_scope(&record)?;
|
self.reclaim_record_scope(&record)?;
|
||||||
let removed =
|
let removed =
|
||||||
|
|||||||
Reference in New Issue
Block a user