fix: route runtime workers through shared bootstrap
This commit is contained in:
@@ -46,6 +46,8 @@ use tokio::runtime::Runtime;
|
|||||||
use tokio::sync::broadcast;
|
use tokio::sync::broadcast;
|
||||||
use workdir::{LocalWorkdirSession, Workdir, WorkdirSessionCapabilities, WorkdirSessionHandle};
|
use workdir::{LocalWorkdirSession, Workdir, WorkdirSessionCapabilities, WorkdirSessionHandle};
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
use worker::WorkerController;
|
||||||
use worker::feature::builtin::{
|
use worker::feature::builtin::{
|
||||||
CompositeWorkerObservationProvider, WorkerObservationError, WorkerObservationProvider,
|
CompositeWorkerObservationProvider, WorkerObservationError, WorkerObservationProvider,
|
||||||
WorkerObservationSubject, WorkerObservationSubjectRef, WorkerSessionCapture,
|
WorkerObservationSubject, WorkerObservationSubjectRef, WorkerSessionCapture,
|
||||||
@@ -54,9 +56,10 @@ use worker::feature::builtin::{
|
|||||||
#[cfg(feature = "ws-server")]
|
#[cfg(feature = "ws-server")]
|
||||||
use worker::ipc::protocol_session::{live_log_entry_event, subscribe_worker_protocol_session};
|
use worker::ipc::protocol_session::{live_log_entry_event, subscribe_worker_protocol_session};
|
||||||
use worker::{
|
use worker::{
|
||||||
PromptCatalogSource, SegmentLogSink, WORKER_INPUT_SUBMISSION_EXTENSION_DOMAIN, Worker,
|
PreparedWorker, PromptCatalogSource, SegmentLogSink, WORKER_INPUT_SUBMISSION_EXTENSION_DOMAIN,
|
||||||
WorkerController, WorkerControllerTransport, WorkerError, WorkerFilesystemAuthority,
|
Worker, WorkerBootstrap, WorkerBootstrapError, WorkerBootstrapLayout,
|
||||||
WorkerHandle, WorkerSharedState, WorkerWorkspaceContext, WorkspaceClient, WorkspaceId,
|
WorkerControllerTransport, WorkerError, WorkerFilesystemAuthority, WorkerHandle,
|
||||||
|
WorkerSharedState, WorkerWorkspaceContext, WorkspaceClient, WorkspaceId,
|
||||||
};
|
};
|
||||||
|
|
||||||
const DEFAULT_BACKEND_ID: &str = "worker-crate";
|
const DEFAULT_BACKEND_ID: &str = "worker-crate";
|
||||||
@@ -881,15 +884,31 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
)?;
|
)?;
|
||||||
let store = CombinedStore::new(session_store, worker_metadata_store);
|
let store = CombinedStore::new(session_store, worker_metadata_store);
|
||||||
|
|
||||||
let mut worker = Worker::from_manifest_with_context(
|
let run_dir = worker_aggregate_dir
|
||||||
|
.join("runs")
|
||||||
|
.join(request.run_generation.to_string());
|
||||||
|
let mut prepared = WorkerBootstrap::new(
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_context,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
|
WorkerBootstrapLayout::RuntimeManagedRun {
|
||||||
|
run_dir: run_dir.clone(),
|
||||||
|
},
|
||||||
|
self.controller_transport,
|
||||||
)
|
)
|
||||||
|
.prepare()
|
||||||
.await
|
.await
|
||||||
.map_err(|err| format!("failed to create Worker from profile: {err}"))?;
|
.map_err(|error| match error {
|
||||||
|
WorkerBootstrapError::Worker(source) => {
|
||||||
|
format!("failed to create Worker from profile: {source}")
|
||||||
|
}
|
||||||
|
WorkerBootstrapError::Controller { source, .. } => {
|
||||||
|
format!("failed to prepare Worker controller: {source}")
|
||||||
|
}
|
||||||
|
})?;
|
||||||
|
let worker = prepared.worker_mut();
|
||||||
validate_worker_memory_settings(worker.manifest(), &request.request)?;
|
validate_worker_memory_settings(worker.manifest(), &request.request)?;
|
||||||
if let Some(binding) = request.working_directory.as_ref() {
|
if let Some(binding) = request.working_directory.as_ref() {
|
||||||
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
||||||
@@ -935,21 +954,16 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let workspace_client = worker.workspace_client_handle();
|
let workspace_client = worker.workspace_client_handle();
|
||||||
let run_dir = worker_aggregate_dir
|
let started = prepared.start().await.map_err(|error| match error {
|
||||||
.join("runs")
|
WorkerBootstrapError::Worker(source) => {
|
||||||
.join(request.run_generation.to_string());
|
format!("failed to prepare Worker before controller start: {source}")
|
||||||
let (handle, shutdown_rx) = WorkerController::spawn_runtime_managed_run_with_transport(
|
}
|
||||||
worker,
|
WorkerBootstrapError::Controller { source, .. } => format!(
|
||||||
&run_dir,
|
"failed to spawn Worker controller in {}: {source}",
|
||||||
self.controller_transport,
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
.map_err(|err| {
|
|
||||||
format!(
|
|
||||||
"failed to spawn Worker controller in {}: {err}",
|
|
||||||
run_dir.display()
|
run_dir.display()
|
||||||
)
|
),
|
||||||
})?;
|
})?;
|
||||||
|
let (handle, shutdown_rx) = (started.handle, started.shutdown);
|
||||||
if flow_transition_enabled {
|
if flow_transition_enabled {
|
||||||
handle.shared_state.enable_flow_transition();
|
handle.shared_state.enable_flow_transition();
|
||||||
}
|
}
|
||||||
@@ -1118,18 +1132,25 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
let run_dir = worker_aggregate_dir
|
let run_dir = worker_aggregate_dir
|
||||||
.join("runs")
|
.join("runs")
|
||||||
.join(request.run_generation.to_string());
|
.join(request.run_generation.to_string());
|
||||||
let (handle, shutdown_rx) = WorkerController::spawn_runtime_managed_run_with_transport(
|
let started = PreparedWorker::new(
|
||||||
worker,
|
worker,
|
||||||
&run_dir,
|
WorkerBootstrapLayout::RuntimeManagedRun {
|
||||||
|
run_dir: run_dir.clone(),
|
||||||
|
},
|
||||||
self.controller_transport,
|
self.controller_transport,
|
||||||
)
|
)
|
||||||
|
.start()
|
||||||
.await
|
.await
|
||||||
.map_err(|err| {
|
.map_err(|error| match error {
|
||||||
format!(
|
WorkerBootstrapError::Worker(source) => {
|
||||||
"failed to spawn restored Worker controller in {}: {err}",
|
format!("failed to prepare restored Worker: {source}")
|
||||||
|
}
|
||||||
|
WorkerBootstrapError::Controller { source, .. } => format!(
|
||||||
|
"failed to spawn restored Worker controller in {}: {source}",
|
||||||
run_dir.display()
|
run_dir.display()
|
||||||
)
|
),
|
||||||
})?;
|
})?;
|
||||||
|
let (handle, shutdown_rx) = (started.handle, started.shutdown);
|
||||||
if flow_transition_enabled {
|
if flow_transition_enabled {
|
||||||
handle.shared_state.enable_flow_transition();
|
handle.shared_state.enable_flow_transition();
|
||||||
}
|
}
|
||||||
@@ -3091,9 +3112,68 @@ mod tests {
|
|||||||
assert!(!socket_path.exists());
|
assert!(!socket_path.exists());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn profile_runtime_factory_uses_shared_worker_bootstrap_seams() {
|
||||||
|
let source = include_str!("worker_backend.rs");
|
||||||
|
let production = source
|
||||||
|
.split_once("#[cfg(test)]\nmod tests")
|
||||||
|
.map(|(production, _)| production)
|
||||||
|
.expect("worker backend test module marker");
|
||||||
|
let factory = production
|
||||||
|
.split_once("impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory")
|
||||||
|
.map(|(_, factory)| factory)
|
||||||
|
.expect("profile runtime factory implementation");
|
||||||
|
let (fresh, restore) = factory
|
||||||
|
.split_once("async fn restore_controller")
|
||||||
|
.expect("fresh and restore factory paths");
|
||||||
|
let assert_in_order = |path: &str, markers: &[&str]| {
|
||||||
|
let mut offset = 0;
|
||||||
|
for marker in markers {
|
||||||
|
let relative = path[offset..]
|
||||||
|
.find(marker)
|
||||||
|
.unwrap_or_else(|| panic!("missing ordered factory marker {marker}"));
|
||||||
|
offset += relative + marker.len();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
assert_in_order(
|
||||||
|
fresh,
|
||||||
|
&[
|
||||||
|
"WorkerBootstrap::new(",
|
||||||
|
".prepare()",
|
||||||
|
"worker.bind_workdir_session(",
|
||||||
|
"worker.bind_worker_observation_provider(",
|
||||||
|
"install_runtime_flow_transition_feature()",
|
||||||
|
"prepared.start()",
|
||||||
|
],
|
||||||
|
);
|
||||||
|
assert_in_order(
|
||||||
|
restore,
|
||||||
|
&[
|
||||||
|
"Worker::restore_from_worker_metadata_with_context(",
|
||||||
|
"worker.bind_workdir_session(",
|
||||||
|
"worker.bind_worker_observation_provider(",
|
||||||
|
"install_runtime_flow_transition_feature()",
|
||||||
|
"PreparedWorker::new(",
|
||||||
|
".start()",
|
||||||
|
],
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
production.contains("WorkerBootstrap::new("),
|
||||||
|
"fresh runtime Workers must use the shared construction bootstrap"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
production.contains("PreparedWorker::new("),
|
||||||
|
"restored runtime Workers must use the shared pre-exposure lifecycle"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!production.contains("WorkerController::spawn_runtime_managed_run_with_transport"),
|
||||||
|
"runtime factory paths must not bypass the shared controller lifecycle"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
#[serial_test::serial(worker_allocation)]
|
#[serial_test::serial(worker_allocation)]
|
||||||
fn in_process_runtime_reopens_persisted_worker_without_overlong_unix_socket() {
|
fn shared_bootstrap_preserves_in_process_transport_for_fresh_and_restored_runtime_workers() {
|
||||||
let root = tempfile::tempdir().unwrap();
|
let root = tempfile::tempdir().unwrap();
|
||||||
let long_component = "embedded-workspace-store-segment".repeat(4);
|
let long_component = "embedded-workspace-store-segment".repeat(4);
|
||||||
let runtime_store_dir = root.path().join(long_component);
|
let runtime_store_dir = root.path().join(long_component);
|
||||||
|
|||||||
@@ -35,6 +35,14 @@ pub struct WorkerBootstrap<St> {
|
|||||||
workdir_session: Option<WorkdirSessionHandle>,
|
workdir_session: Option<WorkdirSessionHandle>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A constructed Worker whose host-owned live bindings can still be installed
|
||||||
|
/// before Feature installation and controller exposure.
|
||||||
|
pub struct PreparedWorker<C: LlmClient, St: Store> {
|
||||||
|
worker: Worker<C, St>,
|
||||||
|
layout: WorkerBootstrapLayout,
|
||||||
|
transport: WorkerControllerTransport,
|
||||||
|
}
|
||||||
|
|
||||||
/// Live controller returned only after Worker construction and feature
|
/// Live controller returned only after Worker construction and feature
|
||||||
/// installation have completed successfully.
|
/// installation have completed successfully.
|
||||||
pub struct BootstrappedWorker {
|
pub struct BootstrappedWorker {
|
||||||
@@ -99,7 +107,12 @@ where
|
|||||||
self
|
self
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn start(self) -> Result<BootstrappedWorker, WorkerBootstrapError> {
|
/// Construct the Worker without exposing a controller handle. Runtime hosts
|
||||||
|
/// use this seam to bind Workdir, observation, Flow, and other live services
|
||||||
|
/// before [`PreparedWorker::start`] performs Feature installation.
|
||||||
|
pub async fn prepare(
|
||||||
|
self,
|
||||||
|
) -> Result<PreparedWorker<Box<dyn LlmClient>, St>, WorkerBootstrapError> {
|
||||||
let mut worker = Worker::from_manifest_with_context_and_model_client(
|
let mut worker = Worker::from_manifest_with_context_and_model_client(
|
||||||
self.manifest,
|
self.manifest,
|
||||||
self.store,
|
self.store,
|
||||||
@@ -114,7 +127,43 @@ where
|
|||||||
if let Some(workdir_session) = self.workdir_session {
|
if let Some(workdir_session) = self.workdir_session {
|
||||||
worker.bind_workdir_session(Some(workdir_session));
|
worker.bind_workdir_session(Some(workdir_session));
|
||||||
}
|
}
|
||||||
start_worker_controller(worker, self.layout, self.transport).await
|
Ok(PreparedWorker::new(worker, self.layout, self.transport))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start(self) -> Result<BootstrappedWorker, WorkerBootstrapError> {
|
||||||
|
self.prepare().await?.start().await
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<C, St> PreparedWorker<C, St>
|
||||||
|
where
|
||||||
|
C: LlmClient + Clone + 'static,
|
||||||
|
St: Store + WorkerMetadataStore + Clone + Send + Sync + 'static,
|
||||||
|
{
|
||||||
|
/// Wrap a restored Worker in the same pre-exposure lifecycle used by fresh
|
||||||
|
/// bootstraps.
|
||||||
|
pub fn new(
|
||||||
|
worker: Worker<C, St>,
|
||||||
|
layout: WorkerBootstrapLayout,
|
||||||
|
transport: WorkerControllerTransport,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
worker,
|
||||||
|
layout,
|
||||||
|
transport,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn worker(&self) -> &Worker<C, St> {
|
||||||
|
&self.worker
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn worker_mut(&mut self) -> &mut Worker<C, St> {
|
||||||
|
&mut self.worker
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start(self) -> Result<BootstrappedWorker, WorkerBootstrapError> {
|
||||||
|
start_worker_controller(self.worker, self.layout, self.transport).await
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -25,8 +25,8 @@ mod permission;
|
|||||||
mod worker;
|
mod worker;
|
||||||
|
|
||||||
pub use bootstrap::{
|
pub use bootstrap::{
|
||||||
BootstrappedWorker, WorkerBootstrap, WorkerBootstrapError, WorkerBootstrapLayout,
|
BootstrappedWorker, PreparedWorker, WorkerBootstrap, WorkerBootstrapError,
|
||||||
start_worker_controller,
|
WorkerBootstrapLayout, start_worker_controller,
|
||||||
};
|
};
|
||||||
pub use compact::token_counter::{EstimateSource, SplitPoint, TokenEstimate};
|
pub use compact::token_counter::{EstimateSource, SplitPoint, TokenEstimate};
|
||||||
pub use controller::{ShutdownReceiver, WorkerController, WorkerControllerTransport, WorkerHandle};
|
pub use controller::{ShutdownReceiver, WorkerController, WorkerControllerTransport, WorkerHandle};
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ worker
|
|||||||
```
|
```
|
||||||
|
|
||||||
`standalone` は `tui`、`worker-runtime`、`yoi-workspace-server` に依存しない。
|
`standalone` は `tui`、`worker-runtime`、`yoi-workspace-server` に依存しない。
|
||||||
`worker` の direct entrypoint と `standalone` は同じ `start_worker_controller` lifecycle を使い、後者は `WorkerBootstrap` で fresh Worker construction も共有する。
|
`worker` の direct entrypoint、`worker-runtime` の fresh/restore factory、`standalone` は同じ controller lifecycle を使う。fresh runtime/standalone construction は `WorkerBootstrap` を共有し、Runtime は `prepare()` 後かつ Feature install/controller exposure 前に Workdir、observation、Flow の live binding を追加する。restore は replay 済み Worker を `PreparedWorker` へ渡して同じ pre-exposure lifecycle を通す。
|
||||||
|
|
||||||
## authority と lifecycle
|
## authority と lifecycle
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user