expose background environment resolution

This commit is contained in:
Sayan Sisodiya
2026-06-16 11:27:34 -07:00
parent 009a2bb93d
commit 2fb0d85dd6

View File

@@ -33,11 +33,31 @@ pub(crate) fn default_thread_environment_selections(
type SnapshotTask = Shared<BoxFuture<'static, TurnEnvironmentSnapshot>>;
/// One environment selection being resolved in the background.
#[derive(Clone)]
pub(crate) struct EnvironmentResolution {
snapshot_task: SnapshotTask,
}
impl EnvironmentResolution {
fn new(snapshot_task: SnapshotTask) -> Self {
Self { snapshot_task }
}
pub(crate) fn snapshot_if_ready(&self) -> Option<TurnEnvironmentSnapshot> {
self.snapshot_task.peek().cloned()
}
pub(crate) async fn snapshot(&self) -> TurnEnvironmentSnapshot {
self.snapshot_task.clone().await
}
}
pub(crate) struct ThreadEnvironments {
environment_manager: Arc<EnvironmentManager>,
local_shell: Shell,
shell_snapshot: ShellSnapshot,
snapshot_task: ArcSwap<SnapshotTask>,
current: ArcSwap<EnvironmentResolution>,
}
impl ThreadEnvironments {
@@ -51,22 +71,23 @@ impl ThreadEnvironments {
environment_manager,
local_shell,
shell_snapshot,
snapshot_task: ArcSwap::from_pointee(futures::future::ready(current).boxed().shared()),
current: ArcSwap::from_pointee(EnvironmentResolution::new(
futures::future::ready(current).boxed().shared(),
)),
}
}
pub(crate) fn update_selections(&self, environments: &[TurnEnvironmentSelection]) {
let previous = self
.snapshot_task
.load()
.peek()
.cloned()
.unwrap_or_default();
/// Returns the work started here so a turn can check it later without blocking startup.
pub(crate) fn update_selections(
&self,
environments: &[TurnEnvironmentSelection],
) -> EnvironmentResolution {
let previous = self.current.load().snapshot_if_ready().unwrap_or_default();
let environment_manager = Arc::clone(&self.environment_manager);
let local_shell = self.local_shell.clone();
let shell_snapshot = self.shell_snapshot.clone();
let environments = environments.to_vec();
let (snapshot_task, snapshot) = async move {
let snapshot_task = async move {
Self::resolve_snapshot(
environment_manager,
local_shell,
@@ -76,10 +97,13 @@ impl ThreadEnvironments {
)
.await
}
.remote_handle();
self.snapshot_task
.store(Arc::new(snapshot.boxed().shared()));
.boxed()
.shared();
let resolution = EnvironmentResolution::new(snapshot_task.clone());
self.current.store(Arc::new(resolution.clone()));
// Keep resolving even when no caller is currently waiting for the result.
drop(tokio::spawn(snapshot_task));
resolution
}
async fn resolve_snapshot(
@@ -176,7 +200,8 @@ impl ThreadEnvironments {
}
pub(crate) async fn snapshot(&self) -> TurnEnvironmentSnapshot {
self.snapshot_task.load_full().as_ref().clone().await
let resolution = self.current.load_full();
resolution.snapshot().await
}
pub(crate) fn environment_manager(&self) -> Arc<EnvironmentManager> {
@@ -359,6 +384,35 @@ url = "ws://127.0.0.1:8765"
);
}
#[tokio::test]
async fn environment_resolution_runs_in_background() {
let cwd = AbsolutePathBuf::current_dir().expect("cwd");
let selection = TurnEnvironmentSelection {
environment_id: LOCAL_ENVIRONMENT_ID.to_string(),
cwd: PathUri::from_abs_path(&cwd),
};
let turn_environments = ThreadEnvironments::new(
Arc::new(EnvironmentManager::default_for_tests()),
crate::shell::default_user_shell(),
ShellSnapshot::disabled(),
TurnEnvironmentSnapshot::default(),
);
let resolution = turn_environments.update_selections(std::slice::from_ref(&selection));
let snapshot = tokio::time::timeout(std::time::Duration::from_secs(5), async {
loop {
if let Some(snapshot) = resolution.snapshot_if_ready() {
break snapshot;
}
tokio::task::yield_now().await;
}
})
.await
.expect("environment resolution should complete in the background");
assert_eq!(snapshot.to_selections(), vec![selection]);
}
#[tokio::test]
async fn resolve_environment_selections_keeps_first_duplicate_id() {
let cwd = AbsolutePathBuf::current_dir().expect("cwd");