runtime: canonicalize Worker aggregates

This commit is contained in:
2026-08-11 22:38:29 +09:00
parent 33d3d1a43f
commit 34f8949e85
16 changed files with 2013 additions and 201 deletions
+4
View File
@@ -245,6 +245,8 @@ impl fmt::Debug for WorkerExecutionContext {
#[derive(Clone, Debug)]
pub struct WorkerExecutionSpawnRequest {
pub worker_ref: WorkerRef,
/// Monotonic execution generation reserved durably before launch.
pub run_generation: u64,
pub request: crate::catalog::CreateWorkerRequest,
pub context: WorkerExecutionContext,
pub working_directory: Option<WorkingDirectoryBinding>,
@@ -255,6 +257,8 @@ pub struct WorkerExecutionSpawnRequest {
#[derive(Clone, Debug)]
pub struct WorkerExecutionRestoreRequest {
pub worker_ref: WorkerRef,
/// Monotonic execution generation reserved durably before restore.
pub run_generation: u64,
pub request: crate::catalog::CreateWorkerRequest,
pub context: WorkerExecutionContext,
pub previous_working_directory: Option<WorkingDirectoryStatus>,
File diff suppressed because it is too large Load Diff
+13 -3
View File
@@ -18,7 +18,7 @@ use worker_runtime::auth::{
RuntimeHttpAuthConfig, RuntimeIdentityMaterial, TrustedServerKey, decode_public_key,
};
use worker_runtime::error::RuntimeError;
use worker_runtime::fs_store::FsRuntimeStoreOptions;
use worker_runtime::fs_store::{FsRuntimeStore, FsRuntimeStoreOptions};
use worker_runtime::http_server::{
RuntimeHttpServerConfig, RuntimeHttpServerError, RuntimeHttpStoreSelection,
};
@@ -80,9 +80,13 @@ fn run() -> Result<(), ProcessError> {
fn build_runtime(config: &ProcessConfig) -> Result<Runtime, ProcessError> {
let fs_paths = config.resolved_fs_paths();
let runtime_store_dir = match &config.http.store {
RuntimeHttpStoreSelection::Memory => fs_paths.runtime_dir.clone(),
RuntimeHttpStoreSelection::Fs { root } => root.clone(),
_ => fs_paths.runtime_dir.clone(),
};
let mut factory = ProfileRuntimeWorkerFactory::new(fs_paths.worker_dir.join("worker-root"))
.with_store_dir(fs_paths.worker_dir.join("sessions"))
.with_worker_metadata_dir(fs_paths.worker_dir.join("metadata"));
.with_runtime_store_dir(runtime_store_dir);
if let Some(endpoint) = config.backend_resource_endpoint.clone() {
factory = factory.with_resource_client(Arc::new(
worker_runtime::resource::HttpBackendResourceClient::new(
@@ -105,6 +109,12 @@ fn build_runtime(config: &ProcessConfig) -> Result<Runtime, ProcessError> {
.map_err(ProcessError::Runtime)
}
RuntimeHttpStoreSelection::Fs { root } => {
FsRuntimeStore::migrate_legacy_worker_aggregates(
root,
fs_paths.worker_dir.join("sessions"),
fs_paths.worker_dir.join("metadata"),
)
.map_err(ProcessError::Runtime)?;
let mut options = FsRuntimeStoreOptions::new(root.clone());
options.display_name = config.http.display_name.clone();
Runtime::with_fs_store_and_execution_backend(options, backend)
+88 -29
View File
@@ -532,12 +532,16 @@ impl Runtime {
status: WorkerStatus::Stopped,
workspace_id: scope.map(|scope| scope.workspace_id.clone()),
request: request.clone(),
run_generation: 1,
working_directory: None,
execution_handle: None,
};
state.workers.insert(worker_id, record);
state.persist_runtime_snapshot()?;
state.persist_worker(&worker_ref.worker_id)?;
let spawn_request = WorkerExecutionSpawnRequest {
worker_ref: worker_ref.clone(),
run_generation: 1,
request,
context: self.execution_context(worker_ref.clone()),
working_directory: None,
@@ -859,35 +863,46 @@ impl Runtime {
/// present this is idempotent; otherwise the configured backend is tried.
pub fn restore_worker(&self, worker_ref: &WorkerRef) -> Result<WorkerDetail, RuntimeError> {
let (backend, request) = {
let state = self.lock()?;
let mut state = self.lock()?;
state.ensure_running()?;
let worker = state.worker(worker_ref)?;
if worker.execution_handle.is_some() {
return Ok(worker.detail());
}
if worker.status == WorkerStatus::Cancelled {
return Err(RuntimeError::InvalidRequest(format!(
"worker {} is cancelled",
worker_ref.worker_id
)));
}
let (worker_request, previous_working_directory, config_bundle, run_generation) = {
let worker = state.worker(worker_ref)?;
if worker.execution_handle.is_some() {
return Ok(worker.detail());
}
if worker.status == WorkerStatus::Cancelled {
return Err(RuntimeError::InvalidRequest(format!(
"worker {} is cancelled",
worker_ref.worker_id
)));
}
let config_bundle = worker
.request
.config_bundle
.as_ref()
.and_then(|bundle_ref| state.config_bundles.get(&bundle_ref.id))
.cloned();
(
worker.request.clone(),
worker.working_directory.clone(),
config_bundle,
worker.run_generation.saturating_add(1).max(1),
)
};
let backend = state.execution_backend.clone().ok_or_else(|| {
RuntimeError::WorkerExecutionUnavailable {
worker_id: worker_ref.worker_id.clone(),
message: "runtime has no execution backend".to_string(),
}
})?;
let config_bundle = worker
.request
.config_bundle
.as_ref()
.and_then(|bundle_ref| state.config_bundles.get(&bundle_ref.id))
.cloned();
state.worker_mut(worker_ref)?.run_generation = run_generation;
state.persist_worker(&worker_ref.worker_id)?;
let request = WorkerExecutionRestoreRequest {
worker_ref: worker_ref.clone(),
request: worker.request.clone(),
run_generation,
request: worker_request,
context: self.execution_context(worker_ref.clone()),
previous_working_directory: worker.working_directory.clone(),
previous_working_directory,
working_directory: None,
config_bundle,
};
@@ -1514,34 +1529,64 @@ impl Runtime {
struct RestoreCandidate {
worker_ref: WorkerRef,
request: CreateWorkerRequest,
run_generation: u64,
previous_working_directory: Option<CatalogWorkingDirectoryStatus>,
config_bundle: Option<ConfigBundle>,
}
let candidates = {
let state = self.lock()?;
let mut state = self.lock()?;
if state.execution_backend.is_none() {
return Ok(());
}
state
let worker_ids = state
.workers
.values()
.filter(|worker| worker.execution_handle.is_none())
.map(|worker| {
.map(|worker| worker.worker_id)
.collect::<Vec<_>>();
let mut candidates = Vec::with_capacity(worker_ids.len());
for worker_id in worker_ids {
let (
worker_ref,
request,
previous_working_directory,
config_bundle,
run_generation,
) = {
let worker = state
.workers
.get(&worker_id)
.expect("collected Worker exists");
let config_bundle = worker
.request
.config_bundle
.as_ref()
.and_then(|bundle_ref| state.config_bundles.get(&bundle_ref.id))
.cloned();
RestoreCandidate {
worker_ref: worker.worker_ref.clone(),
request: worker.request.clone(),
previous_working_directory: worker.working_directory.clone(),
(
worker.worker_ref.clone(),
worker.request.clone(),
worker.working_directory.clone(),
config_bundle,
}
})
.collect::<Vec<_>>()
worker.run_generation.saturating_add(1).max(1),
)
};
state
.workers
.get_mut(&worker_id)
.expect("collected Worker exists")
.run_generation = run_generation;
state.persist_worker(&worker_id)?;
candidates.push(RestoreCandidate {
worker_ref,
request,
run_generation,
previous_working_directory,
config_bundle,
});
}
candidates
};
for candidate in candidates {
@@ -1554,6 +1599,7 @@ impl Runtime {
};
let request = WorkerExecutionRestoreRequest {
worker_ref: candidate.worker_ref.clone(),
run_generation: candidate.run_generation,
request: candidate.request,
context: self.execution_context(candidate.worker_ref.clone()),
previous_working_directory: candidate.previous_working_directory,
@@ -1723,6 +1769,7 @@ impl RuntimeState {
status: WorkerStatus::Stopped,
workspace_id: worker.workspace_id,
request: worker.request,
run_generation: worker.run_generation,
working_directory: worker.working_directory,
execution_handle: None,
},
@@ -2281,6 +2328,7 @@ struct WorkerRecord {
status: WorkerStatus,
workspace_id: Option<String>,
request: CreateWorkerRequest,
run_generation: u64,
working_directory: Option<CatalogWorkingDirectoryStatus>,
execution_handle: Option<WorkerExecutionHandle>,
}
@@ -2324,6 +2372,7 @@ impl WorkerRecord {
worker_ref: self.worker_ref.clone(),
worker_id: self.worker_id.clone(),
request: self.request.clone(),
run_generation: self.run_generation,
workspace_id: self.workspace_id.clone(),
working_directory: self.working_directory.clone(),
}
@@ -2627,6 +2676,7 @@ mod tests {
dispatch_result: Mutex<Option<WorkerExecutionResult>>,
restore_result: Mutex<Option<WorkerExecutionSpawnResult>>,
restore_count: Mutex<u64>,
run_generations: Mutex<Vec<u64>>,
contexts: Mutex<BTreeMap<WorkerId, WorkerExecutionContext>>,
dispatched_inputs: Mutex<Vec<WorkerInput>>,
preserve_commit_ack_submission_id: AtomicBool,
@@ -2670,6 +2720,10 @@ mod tests {
}
fn spawn_worker(&self, request: WorkerExecutionSpawnRequest) -> WorkerExecutionSpawnResult {
self.run_generations
.lock()
.unwrap()
.push(request.run_generation);
self.contexts
.lock()
.unwrap()
@@ -2689,6 +2743,10 @@ mod tests {
request: WorkerExecutionRestoreRequest,
) -> WorkerExecutionSpawnResult {
*self.restore_count.lock().unwrap() += 1;
self.run_generations
.lock()
.unwrap()
.push(request.run_generation);
if let Some(result) = self.restore_result.lock().unwrap().clone() {
return result;
}
@@ -3551,6 +3609,7 @@ mod tests {
.unwrap();
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
assert_eq!(*backend.run_generations.lock().unwrap(), vec![1, 2]);
assert_eq!(
runtime.worker_detail(&detail.worker_ref).unwrap().status,
WorkerStatus::Idle
+118 -126
View File
@@ -10,7 +10,7 @@
use std::collections::HashMap;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, mpsc};
use std::time::Duration;
@@ -30,9 +30,12 @@ use crate::working_directory::{
WorkingDirectoryBinding, WorkingDirectoryDiagnostic, WorkingDirectoryMaterializer,
};
use async_trait::async_trait;
use manifest::paths;
use protocol::{Event, Method, Segment, WorkerStatus};
use session_store::{CombinedStore, FsStore, FsWorkerStore, LogEntry, collect_state};
use session_store::{
CombinedStore, LogEntry, WorkerAggregateStore, WorkerSessionStore, collect_state,
};
#[cfg(test)]
use session_store::{FsStore, FsWorkerStore};
use tokio::runtime::Runtime;
#[cfg(feature = "ws-server")]
use tokio::sync::broadcast;
@@ -57,7 +60,6 @@ const RUNTIME_TASK_TIMEOUT: Duration = Duration::from_secs(10);
// Keep this below the adapter task timeout so a failed acknowledgement task
// returns a typed execution error instead of leaving the outer waiter to time out.
const USER_INPUT_COMMIT_TIMEOUT: Duration = Duration::from_secs(9);
static NEXT_RUNTIME_ARTIFACT_ROOT: AtomicU64 = AtomicU64::new(1);
fn user_input_has_submission(entry: &LogEntry, submission_id: &str) -> bool {
let LogEntry::UserInput { extensions, .. } = entry else {
@@ -69,43 +71,9 @@ fn user_input_has_submission(entry: &LogEntry, submission_id: &str) -> bool {
})
}
#[derive(Clone)]
enum RuntimeArtifactRoot {
Owned(Arc<OwnedRuntimeArtifactRoot>),
External(PathBuf),
}
impl RuntimeArtifactRoot {
fn owned() -> Self {
let sequence = NEXT_RUNTIME_ARTIFACT_ROOT.fetch_add(1, Ordering::Relaxed);
Self::Owned(Arc::new(OwnedRuntimeArtifactRoot {
path: std::env::temp_dir().join(format!(
"yoi-worker-runtime-artifacts-{}-{sequence}",
std::process::id()
)),
}))
}
fn path(&self) -> &std::path::Path {
match self {
Self::Owned(root) => &root.path,
Self::External(path) => path,
}
}
}
struct OwnedRuntimeArtifactRoot {
path: PathBuf,
}
impl Drop for OwnedRuntimeArtifactRoot {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
pub struct RuntimeWorkerController {
pub handle: WorkerHandle,
pub shutdown: Arc<tokio::sync::Mutex<Option<worker::ShutdownReceiver>>>,
pub workspace_client: Arc<dyn WorkspaceClient>,
}
@@ -244,9 +212,7 @@ impl WorkerObservationProvider for RuntimeGrantedWorkerObservationProvider {
pub struct ProfileRuntimeWorkerFactory {
observation_hub: Arc<RuntimeWorkerObservationHub>,
profile_base_dir: PathBuf,
store_dir: Option<PathBuf>,
worker_metadata_dir: Option<PathBuf>,
runtime_base_dir: RuntimeArtifactRoot,
worker_aggregate_root: Option<PathBuf>,
resource_client: Option<Arc<dyn BackendResourceClient>>,
profile_archive_cache: Arc<ProfileSourceArchiveCache>,
}
@@ -257,26 +223,14 @@ impl ProfileRuntimeWorkerFactory {
Self {
observation_hub: Arc::new(RuntimeWorkerObservationHub::default()),
profile_base_dir,
store_dir: None,
worker_metadata_dir: None,
runtime_base_dir: RuntimeArtifactRoot::owned(),
worker_aggregate_root: None,
resource_client: None,
profile_archive_cache: Arc::new(ProfileSourceArchiveCache::default()),
}
}
pub fn with_store_dir(mut self, store_dir: impl Into<PathBuf>) -> Self {
self.store_dir = Some(store_dir.into());
self
}
pub fn with_worker_metadata_dir(mut self, worker_metadata_dir: impl Into<PathBuf>) -> Self {
self.worker_metadata_dir = Some(worker_metadata_dir.into());
self
}
pub fn with_runtime_base_dir(mut self, runtime_base_dir: impl Into<PathBuf>) -> Self {
self.runtime_base_dir = RuntimeArtifactRoot::External(runtime_base_dir.into());
pub fn with_runtime_store_dir(mut self, runtime_store_dir: impl Into<PathBuf>) -> Self {
self.worker_aggregate_root = Some(runtime_store_dir.into().join("workers"));
self
}
@@ -285,28 +239,16 @@ impl ProfileRuntimeWorkerFactory {
self
}
fn store_dir(&self) -> Result<PathBuf, String> {
self.store_dir
.clone()
.or_else(paths::sessions_dir)
fn worker_aggregate_dir(&self, worker_ref: &WorkerRef) -> Result<PathBuf, String> {
self.worker_aggregate_root
.as_ref()
.map(|root| root.join(worker_ref.worker_id.to_string()))
.ok_or_else(|| {
"could not resolve sessions directory (set YOI_DATA_DIR, YOI_HOME, XDG_DATA_HOME, or HOME)"
"Runtime Worker aggregate root is not configured; global Session/metadata roots are migration-only"
.to_string()
})
}
fn worker_metadata_dir(&self, store_dir: &std::path::Path) -> PathBuf {
self.worker_metadata_dir
.clone()
.or_else(|| paths::data_dir().map(|data_dir| data_dir.join("workers")))
.or_else(|| store_dir.parent().map(|parent| parent.join("workers")))
.unwrap_or_else(|| PathBuf::from("workers"))
}
fn runtime_base_dir(&self) -> Result<PathBuf, String> {
Ok(self.runtime_base_dir.path().to_path_buf())
}
fn runtime_worker_name_for_ref(worker_ref: &crate::identity::WorkerRef) -> String {
format!("worker-runtime-{}", worker_ref.worker_id)
}
@@ -575,20 +517,23 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
};
let flow_transition_enabled = manifest.feature.flow.enabled;
let store_dir = self.store_dir()?;
let session_store = FsStore::new(&store_dir).map_err(|err| {
let worker_aggregate_dir = self.worker_aggregate_dir(&request.worker_ref)?;
let session_dir = worker_aggregate_dir.join("session");
let session_store = WorkerSessionStore::new(&session_dir).map_err(|err| {
format!(
"failed to initialize session store at {}: {err}",
store_dir.display()
)
})?;
let worker_metadata_dir = self.worker_metadata_dir(&store_dir);
let worker_metadata_store = FsWorkerStore::new(&worker_metadata_dir).map_err(|err| {
format!(
"failed to initialize worker metadata store at {}: {err}",
worker_metadata_dir.display()
"failed to initialize canonical Worker Session store at {}: {err}",
session_dir.display()
)
})?;
let worker_metadata_store =
WorkerAggregateStore::new(&worker_aggregate_dir, worker_name.clone()).map_err(
|err| {
format!(
"failed to initialize canonical Worker metadata store at {}: {err}",
worker_aggregate_dir.display()
)
},
)?;
let store = CombinedStore::new(session_store, worker_metadata_store);
let mut worker = Worker::from_manifest_with_context(
@@ -642,10 +587,17 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
}
let workspace_client = worker.workspace_client_handle();
let runtime_base = self.runtime_base_dir()?;
let (handle, _shutdown_rx) = WorkerController::spawn_runtime_managed(worker, &runtime_base)
let run_dir = worker_aggregate_dir
.join("runs")
.join(request.run_generation.to_string());
let (handle, shutdown_rx) = WorkerController::spawn_runtime_managed_run(worker, &run_dir)
.await
.map_err(|err| format!("failed to spawn Worker controller: {err}"))?;
.map_err(|err| {
format!(
"failed to spawn Worker controller in {}: {err}",
run_dir.display()
)
})?;
if flow_transition_enabled {
handle.shared_state.enable_flow_transition();
}
@@ -656,6 +608,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
);
Ok(RuntimeWorkerController {
handle,
shutdown: Arc::new(tokio::sync::Mutex::new(Some(shutdown_rx))),
workspace_client,
})
}
@@ -692,26 +645,29 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
let workspace_context = workspace_backend_ref.worker_context(&request.worker_ref);
let (manifest, loader) = Self::restore_fallback_manifest(&worker_name)?;
let store_dir = self.store_dir()?;
let session_store = FsStore::new(&store_dir).map_err(|err| {
let worker_aggregate_dir = self.worker_aggregate_dir(&request.worker_ref)?;
let session_dir = worker_aggregate_dir.join("session");
let session_store = WorkerSessionStore::new(&session_dir).map_err(|err| {
format!(
"failed to initialize session store at {}: {err}",
store_dir.display()
)
})?;
let worker_metadata_dir = self.worker_metadata_dir(&store_dir);
let worker_metadata_store = FsWorkerStore::new(&worker_metadata_dir).map_err(|err| {
format!(
"failed to initialize worker metadata store at {}: {err}",
worker_metadata_dir.display()
"failed to initialize canonical Worker Session store at {}: {err}",
session_dir.display()
)
})?;
let worker_metadata_store =
WorkerAggregateStore::new(&worker_aggregate_dir, worker_name.clone()).map_err(
|err| {
format!(
"failed to initialize canonical Worker metadata store at {}: {err}",
worker_aggregate_dir.display()
)
},
)?;
let store = CombinedStore::new(session_store, worker_metadata_store);
let mut worker = match Worker::restore_from_worker_metadata_with_context(
&worker_name,
manifest.clone(),
store,
store.clone(),
loader.clone(),
workspace_context.clone(),
filesystem_authority.clone(),
@@ -722,20 +678,6 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
Err(WorkerError::WorkerMetadataPending { .. })
if request.request.initial_input.is_none() =>
{
let session_store = FsStore::new(&store_dir).map_err(|err| {
format!(
"failed to initialize session store at {}: {err}",
store_dir.display()
)
})?;
let worker_metadata_store =
FsWorkerStore::new(&worker_metadata_dir).map_err(|err| {
format!(
"failed to initialize worker metadata store at {}: {err}",
worker_metadata_dir.display()
)
})?;
let store = CombinedStore::new(session_store, worker_metadata_store);
Worker::restore_pending_from_worker_metadata_with_context(
&worker_name,
manifest.clone(),
@@ -792,10 +734,17 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
}
let workspace_client = worker.workspace_client_handle();
let runtime_base = self.runtime_base_dir()?;
let (handle, _shutdown_rx) = WorkerController::spawn_runtime_managed(worker, &runtime_base)
let run_dir = worker_aggregate_dir
.join("runs")
.join(request.run_generation.to_string());
let (handle, shutdown_rx) = WorkerController::spawn_runtime_managed_run(worker, &run_dir)
.await
.map_err(|err| format!("failed to spawn restored Worker controller: {err}"))?;
.map_err(|err| {
format!(
"failed to spawn restored Worker controller in {}: {err}",
run_dir.display()
)
})?;
if flow_transition_enabled {
handle.shared_state.enable_flow_transition();
}
@@ -806,6 +755,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
);
Ok(RuntimeWorkerController {
handle,
shutdown: Arc::new(tokio::sync::Mutex::new(Some(shutdown_rx))),
workspace_client,
})
}
@@ -813,6 +763,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
struct RuntimeWorkerExecution {
handle: WorkerHandle,
shutdown: Arc<tokio::sync::Mutex<Option<worker::ShutdownReceiver>>>,
busy: Arc<AtomicBool>,
workspace_client: Option<Arc<dyn WorkspaceClient>>,
}
@@ -828,7 +779,10 @@ pub struct WorkerRuntimeExecutionBackend<F = ProfileRuntimeWorkerFactory> {
impl WorkerRuntimeExecutionBackend<ProfileRuntimeWorkerFactory> {
pub fn from_workspace(workspace_root: impl Into<PathBuf>) -> Result<Self, String> {
Self::new(ProfileRuntimeWorkerFactory::new(workspace_root))
let workspace_root = workspace_root.into();
let factory = ProfileRuntimeWorkerFactory::new(&workspace_root)
.with_runtime_store_dir(workspace_root.join(".yoi/runtime-store"));
Self::new(factory)
}
}
@@ -1094,6 +1048,7 @@ where
worker_ref: crate::identity::WorkerRef,
bridge_context: crate::execution::WorkerExecutionContext,
handle: WorkerHandle,
shutdown: Arc<tokio::sync::Mutex<Option<worker::ShutdownReceiver>>>,
working_directory: Option<WorkingDirectoryBinding>,
workspace_client: Option<Arc<dyn WorkspaceClient>>,
) -> WorkerExecutionSpawnResult {
@@ -1157,6 +1112,7 @@ where
worker_ref.clone(),
RuntimeWorkerExecution {
handle,
shutdown,
busy,
workspace_client,
},
@@ -1393,6 +1349,7 @@ where
worker_ref,
bridge_context,
controller.handle,
controller.shutdown,
working_directory,
Some(controller.workspace_client),
)
@@ -1489,6 +1446,7 @@ where
worker_ref,
bridge_context,
controller.handle,
controller.shutdown,
working_directory,
Some(controller.workspace_client),
)
@@ -1698,12 +1656,28 @@ where
"execution handle does not reference a live Worker",
);
};
self.send_method(
let shutdown = execution.shutdown.clone();
let result = self.send_method(
WorkerExecutionOperation::Stop,
execution.handle,
Method::Shutdown,
WorkerExecutionRunState::Stopped,
)
);
if result.outcome != crate::execution::WorkerExecutionOutcome::Accepted {
return result;
}
match self.run_on_adapter_runtime(async move {
let receiver = shutdown.lock().await.take();
if let Some(receiver) = receiver {
receiver
.await
.map_err(|_| "Worker shutdown completion channel closed".to_string())?;
}
Ok(())
}) {
Ok(()) => result,
Err(message) => WorkerExecutionResult::errored(WorkerExecutionOperation::Stop, message),
}
}
fn cancel_worker(&self, handle: &WorkerExecutionHandle) -> WorkerExecutionResult {
@@ -1924,12 +1898,13 @@ mod tests {
)
.await
.map_err(|err| err.to_string())?;
let (handle, _shutdown_rx) =
let (handle, shutdown_rx) =
WorkerController::spawn_runtime_managed(worker, &self.runtime_base)
.await
.map_err(|err| err.to_string())?;
Ok(RuntimeWorkerController {
handle,
shutdown: Arc::new(tokio::sync::Mutex::new(Some(shutdown_rx))),
workspace_client,
})
}
@@ -1939,6 +1914,7 @@ mod tests {
) -> Result<RuntimeWorkerController, String> {
let request = WorkerExecutionSpawnRequest {
worker_ref: request.worker_ref,
run_generation: request.run_generation,
request: request.request,
context: request.context,
working_directory: request.working_directory,
@@ -2210,6 +2186,7 @@ mod tests {
let worker_ref = crate::identity::WorkerRef::new(crate::identity::WorkerId::new(1));
let request = WorkerExecutionSpawnRequest {
worker_ref: worker_ref.clone(),
run_generation: 1,
request: create_request("1"),
context: test_execution_context(worker_ref),
working_directory: None,
@@ -2263,9 +2240,9 @@ mod tests {
#[tokio::test]
async fn restore_pending_worker_uses_saved_manifest_snapshot() {
let root = tempfile::tempdir().unwrap();
let store_dir = root.path().join("sessions");
let worker_metadata_dir = root.path().join("workers");
let runtime_store_dir = root.path().join("runtime");
let worker_ref = WorkerRef::new(crate::identity::WorkerId::new(1));
let worker_aggregate_dir = runtime_store_dir.join("workers/1");
let worker_name = ProfileRuntimeWorkerFactory::runtime_worker_name_for_ref(&worker_ref);
let session_id = session_store::new_session_id();
let manifest = manifest::WorkerManifest::from_toml(&format!(
@@ -2294,7 +2271,7 @@ mod tests {
root.path().display(),
))
.unwrap();
FsWorkerStore::new(&worker_metadata_dir)
WorkerAggregateStore::new(&worker_aggregate_dir, &worker_name)
.unwrap()
.set_active(
&worker_name,
@@ -2312,10 +2289,10 @@ mod tests {
runtime_id: Some("runtime-restore".to_string()),
});
let controller = ProfileRuntimeWorkerFactory::new(root.path())
.with_store_dir(&store_dir)
.with_worker_metadata_dir(&worker_metadata_dir)
.with_runtime_store_dir(&runtime_store_dir)
.restore_controller(WorkerExecutionRestoreRequest {
worker_ref: worker_ref.clone(),
run_generation: 1,
request,
context: test_execution_context(worker_ref),
previous_working_directory: None,
@@ -2325,8 +2302,23 @@ mod tests {
.await
.expect("pending restore should use the saved manifest snapshot");
assert!(controller.handle.shared_state.flow_transition_enabled());
let run_dir = runtime_store_dir.join("workers/1/runs/1");
assert!(run_dir.join("worker.sock").exists());
assert!(run_dir.join("worker.out.log").is_file());
assert!(run_dir.join("worker.err.log").is_file());
assert!(run_dir.join("artifacts").is_dir());
assert!(run_dir.join("spawned").is_dir());
let shutdown = controller.shutdown.clone();
controller.handle.send(Method::Shutdown).await.unwrap();
if let Some(receiver) = shutdown.lock().await.take() {
receiver.await.unwrap();
}
assert!(
run_dir.is_dir(),
"run evidence remains until a separate retention policy disposes it"
);
assert!(!run_dir.join("worker.sock").exists());
}
#[test]