runtime: avoid profile archive during restore
This commit is contained in:
parent
7caf6cd6e6
commit
4de801f8d0
|
|
@ -344,13 +344,7 @@ impl ProfileRuntimeWorkerFactory {
|
||||||
&self,
|
&self,
|
||||||
request: WorkerExecutionRestoreDryRequest,
|
request: WorkerExecutionRestoreDryRequest,
|
||||||
) -> Result<(), String> {
|
) -> Result<(), String> {
|
||||||
let WorkerExecutionRestoreDryRequest {
|
let WorkerExecutionRestoreDryRequest { worker_ref, .. } = request;
|
||||||
worker_ref,
|
|
||||||
request,
|
|
||||||
previous_execution: _,
|
|
||||||
working_directory,
|
|
||||||
config_bundle: _,
|
|
||||||
} = request;
|
|
||||||
let store_dir = self.store_dir()?;
|
let store_dir = self.store_dir()?;
|
||||||
let session_store = FsStore::new(&store_dir).map_err(|err| {
|
let session_store = FsStore::new(&store_dir).map_err(|err| {
|
||||||
format!(
|
format!(
|
||||||
|
|
@ -387,31 +381,6 @@ impl ProfileRuntimeWorkerFactory {
|
||||||
if state.entries_count == 0 {
|
if state.entries_count == 0 {
|
||||||
return Err(format!("active segment {active_segment_id} is empty"));
|
return Err(format!("active segment {active_segment_id} is empty"));
|
||||||
}
|
}
|
||||||
let worker_root = working_directory
|
|
||||||
.as_ref()
|
|
||||||
.map(|binding| binding.root.clone())
|
|
||||||
.unwrap_or_else(|| self.profile_base_dir.clone());
|
|
||||||
let profile = self.runtime_profile_for_request(&request);
|
|
||||||
let selector = profile.as_deref().unwrap_or("builtin:default");
|
|
||||||
let archive = self
|
|
||||||
.resolve_profile_source_archive(&request.profile_source)
|
|
||||||
.await?;
|
|
||||||
let manifest = archive
|
|
||||||
.resolve_profile(selector, &worker_root, &worker_name)
|
|
||||||
.map_err(|err| format!("failed to resolve profile source archive: {err}"))?;
|
|
||||||
if working_directory.is_some() {
|
|
||||||
worker::entrypoint::resolve_runtime_profile_manifest_from_manifest(
|
|
||||||
manifest,
|
|
||||||
&worker_root,
|
|
||||||
&worker_name,
|
|
||||||
)?;
|
|
||||||
} else {
|
|
||||||
worker::entrypoint::resolve_runtime_profile_manifest_from_manifest_without_filesystem(
|
|
||||||
manifest,
|
|
||||||
&worker_root,
|
|
||||||
&worker_name,
|
|
||||||
)?;
|
|
||||||
}
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -514,11 +483,6 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
request: WorkerExecutionRestoreRequest,
|
request: WorkerExecutionRestoreRequest,
|
||||||
) -> Result<WorkerHandle, String> {
|
) -> Result<WorkerHandle, String> {
|
||||||
let worker_name = Self::runtime_worker_name_for_ref(&request.worker_ref);
|
let worker_name = Self::runtime_worker_name_for_ref(&request.worker_ref);
|
||||||
let worker_root = request
|
|
||||||
.working_directory
|
|
||||||
.as_ref()
|
|
||||||
.map(|binding| binding.root().to_path_buf())
|
|
||||||
.unwrap_or_else(|| self.profile_base_dir.clone());
|
|
||||||
let filesystem_authority = request
|
let filesystem_authority = request
|
||||||
.working_directory
|
.working_directory
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -564,24 +528,6 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
Err(WorkerError::WorkerMetadataPending { .. })
|
Err(WorkerError::WorkerMetadataPending { .. })
|
||||||
if request.request.initial_input.is_none() =>
|
if request.request.initial_input.is_none() =>
|
||||||
{
|
{
|
||||||
let profile = self.runtime_profile_for_request(&request.request);
|
|
||||||
let selector = profile.as_deref().unwrap_or("builtin:default");
|
|
||||||
let archive = self
|
|
||||||
.resolve_profile_source_archive(&request.request.profile_source)
|
|
||||||
.await?;
|
|
||||||
let (mut pending_manifest, pending_loader) = {
|
|
||||||
let manifest = archive
|
|
||||||
.resolve_profile(selector, &worker_root, &worker_name)
|
|
||||||
.map_err(|err| {
|
|
||||||
format!("failed to resolve profile source archive: {err}")
|
|
||||||
})?;
|
|
||||||
worker::entrypoint::resolve_runtime_profile_manifest_from_manifest(
|
|
||||||
manifest,
|
|
||||||
&worker_root,
|
|
||||||
&worker_name,
|
|
||||||
)?
|
|
||||||
};
|
|
||||||
pending_manifest.worker.name = worker_name.clone();
|
|
||||||
let session_store = FsStore::new(&store_dir).map_err(|err| {
|
let session_store = FsStore::new(&store_dir).map_err(|err| {
|
||||||
format!(
|
format!(
|
||||||
"failed to initialize session store at {}: {err}",
|
"failed to initialize session store at {}: {err}",
|
||||||
|
|
@ -596,15 +542,16 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
)
|
)
|
||||||
})?;
|
})?;
|
||||||
let store = CombinedStore::new(session_store, worker_metadata_store);
|
let store = CombinedStore::new(session_store, worker_metadata_store);
|
||||||
Worker::from_manifest_with_context(
|
Worker::restore_pending_from_worker_metadata_with_context(
|
||||||
pending_manifest,
|
&worker_name,
|
||||||
|
manifest.clone(),
|
||||||
store,
|
store,
|
||||||
pending_loader,
|
loader,
|
||||||
workspace_context,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| format!("failed to recreate pending Worker from profile: {err}"))?
|
.map_err(|err| format!("failed to recreate pending Worker from metadata: {err}"))?
|
||||||
}
|
}
|
||||||
Err(err) => return Err(format!("failed to restore Worker from metadata: {err}")),
|
Err(err) => return Err(format!("failed to restore Worker from metadata: {err}")),
|
||||||
};
|
};
|
||||||
|
|
@ -1468,6 +1415,7 @@ mod tests {
|
||||||
[model]
|
[model]
|
||||||
scheme = "anthropic"
|
scheme = "anthropic"
|
||||||
model_id = "test-model"
|
model_id = "test-model"
|
||||||
|
auth = { kind = "none" }
|
||||||
|
|
||||||
[engine]
|
[engine]
|
||||||
max_tokens = 100
|
max_tokens = 100
|
||||||
|
|
@ -1740,6 +1688,122 @@ mod tests {
|
||||||
.expect("embedded archive should resolve without Backend resource client");
|
.expect("embedded archive should resolve without Backend resource client");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn dry_restore_validates_saved_worker_state_without_profile_source_resolution() {
|
||||||
|
let root = tempfile::tempdir().unwrap();
|
||||||
|
let store_dir = root.path().join("sessions");
|
||||||
|
let worker_metadata_dir = root.path().join("workers");
|
||||||
|
let worker_ref = WorkerRef::new(crate::identity::WorkerId::new(1));
|
||||||
|
let worker_name = ProfileRuntimeWorkerFactory::runtime_worker_name_for_ref(&worker_ref);
|
||||||
|
let session_id = session_store::new_session_id();
|
||||||
|
let segment_id = session_store::new_segment_id();
|
||||||
|
FsStore::new(&store_dir)
|
||||||
|
.unwrap()
|
||||||
|
.create_segment(
|
||||||
|
session_id,
|
||||||
|
segment_id,
|
||||||
|
&[session_store::LogEntry::Invoke {
|
||||||
|
ts: 1,
|
||||||
|
trigger: protocol::InvokeKind::UserSend,
|
||||||
|
}],
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
FsWorkerStore::new(&worker_metadata_dir)
|
||||||
|
.unwrap()
|
||||||
|
.set_active(
|
||||||
|
&worker_name,
|
||||||
|
Some(session_store::WorkerActiveSegmentRef::active_segment(
|
||||||
|
session_id, segment_id,
|
||||||
|
)),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let request = create_request("restore");
|
||||||
|
let result = ProfileRuntimeWorkerFactory::new(root.path())
|
||||||
|
.with_store_dir(&store_dir)
|
||||||
|
.with_worker_metadata_dir(&worker_metadata_dir)
|
||||||
|
.dry_restore_controller(WorkerExecutionRestoreDryRequest {
|
||||||
|
worker_ref,
|
||||||
|
request,
|
||||||
|
previous_execution: crate::execution::WorkerExecutionStatus::alive(
|
||||||
|
WorkerExecutionRunState::Idle,
|
||||||
|
),
|
||||||
|
working_directory: None,
|
||||||
|
config_bundle: None,
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
result.status,
|
||||||
|
crate::execution::WorkerRestoreDryCheckStatus::Valid,
|
||||||
|
"{}",
|
||||||
|
result.message
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[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 worker_ref = WorkerRef::new(crate::identity::WorkerId::new(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!(
|
||||||
|
r#"
|
||||||
|
[worker]
|
||||||
|
name = "{}"
|
||||||
|
pwd = "{}"
|
||||||
|
|
||||||
|
[model]
|
||||||
|
scheme = "anthropic"
|
||||||
|
model_id = "test-model"
|
||||||
|
auth = {{ kind = "none" }}
|
||||||
|
|
||||||
|
[engine]
|
||||||
|
max_tokens = 100
|
||||||
|
|
||||||
|
[[scope.allow]]
|
||||||
|
target = "{}"
|
||||||
|
permission = "write"
|
||||||
|
"#,
|
||||||
|
worker_name,
|
||||||
|
root.path().display(),
|
||||||
|
root.path().display(),
|
||||||
|
))
|
||||||
|
.unwrap();
|
||||||
|
FsWorkerStore::new(&worker_metadata_dir)
|
||||||
|
.unwrap()
|
||||||
|
.set_active(
|
||||||
|
&worker_name,
|
||||||
|
Some(session_store::WorkerActiveSegmentRef::pending_segment(
|
||||||
|
session_id,
|
||||||
|
)),
|
||||||
|
Some(serde_json::to_value(&manifest).unwrap()),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let request = create_request("restore");
|
||||||
|
let handle = ProfileRuntimeWorkerFactory::new(root.path())
|
||||||
|
.with_store_dir(&store_dir)
|
||||||
|
.with_worker_metadata_dir(&worker_metadata_dir)
|
||||||
|
.restore_controller(WorkerExecutionRestoreRequest {
|
||||||
|
worker_ref: worker_ref.clone(),
|
||||||
|
request,
|
||||||
|
context: test_execution_context(worker_ref),
|
||||||
|
previous_execution: crate::execution::WorkerExecutionStatus::alive(
|
||||||
|
WorkerExecutionRunState::Idle,
|
||||||
|
),
|
||||||
|
working_directory: None,
|
||||||
|
config_bundle: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("pending restore should use the saved manifest snapshot");
|
||||||
|
|
||||||
|
handle.send(Method::Shutdown).await.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn builtin_profile_selector_is_not_double_prefixed() {
|
fn builtin_profile_selector_is_not_double_prefixed() {
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|
|
||||||
|
|
@ -3769,6 +3769,65 @@ where
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Recreate a pending Worker whose metadata has a session id but no
|
||||||
|
/// materialized segment yet.
|
||||||
|
///
|
||||||
|
/// Pending Workers have already had their profile source resolved at creation
|
||||||
|
/// time, but they have not rendered the system prompt or written
|
||||||
|
/// `SegmentStart`. Restore therefore uses only the resolved manifest snapshot
|
||||||
|
/// stored in Worker metadata and never re-resolves the profile source.
|
||||||
|
pub async fn restore_pending_from_worker_metadata_with_context(
|
||||||
|
worker_name: &str,
|
||||||
|
fallback: WorkerManifest,
|
||||||
|
store: St,
|
||||||
|
loader: PromptLoader,
|
||||||
|
workspace_context: WorkerWorkspaceContext,
|
||||||
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
|
) -> Result<Self, WorkerError> {
|
||||||
|
let metadata =
|
||||||
|
store
|
||||||
|
.read_by_name(worker_name)?
|
||||||
|
.ok_or_else(|| WorkerError::WorkerMetadataMissing {
|
||||||
|
worker_name: worker_name.to_string(),
|
||||||
|
})?;
|
||||||
|
let active = metadata
|
||||||
|
.active
|
||||||
|
.ok_or_else(|| WorkerError::WorkerMetadataInactive {
|
||||||
|
worker_name: worker_name.to_string(),
|
||||||
|
})?;
|
||||||
|
if let Some(segment_id) = active.segment_id {
|
||||||
|
return Self::restore_from_manifest_with_context(
|
||||||
|
active.session_id,
|
||||||
|
segment_id,
|
||||||
|
restore_manifest_from_worker_metadata_snapshot(
|
||||||
|
worker_name,
|
||||||
|
metadata.resolved_manifest_snapshot,
|
||||||
|
fallback,
|
||||||
|
)?,
|
||||||
|
store,
|
||||||
|
loader,
|
||||||
|
workspace_context,
|
||||||
|
filesystem_authority,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
let snapshot = metadata.resolved_manifest_snapshot.ok_or_else(|| {
|
||||||
|
WorkerError::WorkerMetadataManifestSnapshotMissing {
|
||||||
|
worker_name: worker_name.to_string(),
|
||||||
|
}
|
||||||
|
})?;
|
||||||
|
let manifest =
|
||||||
|
restore_manifest_from_worker_metadata_snapshot(worker_name, Some(snapshot), fallback)?;
|
||||||
|
Self::from_manifest_with_context(
|
||||||
|
manifest,
|
||||||
|
store,
|
||||||
|
loader,
|
||||||
|
workspace_context,
|
||||||
|
filesystem_authority,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
/// Restore a Worker from an existing session log.
|
/// Restore a Worker from an existing session log.
|
||||||
///
|
///
|
||||||
/// Uses the resolved manifest supplied by the caller, seeds a
|
/// Uses the resolved manifest supplied by the caller, seeds a
|
||||||
|
|
@ -4548,6 +4607,9 @@ pub enum WorkerError {
|
||||||
session_id: SessionId,
|
session_id: SessionId,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
#[error("worker metadata for {worker_name} does not include a resolved manifest snapshot")]
|
||||||
|
WorkerMetadataManifestSnapshotMissing { worker_name: String },
|
||||||
|
|
||||||
#[error(
|
#[error(
|
||||||
"worker metadata for {worker_name} contains an invalid resolved manifest snapshot: {source}"
|
"worker metadata for {worker_name} contains an invalid resolved manifest snapshot: {source}"
|
||||||
)]
|
)]
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user