fix: make zero-input workers durably restorable
This commit is contained in:
@@ -974,6 +974,10 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
worker
|
||||||
|
.materialize_durable_session_head()
|
||||||
|
.await
|
||||||
|
.map_err(|error| format!("materialize durable Worker session head: {error}"))?;
|
||||||
let workspace_client = worker.workspace_client_handle();
|
let workspace_client = worker.workspace_client_handle();
|
||||||
let started = prepared.start().await.map_err(|error| match error {
|
let started = prepared.start().await.map_err(|error| match error {
|
||||||
WorkerBootstrapError::Worker(source) => {
|
WorkerBootstrapError::Worker(source) => {
|
||||||
|
|||||||
@@ -3617,6 +3617,21 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
|||||||
.is_some_and(|state| state.pre_run_eligible(self.total_tokens().tokens))
|
.is_some_and(|state| state.pre_run_eligible(self.total_tokens().tokens))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Materialize the initial durable session head without running the model.
|
||||||
|
///
|
||||||
|
/// Runtime-managed Workers call this before creation is accepted so even a
|
||||||
|
/// zero-input Worker has an authoritative restore snapshot independent of
|
||||||
|
/// operation-owned launch configuration.
|
||||||
|
pub async fn materialize_durable_session_head(&mut self) -> Result<(), WorkerError>
|
||||||
|
where
|
||||||
|
St: Clone + 'static,
|
||||||
|
{
|
||||||
|
self.refresh_prompt_projection_for_future_operations()?;
|
||||||
|
self.ensure_system_prompt_materialized().await?;
|
||||||
|
self.ensure_segment_head().await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// Prelude shared by `run` / `run_for_notification` / `resume`.
|
/// Prelude shared by `run` / `run_for_notification` / `resume`.
|
||||||
/// Wires up worker hooks, ensures the session is materialized on the
|
/// Wires up worker hooks, ensures the session is materialized on the
|
||||||
/// store, and runs pre-run compact (joining any in-flight memory task
|
/// store, and runs pre-run compact (joining any in-flight memory task
|
||||||
@@ -3625,10 +3640,8 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
|||||||
where
|
where
|
||||||
St: Clone + 'static,
|
St: Clone + 'static,
|
||||||
{
|
{
|
||||||
self.refresh_prompt_projection_for_future_operations()?;
|
self.materialize_durable_session_head().await?;
|
||||||
self.ensure_interceptor_installed();
|
self.ensure_interceptor_installed();
|
||||||
self.ensure_system_prompt_materialized().await?;
|
|
||||||
self.ensure_segment_head().await?;
|
|
||||||
if self.should_pre_run_compact() {
|
if self.should_pre_run_compact() {
|
||||||
self.try_pre_run_compact().await;
|
self.try_pre_run_compact().await;
|
||||||
}
|
}
|
||||||
@@ -9445,6 +9458,29 @@ mod build_summary_prompt_tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn materialize_durable_session_head_commits_without_model_run() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let cwd = dir.path().join("workspace");
|
||||||
|
std::fs::create_dir_all(&cwd).unwrap();
|
||||||
|
let store = session_store::FsStore::new(dir.path().join("sessions")).unwrap();
|
||||||
|
let mut worker = Worker::new(
|
||||||
|
minimal_manifest(),
|
||||||
|
Engine::<_, Mutable, SessionHistoryMetadata>::new_annotated(NoopClient),
|
||||||
|
store,
|
||||||
|
WorkerWorkspaceContext::no_workspace(),
|
||||||
|
WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()),
|
||||||
|
Scope::writable(&cwd).unwrap(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
worker.materialize_durable_session_head().await.unwrap();
|
||||||
|
|
||||||
|
assert_eq!(worker.segment_state.entries_written(), 1);
|
||||||
|
assert!(worker.history().is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn worker_without_initial_system_prompt_does_not_load_feature_contribution() {
|
async fn worker_without_initial_system_prompt_does_not_load_feature_contribution() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
|||||||
@@ -10398,6 +10398,17 @@ async fn scoped_capture_worker_observation_session(
|
|||||||
})))
|
})))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn diagnostics_indicate_unrecoverable_pending_workspace_restore(
|
||||||
|
diagnostics: &[RuntimeDiagnostic],
|
||||||
|
) -> bool {
|
||||||
|
diagnostics.iter().any(|diagnostic| {
|
||||||
|
diagnostic.code == "embedded_worker_execution_rejected"
|
||||||
|
&& diagnostic.message.to_ascii_lowercase().contains(
|
||||||
|
"pending workspace worker restore requires operation-owned launch material",
|
||||||
|
)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
async fn scoped_start_workspace_orchestrator(
|
async fn scoped_start_workspace_orchestrator(
|
||||||
State(api): State<WorkspaceApi>,
|
State(api): State<WorkspaceApi>,
|
||||||
AxumPath(path): AxumPath<ScopedWorkspacePath>,
|
AxumPath(path): AxumPath<ScopedWorkspacePath>,
|
||||||
@@ -10414,6 +10425,7 @@ async fn scoped_start_workspace_orchestrator(
|
|||||||
)
|
)
|
||||||
})?;
|
})?;
|
||||||
|
|
||||||
|
let mut disposition = "created";
|
||||||
if let Some(existing) = find_workspace_orchestrator(&api) {
|
if let Some(existing) = find_workspace_orchestrator(&api) {
|
||||||
if workspace_orchestrator_is_online(&existing) {
|
if workspace_orchestrator_is_online(&existing) {
|
||||||
return Ok(Json(workspace_orchestrator_response(&api, "existing")));
|
return Ok(Json(workspace_orchestrator_response(&api, "existing")));
|
||||||
@@ -10422,7 +10434,14 @@ async fn scoped_start_workspace_orchestrator(
|
|||||||
.runtime
|
.runtime
|
||||||
.restore_worker(&existing.worker)
|
.restore_worker(&existing.worker)
|
||||||
.map_err(|error| error.into_error())?;
|
.map_err(|error| error.into_error())?;
|
||||||
if restored.state != WorkerOperationState::Accepted {
|
if restored.state == WorkerOperationState::Accepted {
|
||||||
|
*api.orchestrator_attention_fingerprint
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
|
||||||
|
dispatch_orchestrator_queue_attention(&api);
|
||||||
|
return Ok(Json(workspace_orchestrator_response(&api, "restored")));
|
||||||
|
}
|
||||||
|
if !diagnostics_indicate_unrecoverable_pending_workspace_restore(&restored.diagnostics) {
|
||||||
return Err(ApiError::with_diagnostics(
|
return Err(ApiError::with_diagnostics(
|
||||||
Error::RuntimeOperationFailed {
|
Error::RuntimeOperationFailed {
|
||||||
runtime_id: existing.worker.runtime_id.clone(),
|
runtime_id: existing.worker.runtime_id.clone(),
|
||||||
@@ -10432,11 +10451,21 @@ async fn scoped_start_workspace_orchestrator(
|
|||||||
restored.diagnostics,
|
restored.diagnostics,
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
*api.orchestrator_attention_fingerprint
|
let deleted = api
|
||||||
.lock()
|
.runtime
|
||||||
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
|
.delete_worker(&existing.worker)
|
||||||
dispatch_orchestrator_queue_attention(&api);
|
.map_err(|error| error.into_error())?;
|
||||||
return Ok(Json(workspace_orchestrator_response(&api, "restored")));
|
if deleted.state != WorkerOperationState::Accepted {
|
||||||
|
return Err(ApiError::with_diagnostics(
|
||||||
|
Error::RuntimeOperationFailed {
|
||||||
|
runtime_id: existing.worker.runtime_id.clone(),
|
||||||
|
code: "workspace_orchestrator_replacement_delete_rejected".to_string(),
|
||||||
|
message: "Embedded Runtime rejected replacement of an unrecoverable pending Workspace Orchestrator".to_string(),
|
||||||
|
},
|
||||||
|
deleted.diagnostics,
|
||||||
|
));
|
||||||
|
}
|
||||||
|
disposition = "recreated";
|
||||||
}
|
}
|
||||||
|
|
||||||
let result = api.spawn_workspace_worker(
|
let result = api.spawn_workspace_worker(
|
||||||
@@ -10485,7 +10514,7 @@ async fn scoped_start_workspace_orchestrator(
|
|||||||
.lock()
|
.lock()
|
||||||
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
|
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
|
||||||
dispatch_orchestrator_queue_attention(&api);
|
dispatch_orchestrator_queue_attention(&api);
|
||||||
Ok(Json(workspace_orchestrator_response(&api, "created")))
|
Ok(Json(workspace_orchestrator_response(&api, disposition)))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn worker_launch_worker_summary(worker: WorkerSummary) -> WorkerLaunchWorkerSummary {
|
fn worker_launch_worker_summary(worker: WorkerSummary) -> WorkerLaunchWorkerSummary {
|
||||||
@@ -20722,7 +20751,7 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn production_profile_backend_rejects_unrecoverable_pending_orchestrator_restore() {
|
async fn production_profile_backend_restores_zero_input_orchestrator() {
|
||||||
let workspace = tempfile::tempdir().unwrap();
|
let workspace = tempfile::tempdir().unwrap();
|
||||||
init_clean_git_workspace(workspace.path());
|
init_clean_git_workspace(workspace.path());
|
||||||
let config = test_server_config(workspace.path());
|
let config = test_server_config(workspace.path());
|
||||||
@@ -20765,20 +20794,19 @@ mod tests {
|
|||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert_eq!(stopped.state, WorkerOperationState::Accepted);
|
assert_eq!(stopped.state, WorkerOperationState::Accepted);
|
||||||
let error = scoped_start_workspace_orchestrator(
|
let Json(restored) = scoped_start_workspace_orchestrator(
|
||||||
State(api.clone()),
|
State(api.clone()),
|
||||||
AxumPath(ScopedWorkspacePath {
|
AxumPath(ScopedWorkspacePath {
|
||||||
workspace_id: workspace_id.clone(),
|
workspace_id: workspace_id.clone(),
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect_err("pending Workspace Orchestrator restore without durable Prompt must fail");
|
.expect("zero-input Workspace Orchestrator should restore from its durable session head");
|
||||||
assert!(
|
assert_eq!(restored.disposition, "restored");
|
||||||
format!("{error:?}").contains(
|
assert!(restored.online);
|
||||||
"pending Workspace Worker restore requires operation-owned launch material"
|
let restored_worker = restored.worker.expect("restored Orchestrator Worker");
|
||||||
),
|
assert_eq!(restored_worker.runtime_id, worker.runtime_id);
|
||||||
"unexpected restore error: {error:?}"
|
assert_eq!(restored_worker.worker_id, worker.worker_id);
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -21699,6 +21727,17 @@ mod tests {
|
|||||||
assert_eq!(sanitized, "failed to open server database");
|
assert_eq!(sanitized, "failed to open server database");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn unrecoverable_pending_workspace_restore_is_typed_for_replacement() {
|
||||||
|
let diagnostics = [RuntimeDiagnostic {
|
||||||
|
code: "embedded_worker_execution_rejected".to_string(),
|
||||||
|
severity: DiagnosticSeverity::Error,
|
||||||
|
message: "Restore Errored: pending Workspace Worker restore requires operation-owned launch material; generic restore must not reconstruct it from current Workspace config".to_string(),
|
||||||
|
}];
|
||||||
|
|
||||||
|
assert!(diagnostics_indicate_unrecoverable_pending_workspace_restore(&diagnostics));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn workdir_runtime_miss_uses_exact_typed_code() {
|
fn workdir_runtime_miss_uses_exact_typed_code() {
|
||||||
let typed_not_found = [RuntimeDiagnostic {
|
let typed_not_found = [RuntimeDiagnostic {
|
||||||
|
|||||||
Reference in New Issue
Block a user