Merge commit '80221289935e227820aec902988a4f987273e05e' into work/T-549-provider-published-ref
# Conflicts: # crates/worker-runtime/src/runtime.rs
This commit is contained in:
@@ -4,7 +4,7 @@
|
||||
|
||||
use agen::llm_client::scheme::{Scheme, anthropic::AnthropicScheme};
|
||||
use agen::llm_client::transport::{HttpTransport, ResolvedAuth};
|
||||
use agen::{Engine, EngineRunExit, StopReason};
|
||||
use agen::{Engine, EngineRunExit, RunInterruptionReason};
|
||||
use std::time::Duration;
|
||||
|
||||
#[tokio::main]
|
||||
@@ -51,7 +51,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
EngineRunExit::Finished => println!("✅ Task completed normally"),
|
||||
EngineRunExit::Paused => println!("⏸️ Task paused"),
|
||||
EngineRunExit::Yielded => println!("↩️ Task yielded"),
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached) => {
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached) => {
|
||||
println!("🔒 Turn limit reached")
|
||||
}
|
||||
EngineRunExit::Interrupted(reason) => println!("❌ Task interrupted: {reason:?}"),
|
||||
|
||||
@@ -39,7 +39,7 @@ use tracing::info;
|
||||
use tracing_subscriber::EnvFilter;
|
||||
|
||||
use agen::{
|
||||
Engine, EngineRunExit, StopReason,
|
||||
Engine, EngineRunExit, RunInterruptionReason,
|
||||
interceptor::{Interceptor, PostToolAction, ToolResultInfo},
|
||||
llm_client::{
|
||||
LlmClient,
|
||||
@@ -478,7 +478,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
// One-shot mode
|
||||
if let Some(prompt) = args.prompt {
|
||||
let output = engine.run(&mut history, &prompt).await;
|
||||
if let EngineRunExit::Interrupted(StopReason::Unexpected(error)) = output.result {
|
||||
if let EngineRunExit::Interrupted(RunInterruptionReason::Unexpected(error)) = output.result
|
||||
{
|
||||
eprintln!("\n❌ Error: {error}");
|
||||
}
|
||||
|
||||
@@ -518,7 +519,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
break;
|
||||
}
|
||||
|
||||
if let EngineRunExit::Interrupted(StopReason::Unexpected(error)) =
|
||||
if let EngineRunExit::Interrupted(RunInterruptionReason::Unexpected(error)) =
|
||||
locked.run(&mut history, input).await
|
||||
{
|
||||
eprintln!("\n❌ Error: {error}");
|
||||
|
||||
@@ -147,12 +147,12 @@ pub enum EngineRunExit {
|
||||
Finished,
|
||||
Paused,
|
||||
Yielded,
|
||||
Interrupted(StopReason),
|
||||
Interrupted(RunInterruptionReason),
|
||||
}
|
||||
|
||||
/// A typed reason why an engine run could not finish normally.
|
||||
#[derive(Debug)]
|
||||
pub enum StopReason {
|
||||
pub enum RunInterruptionReason {
|
||||
LimitReached,
|
||||
ContextWindowExceeded,
|
||||
Cancelled,
|
||||
@@ -165,13 +165,15 @@ impl From<Result<EngineResult, EngineError>> for EngineRunExit {
|
||||
Ok(EngineResult::Finished) => Self::Finished,
|
||||
Ok(EngineResult::Paused) => Self::Paused,
|
||||
Ok(EngineResult::Yielded) => Self::Yielded,
|
||||
Ok(EngineResult::LimitReached) => Self::Interrupted(StopReason::LimitReached),
|
||||
Err(EngineError::Client(ClientError::ContextWindowExceeded)) => {
|
||||
Self::Interrupted(StopReason::ContextWindowExceeded)
|
||||
Ok(EngineResult::LimitReached) => {
|
||||
Self::Interrupted(RunInterruptionReason::LimitReached)
|
||||
}
|
||||
Err(EngineError::Cancelled) => Self::Interrupted(StopReason::Cancelled),
|
||||
Err(EngineError::Client(ClientError::ContextWindowExceeded)) => {
|
||||
Self::Interrupted(RunInterruptionReason::ContextWindowExceeded)
|
||||
}
|
||||
Err(EngineError::Cancelled) => Self::Interrupted(RunInterruptionReason::Cancelled),
|
||||
Err(EngineError::PauseRequested) => Self::Paused,
|
||||
Err(error) => Self::Interrupted(StopReason::Unexpected(error)),
|
||||
Err(error) => Self::Interrupted(RunInterruptionReason::Unexpected(error)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ pub use agen_macros::{description, tool, tool_registry};
|
||||
pub use callback::{TextBlockScope, ThinkingBlockScope, ToolUseBlockScope};
|
||||
pub use engine::{
|
||||
Engine, EngineConfig, EngineError, EngineResult, EngineRunExit, EngineRunOutput,
|
||||
LlmRetryNotice, StopReason, ToolRegistryError,
|
||||
LlmRetryNotice, RunInterruptionReason, ToolRegistryError,
|
||||
};
|
||||
pub use handler::ToolUseBlockStart;
|
||||
pub use history::{History, HistoryEntry};
|
||||
|
||||
@@ -14,7 +14,7 @@ use agen::interceptor::{
|
||||
};
|
||||
use agen::llm_client::event::{Event, ResponseStatus, StatusEvent};
|
||||
use agen::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput};
|
||||
use agen::{Engine, EngineError, EngineRunExit, History, StopReason};
|
||||
use agen::{Engine, EngineError, EngineRunExit, History, RunInterruptionReason};
|
||||
use async_trait::async_trait;
|
||||
use common::MockLlmClient;
|
||||
|
||||
@@ -205,7 +205,7 @@ async fn history_append_failure_stops_before_tool_execution() {
|
||||
let exit = engine.run(&mut history, "use the tool").await;
|
||||
|
||||
assert!(
|
||||
matches!(exit, EngineRunExit::Interrupted(StopReason::Unexpected(EngineError::HistoryAppend(ref message))) if message == "simulated ENOSPC")
|
||||
matches!(exit, EngineRunExit::Interrupted(RunInterruptionReason::Unexpected(EngineError::HistoryAppend(ref message))) if message == "simulated ENOSPC")
|
||||
);
|
||||
assert_eq!(tool.call_count(), 0);
|
||||
assert_eq!(history.len(), 1);
|
||||
@@ -730,7 +730,7 @@ async fn paused_tool_resume_does_not_reset_the_consumed_turn_budget() {
|
||||
|
||||
assert!(matches!(
|
||||
engine.resume(&mut history).await,
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached)
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached)
|
||||
));
|
||||
assert_eq!(engine.turn_count(), 1);
|
||||
assert_eq!(engine.active_run_turn_count(), None);
|
||||
@@ -785,7 +785,7 @@ async fn interceptor_continuation_consumes_the_logical_run_budget() {
|
||||
|
||||
assert!(matches!(
|
||||
engine.run(&mut history, "start").await,
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached)
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached)
|
||||
));
|
||||
assert_eq!(engine.turn_count(), 1);
|
||||
assert_eq!(engine.llm_call_count(), 1);
|
||||
@@ -803,7 +803,7 @@ async fn restored_active_run_budget_is_enforced_before_another_llm_call() {
|
||||
|
||||
assert!(matches!(
|
||||
engine.resume(&mut history).await,
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached)
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached)
|
||||
));
|
||||
assert_eq!(engine.turn_count(), 7);
|
||||
assert_eq!(engine.llm_call_count(), 0);
|
||||
|
||||
@@ -580,7 +580,7 @@ async fn cooperative_cancellation_commits_bounded_terminal_output() {
|
||||
);
|
||||
assert!(matches!(
|
||||
output.result,
|
||||
agen::EngineRunExit::Interrupted(agen::StopReason::Cancelled)
|
||||
agen::EngineRunExit::Interrupted(agen::RunInterruptionReason::Cancelled)
|
||||
));
|
||||
}
|
||||
|
||||
@@ -1214,7 +1214,7 @@ async fn post_tool_abort_commits_confirmed_result_before_stopping_run() {
|
||||
);
|
||||
assert!(matches!(
|
||||
output.result,
|
||||
agen::EngineRunExit::Interrupted(agen::StopReason::Unexpected(
|
||||
agen::EngineRunExit::Interrupted(agen::RunInterruptionReason::Unexpected(
|
||||
agen::EngineError::Aborted(ref reason)
|
||||
)) if reason == "policy stopped the run"
|
||||
));
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use reqwest::Method;
|
||||
use serde::Serialize;
|
||||
use serde::de::DeserializeOwned;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use ticket::{
|
||||
MarkdownText, NewOrchestrationPlanRecord, NewTicket, NewTicketEvent, NewTicketRelation,
|
||||
OrchestrationPlanKind, OrchestrationPlanRecord, Ticket, TicketBackend, TicketDependencyCheck,
|
||||
@@ -9,39 +9,17 @@ use ticket::{
|
||||
TicketRelationKind, TicketRelationView, TicketStateChange, TicketStateSelector, TicketSummary,
|
||||
};
|
||||
use workspace_api::{
|
||||
ListResponse, ObjectiveCreateRequest, ObjectiveDetail, ObjectiveEditRequest,
|
||||
ObjectiveLinkTicketRequest, ObjectiveStateRequest, ObjectiveSummary,
|
||||
BrowserCreateWorkerResponse, BrowserWorkspaceOrchestratorResponse,
|
||||
CreateWorkspaceWorkerRequest, ListResponse, ObjectiveCreateRequest, ObjectiveDetail,
|
||||
ObjectiveEditRequest, ObjectiveLinkTicketRequest, ObjectiveStateRequest, ObjectiveSummary,
|
||||
TICKET_ORCHESTRATION_PLANS_QUERY_PATH, TICKET_RELATIONS_QUERY_PATH,
|
||||
WorkerLaunchOptionsResponse,
|
||||
};
|
||||
|
||||
use crate::{BackendApiClient, BackendWorkspaceClientError};
|
||||
|
||||
const DEFAULT_PRODUCT_LIST_LIMIT: usize = 1_000;
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct BackendWorkerLaunchOptions {
|
||||
runtimes: Vec<BackendWorkerLaunchRuntime>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct BackendWorkerLaunchRuntime {
|
||||
runtime_id: String,
|
||||
worker_creation_available: bool,
|
||||
working_directory_required: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct BackendCreateWorkerResponse {
|
||||
runtime_id: String,
|
||||
worker_id: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct BackendWorkspaceOrchestratorResponse {
|
||||
disposition: String,
|
||||
worker: Option<BackendCreateWorkerResponse>,
|
||||
}
|
||||
|
||||
/// Workspace-scoped Backend client for Ticket and Objective product state.
|
||||
///
|
||||
/// Construction requires both the selected Backend URL and Workspace identity.
|
||||
@@ -267,7 +245,7 @@ impl BackendWorkspaceProductClient {
|
||||
&self,
|
||||
ticket_id: &str,
|
||||
) -> Result<String, BackendWorkspaceClientError> {
|
||||
let options: BackendWorkerLaunchOptions = self.get_json("/workers/launch-options")?;
|
||||
let options: WorkerLaunchOptionsResponse = self.get_json("/workers/launch-options")?;
|
||||
let runtime = options
|
||||
.runtimes
|
||||
.iter()
|
||||
@@ -278,19 +256,19 @@ impl BackendWorkspaceProductClient {
|
||||
.to_string(),
|
||||
)
|
||||
})?;
|
||||
let response: BackendCreateWorkerResponse = self.send_json(
|
||||
Method::POST,
|
||||
"/workers",
|
||||
Some(&serde_json::json!({
|
||||
"runtime_id": runtime.runtime_id,
|
||||
"display_name": format!("intake-{ticket_id}"),
|
||||
"profile": "builtin:intake",
|
||||
"initial_submit": [{
|
||||
"kind": "text",
|
||||
"content": format!("Please handle intake for Ticket {ticket_id}.")
|
||||
}]
|
||||
})),
|
||||
)?;
|
||||
let request = CreateWorkspaceWorkerRequest {
|
||||
runtime_id: runtime.runtime_id.clone(),
|
||||
display_name: format!("intake-{ticket_id}"),
|
||||
profile: Some("builtin:intake".to_string()),
|
||||
ticket_assignment: None,
|
||||
initial_submit: vec![protocol::Segment::Text {
|
||||
content: format!("Please handle intake for Ticket {ticket_id}."),
|
||||
}],
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
};
|
||||
let response: BrowserCreateWorkerResponse =
|
||||
self.send_json(Method::POST, "/workers", Some(&request))?;
|
||||
Ok(format!(
|
||||
"Started Intake Worker {}/{} for Ticket {ticket_id}",
|
||||
response.runtime_id, response.worker_id
|
||||
@@ -298,7 +276,7 @@ impl BackendWorkspaceProductClient {
|
||||
}
|
||||
|
||||
pub fn start_workspace_orchestrator(&self) -> Result<String, BackendWorkspaceClientError> {
|
||||
let response: BackendWorkspaceOrchestratorResponse =
|
||||
let response: BrowserWorkspaceOrchestratorResponse =
|
||||
self.send_json::<(), _>(Method::POST, "/orchestrator", None)?;
|
||||
let worker = response.worker.ok_or_else(|| {
|
||||
BackendWorkspaceClientError::InvalidTarget(
|
||||
@@ -792,11 +770,11 @@ mod tests {
|
||||
let (base_url, requests, handle) = response_sequence_server(vec![
|
||||
(
|
||||
"200 OK",
|
||||
r#"{"runtimes":[{"runtime_id":"embedded","worker_creation_available":true,"working_directory_required":false}]}"#,
|
||||
r#"{"workspace_id":"workspace-a","runtimes":[{"runtime_id":"embedded","display_name":"Embedded","built_in":true,"worker_creation_available":true,"working_directory_required":false,"status":"connected","diagnostics":[]}],"default_profile":null,"profiles":[],"repositories":[],"working_directories":[],"diagnostics":[]}"#,
|
||||
),
|
||||
(
|
||||
"200 OK",
|
||||
r#"{"runtime_id":"embedded","worker_id":"worker-1"}"#,
|
||||
r#"{"workspace_id":"workspace-a","runtime_id":"embedded","worker_id":"worker-1","console_href":"/w/workspace-a/workers/worker-1","worker":{"runtime_id":"embedded","worker_id":"worker-1","host_id":"embedded","display_name":"Intake","label":"worker-1","profile":"builtin:intake","singleton_key":null,"tags":[],"workspace":{"visibility":"workspace","identity":"workspace-a","workspace_id":"workspace-a"},"state":"idle","last_seen_at":null,"pinned":false,"retention_state":"active","implementation":{"kind":"runtime","display_hint":"Runtime Worker"},"capabilities":{"can_stop":true,"can_spawn_followup":false},"diagnostics":[]},"diagnostics":[]}"#,
|
||||
),
|
||||
]);
|
||||
let client = BackendWorkspaceProductClient::new_with_access_token(
|
||||
@@ -824,7 +802,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn workspace_orchestrator_launch_uses_scoped_backend_route() {
|
||||
let body = r#"{"disposition":"created","worker":{"runtime_id":"embedded","worker_id":"worker-2"}}"#;
|
||||
let body = r#"{"workspace_id":"workspace-a","online":true,"disposition":"created","worker":{"runtime_id":"embedded","worker_id":"worker-2","host_id":"embedded","display_name":"Orchestrator","label":"worker-2","profile":"builtin:orchestrator","singleton_key":"workspace-orchestrator","tags":[],"workspace":{"visibility":"workspace","identity":"workspace-a","workspace_id":"workspace-a"},"state":"idle","last_seen_at":null,"pinned":true,"retention_state":"active","implementation":{"kind":"runtime","display_hint":"Runtime Worker"},"capabilities":{"can_stop":true,"can_spawn_followup":false},"diagnostics":[]},"diagnostics":[]}"#;
|
||||
let (base_url, request, handle) = one_response_server("200 OK", body);
|
||||
let client = BackendWorkspaceProductClient::new_with_access_token(
|
||||
base_url,
|
||||
|
||||
@@ -557,7 +557,6 @@ pub enum SubscriptionWorkerState {
|
||||
Running,
|
||||
Paused,
|
||||
Stopped,
|
||||
Cancelled,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
@@ -1110,6 +1109,25 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn worker_subscription_state_has_exactly_four_lifecycle_values() {
|
||||
for (state, wire) in [
|
||||
(SubscriptionWorkerState::Idle, "idle"),
|
||||
(SubscriptionWorkerState::Running, "running"),
|
||||
(SubscriptionWorkerState::Paused, "paused"),
|
||||
(SubscriptionWorkerState::Stopped, "stopped"),
|
||||
] {
|
||||
assert_eq!(
|
||||
serde_json::to_value(state).unwrap(),
|
||||
serde_json::json!(wire)
|
||||
);
|
||||
}
|
||||
assert!(
|
||||
serde_json::from_value::<SubscriptionWorkerState>(serde_json::json!("cancelled"))
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn client_selector_has_no_workspace_scope_field() {
|
||||
let json = serde_json::to_value(EventSubscriptionSelector::WorkspaceWorkers).unwrap();
|
||||
|
||||
@@ -195,7 +195,7 @@ async fn run_and_persist(
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
agen::EngineRunExit::Interrupted(agen::StopReason::LimitReached) => {
|
||||
agen::EngineRunExit::Interrupted(agen::RunInterruptionReason::LimitReached) => {
|
||||
session_store::save_run_completed(
|
||||
store,
|
||||
session_id,
|
||||
|
||||
@@ -274,6 +274,10 @@ pub struct CreateWorkerRequest {
|
||||
}
|
||||
|
||||
/// Worker lifecycle status for the in-memory embedded runtime.
|
||||
///
|
||||
/// Run termination details are carried separately by the Worker protocol. In
|
||||
/// particular, cancellation returns a Worker to `Idle`; it is not a lifecycle
|
||||
/// state of its own.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum WorkerStatus {
|
||||
@@ -281,7 +285,6 @@ pub enum WorkerStatus {
|
||||
Running,
|
||||
Paused,
|
||||
Stopped,
|
||||
Cancelled,
|
||||
}
|
||||
|
||||
impl WorkerStatus {
|
||||
@@ -290,6 +293,13 @@ impl WorkerStatus {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub(crate) enum WorkerRestoreIntent {
|
||||
Automatic,
|
||||
Explicit,
|
||||
}
|
||||
|
||||
/// Lightweight catalog row.
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct WorkerSummary {
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
use crate::catalog::{CreateWorkerRequest, WorkingDirectoryStatus};
|
||||
use crate::catalog::{
|
||||
CreateWorkerRequest, WorkerRestoreIntent, WorkerStatus, WorkingDirectoryStatus,
|
||||
};
|
||||
use crate::config_bundle::ConfigBundle;
|
||||
use crate::diagnostics::{DiagnosticSeverity, RuntimeDiagnostic};
|
||||
use crate::error::RuntimeError;
|
||||
@@ -13,7 +15,7 @@ use std::io::{BufReader, Write};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
const SCHEMA_VERSION: u32 = 3;
|
||||
const SCHEMA_VERSION: u32 = 4;
|
||||
const RUNTIME_FILE: &str = "runtime.json";
|
||||
const WORKERS_DIR: &str = "workers";
|
||||
const WORKER_FILE: &str = "worker.json";
|
||||
@@ -274,13 +276,24 @@ pub(crate) struct PersistedRuntimeState {
|
||||
pub(crate) diagnostics: Vec<RuntimeDiagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub(crate) struct PersistedWorkerExecutionBinding {
|
||||
pub(crate) run_generation: u64,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub(crate) struct PersistedWorkerExecution {
|
||||
pub(crate) binding: Option<PersistedWorkerExecutionBinding>,
|
||||
pub(crate) restore_intent: WorkerRestoreIntent,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub(crate) struct PersistedWorkerRecord {
|
||||
pub(crate) worker_ref: WorkerRef,
|
||||
pub(crate) worker_id: WorkerId,
|
||||
pub(crate) request: CreateWorkerRequest,
|
||||
/// Last generation durably reserved for this Worker's execution.
|
||||
pub(crate) run_generation: u64,
|
||||
pub(crate) status: WorkerStatus,
|
||||
pub(crate) execution: PersistedWorkerExecution,
|
||||
pub(crate) workspace_id: Option<String>,
|
||||
pub(crate) working_directory: Option<WorkingDirectoryStatus>,
|
||||
}
|
||||
@@ -357,8 +370,8 @@ fn plan_runtime_store_migration(
|
||||
format!("Runtime store schema version {schema_version} is out of range"),
|
||||
)
|
||||
})?;
|
||||
let staging = migration_sibling(root, "schema-v3-staging")?;
|
||||
let backup = migration_sibling(root, "pre-schema-v3-backup")?;
|
||||
let staging = migration_sibling(root, "schema-v4-staging")?;
|
||||
let backup = migration_sibling(root, "pre-schema-v4-backup")?;
|
||||
if staging.exists() || backup.exists() {
|
||||
return Err(runtime_store_corrupt(
|
||||
root,
|
||||
@@ -384,11 +397,11 @@ fn plan_runtime_store_migration(
|
||||
};
|
||||
return Ok((plan, Vec::new()));
|
||||
}
|
||||
if !matches!(current_schema_version, 1 | 2) {
|
||||
if current_schema_version != 3 {
|
||||
return Err(runtime_store_corrupt(
|
||||
&runtime_path,
|
||||
format!(
|
||||
"unsupported Runtime store schema version {schema_version}; expected 1, 2, or {SCHEMA_VERSION}"
|
||||
"unsupported Runtime store schema version {schema_version}; expected 3 or {SCHEMA_VERSION}"
|
||||
),
|
||||
));
|
||||
}
|
||||
@@ -448,7 +461,7 @@ fn plan_runtime_store_migration(
|
||||
let worker_id = name.parse::<WorkerId>().map_err(|_| {
|
||||
runtime_store_corrupt(
|
||||
&source_dir,
|
||||
format!("schema-v2 Worker directory name must be a UUIDv7, found {name}"),
|
||||
format!("pre-v4 Worker directory name must be a UUIDv7, found {name}"),
|
||||
)
|
||||
})?;
|
||||
(worker_id, None, None)
|
||||
@@ -610,7 +623,7 @@ fn migrate_worker_document(
|
||||
snapshot_path: &Path,
|
||||
) -> Result<serde_json::Value, RuntimeError> {
|
||||
if source_schema_version == 1 {
|
||||
return migrate_v1_worker_document(
|
||||
document = migrate_v1_worker_document(
|
||||
document,
|
||||
mapping.ok_or_else(|| {
|
||||
runtime_store_corrupt(
|
||||
@@ -619,7 +632,7 @@ fn migrate_worker_document(
|
||||
)
|
||||
})?,
|
||||
snapshot_path,
|
||||
);
|
||||
)?;
|
||||
}
|
||||
let object = document.as_object_mut().ok_or_else(|| {
|
||||
runtime_store_corrupt(
|
||||
@@ -627,10 +640,46 @@ fn migrate_worker_document(
|
||||
"Worker snapshot must be an object".to_string(),
|
||||
)
|
||||
})?;
|
||||
let run_generation = object
|
||||
.remove("run_generation")
|
||||
.map(|value| {
|
||||
value.as_u64().ok_or_else(|| {
|
||||
runtime_store_corrupt(
|
||||
snapshot_path,
|
||||
"Worker snapshot run_generation must be an unsigned integer".to_string(),
|
||||
)
|
||||
})
|
||||
})
|
||||
.transpose()?
|
||||
.filter(|generation| *generation > 0);
|
||||
let legacy_execution = object.remove("execution");
|
||||
if !object.contains_key("working_directory") {
|
||||
if let Some(working_directory) = legacy_execution
|
||||
.as_ref()
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.and_then(|execution| execution.get("working_directory"))
|
||||
.cloned()
|
||||
{
|
||||
object.insert("working_directory".to_string(), working_directory);
|
||||
}
|
||||
}
|
||||
object.insert(
|
||||
"schema_version".to_string(),
|
||||
serde_json::Value::from(SCHEMA_VERSION),
|
||||
);
|
||||
object.insert(
|
||||
"status".to_string(),
|
||||
serde_json::Value::String("stopped".to_string()),
|
||||
);
|
||||
object.insert(
|
||||
"execution".to_string(),
|
||||
serde_json::json!({
|
||||
"binding": run_generation.map(|run_generation| {
|
||||
serde_json::json!({ "run_generation": run_generation })
|
||||
}),
|
||||
"restore_intent": "explicit",
|
||||
}),
|
||||
);
|
||||
Ok(document)
|
||||
}
|
||||
|
||||
@@ -1005,8 +1054,8 @@ fn migrate_runtime_store(
|
||||
if !plan.migration_required {
|
||||
return Ok(plan);
|
||||
}
|
||||
let staging = migration_sibling(root, "schema-v3-staging")?;
|
||||
let backup = migration_sibling(root, "pre-schema-v3-backup")?;
|
||||
let staging = migration_sibling(root, "schema-v4-staging")?;
|
||||
let backup = migration_sibling(root, "pre-schema-v4-backup")?;
|
||||
if staging.exists() || backup.exists() {
|
||||
return Err(runtime_store_corrupt(
|
||||
root,
|
||||
@@ -1236,22 +1285,12 @@ struct WorkerSnapshot {
|
||||
worker_ref: WorkerRef,
|
||||
worker_id: WorkerId,
|
||||
request: CreateWorkerRequest,
|
||||
#[serde(default)]
|
||||
run_generation: u64,
|
||||
status: WorkerStatus,
|
||||
execution: PersistedWorkerExecution,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
workspace_id: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
working_directory: Option<WorkingDirectoryStatus>,
|
||||
/// One-way migration input for schema-v1 snapshots. New snapshots never
|
||||
/// write the removed execution projection.
|
||||
#[serde(default, rename = "execution", skip_serializing)]
|
||||
legacy_execution: Option<LegacyWorkerExecutionProjection>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
struct LegacyWorkerExecutionProjection {
|
||||
#[serde(default)]
|
||||
working_directory: Option<WorkingDirectoryStatus>,
|
||||
}
|
||||
|
||||
impl WorkerSnapshot {
|
||||
@@ -1261,10 +1300,10 @@ impl WorkerSnapshot {
|
||||
worker_ref: worker.worker_ref.clone(),
|
||||
worker_id: worker.worker_id.clone(),
|
||||
request: worker.request.clone(),
|
||||
run_generation: worker.run_generation,
|
||||
status: worker.status,
|
||||
execution: worker.execution.clone(),
|
||||
workspace_id: worker.workspace_id.clone(),
|
||||
working_directory: worker.working_directory.clone(),
|
||||
legacy_execution: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1289,6 +1328,51 @@ impl WorkerSnapshot {
|
||||
),
|
||||
});
|
||||
}
|
||||
match (self.status, self.execution.restore_intent) {
|
||||
(status, WorkerRestoreIntent::Automatic) if status.is_active() => {
|
||||
let Some(binding) = self.execution.binding.as_ref() else {
|
||||
return Err(RuntimeError::StoreCorrupt {
|
||||
operation: "read worker snapshot",
|
||||
path: path.to_path_buf(),
|
||||
message: "automatic restore intent requires an execution binding"
|
||||
.to_string(),
|
||||
});
|
||||
};
|
||||
if binding.run_generation == 0 {
|
||||
return Err(RuntimeError::StoreCorrupt {
|
||||
operation: "read worker snapshot",
|
||||
path: path.to_path_buf(),
|
||||
message: "execution binding run_generation must be greater than zero"
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
(WorkerStatus::Stopped, WorkerRestoreIntent::Explicit) => {
|
||||
if self
|
||||
.execution
|
||||
.binding
|
||||
.as_ref()
|
||||
.is_some_and(|binding| binding.run_generation == 0)
|
||||
{
|
||||
return Err(RuntimeError::StoreCorrupt {
|
||||
operation: "read worker snapshot",
|
||||
path: path.to_path_buf(),
|
||||
message: "execution binding run_generation must be greater than zero"
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
_ => {
|
||||
return Err(RuntimeError::StoreCorrupt {
|
||||
operation: "read worker snapshot",
|
||||
path: path.to_path_buf(),
|
||||
message: format!(
|
||||
"worker status {:?} conflicts with restore intent {:?}",
|
||||
self.status, self.execution.restore_intent
|
||||
),
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1303,12 +1387,10 @@ impl WorkerSnapshot {
|
||||
worker_ref: self.worker_ref,
|
||||
worker_id: self.worker_id,
|
||||
request: self.request,
|
||||
run_generation: self.run_generation,
|
||||
status: self.status,
|
||||
execution: self.execution,
|
||||
workspace_id,
|
||||
working_directory: self.working_directory.or_else(|| {
|
||||
self.legacy_execution
|
||||
.and_then(|execution| execution.working_directory)
|
||||
}),
|
||||
working_directory: self.working_directory,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1028,14 +1028,14 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migration_dry_run_accepts_real_v1_document_without_workers_field() {
|
||||
fn migration_dry_run_accepts_previous_schema_document_without_workers_field() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let root = temp.path().join("runtime");
|
||||
std::fs::create_dir_all(root.join("workers")).unwrap();
|
||||
std::fs::write(
|
||||
root.join("runtime.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"schema_version": 1,
|
||||
"schema_version": 3,
|
||||
"display_name": "local",
|
||||
"backend": "fs_store",
|
||||
"status": "running",
|
||||
@@ -1067,14 +1067,14 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migration_dry_run_rejects_v1_document_that_cannot_decode_as_v3() {
|
||||
fn migration_dry_run_rejects_previous_schema_document_that_cannot_decode_as_v4() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let root = temp.path().join("runtime");
|
||||
std::fs::create_dir_all(root.join("workers")).unwrap();
|
||||
std::fs::write(
|
||||
root.join("runtime.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"schema_version": 1,
|
||||
"schema_version": 3,
|
||||
"display_name": "local",
|
||||
"backend": "fs_store",
|
||||
"status": 3,
|
||||
|
||||
@@ -43,7 +43,6 @@ pub struct RuntimeSummary {
|
||||
pub worker_count: usize,
|
||||
pub active_worker_count: usize,
|
||||
pub stopped_worker_count: usize,
|
||||
pub cancelled_worker_count: usize,
|
||||
pub diagnostic_count: usize,
|
||||
#[serde(default = "unknown_platform_component")]
|
||||
pub os: String,
|
||||
|
||||
@@ -303,7 +303,12 @@ impl FsWorkerRetentionProvider {
|
||||
));
|
||||
continue;
|
||||
}
|
||||
match self.inventory(workspace_id, runtime_id, worker_id, snapshot.run_generation) {
|
||||
match self.inventory(
|
||||
workspace_id,
|
||||
runtime_id,
|
||||
worker_id,
|
||||
snapshot.run_generation(),
|
||||
) {
|
||||
Ok(item) => workers.push(item),
|
||||
Err(_) => diagnostics.push(runtime_aggregate_diagnostic(
|
||||
&bounded_id,
|
||||
@@ -388,10 +393,11 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
|
||||
if worker.workspace_id.as_deref() != Some(workspace_id) {
|
||||
return Err(RuntimeError::WorkerNotFound { worker_id });
|
||||
}
|
||||
if worker.run_generation != run_generation {
|
||||
let current_run_generation = worker.run_generation();
|
||||
if current_run_generation != run_generation {
|
||||
return Err(RuntimeError::InvalidRequest(format!(
|
||||
"Worker retention inventory expected generation {run_generation}, current generation is {}",
|
||||
worker.run_generation
|
||||
current_run_generation
|
||||
)));
|
||||
}
|
||||
let session_dir = worker_dir.join("session");
|
||||
@@ -498,10 +504,11 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
|
||||
worker_id: request.worker_id,
|
||||
});
|
||||
}
|
||||
if snapshot.run_generation != request.expected_run_generation {
|
||||
let run_generation = snapshot.run_generation();
|
||||
if run_generation != request.expected_run_generation {
|
||||
return Err(RuntimeError::InvalidRequest(format!(
|
||||
"Worker retention plan expected generation {}, current generation is {}",
|
||||
request.expected_run_generation, snapshot.run_generation
|
||||
request.expected_run_generation, run_generation
|
||||
)));
|
||||
}
|
||||
|
||||
@@ -568,10 +575,29 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
|
||||
struct WorkerGenerationSnapshot {
|
||||
#[serde(default)]
|
||||
workspace_id: Option<String>,
|
||||
#[serde(default)]
|
||||
execution: WorkerGenerationExecution,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct WorkerGenerationExecution {
|
||||
binding: Option<WorkerGenerationBinding>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct WorkerGenerationBinding {
|
||||
run_generation: u64,
|
||||
}
|
||||
|
||||
impl WorkerGenerationSnapshot {
|
||||
fn run_generation(&self) -> u64 {
|
||||
self.execution
|
||||
.binding
|
||||
.as_ref()
|
||||
.map(|binding| binding.run_generation)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct CanonicalSessionManifest {
|
||||
session_id: String,
|
||||
@@ -1264,7 +1290,10 @@ mod tests {
|
||||
let worker = root.join("workers").join(worker_id.to_string());
|
||||
write_json(
|
||||
&worker.join("worker.json"),
|
||||
&serde_json::json!({"workspace_id": "workspace-a", "run_generation": generation}),
|
||||
&serde_json::json!({
|
||||
"workspace_id": "workspace-a",
|
||||
"execution": {"binding": {"run_generation": generation}}
|
||||
}),
|
||||
);
|
||||
write_json(
|
||||
&worker.join("session/session.json"),
|
||||
@@ -1462,7 +1491,10 @@ mod tests {
|
||||
.join("workers")
|
||||
.join(other_worker.to_string())
|
||||
.join("worker.json"),
|
||||
&serde_json::json!({"workspace_id": "other-workspace", "run_generation": 1}),
|
||||
&serde_json::json!({
|
||||
"workspace_id": "other-workspace",
|
||||
"execution": {"binding": {"run_generation": 1}}
|
||||
}),
|
||||
);
|
||||
fs::create_dir_all(temp.path().join("workers/not-a-worker")).unwrap();
|
||||
fs::write(
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
use crate::catalog::{
|
||||
ConfigBundleRef, CreateWorkerRequest, ProfileSelector, RepositoryRefObservation,
|
||||
RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck, WorkerStatus, WorkerSummary,
|
||||
WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest,
|
||||
RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck, WorkerRestoreIntent,
|
||||
WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest,
|
||||
WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, WorkspaceApiRef,
|
||||
};
|
||||
use crate::config_bundle::{
|
||||
@@ -18,7 +18,8 @@ use crate::execution::{
|
||||
};
|
||||
#[cfg(feature = "fs-store")]
|
||||
use crate::fs_store::{
|
||||
FsRuntimeStore, FsRuntimeStoreOptions, PersistedRuntimeState, PersistedWorkerRecord,
|
||||
FsRuntimeStore, FsRuntimeStoreOptions, PersistedRuntimeState, PersistedWorkerExecution,
|
||||
PersistedWorkerExecutionBinding, PersistedWorkerRecord,
|
||||
};
|
||||
use crate::identity::{WorkerId, WorkerRef};
|
||||
use crate::interaction::{WorkerInput, WorkerInputKind, WorkerInteractionAck};
|
||||
@@ -230,14 +231,12 @@ impl Runtime {
|
||||
let state = self.lock()?;
|
||||
let mut active_worker_count = 0;
|
||||
let mut stopped_worker_count = 0;
|
||||
let mut cancelled_worker_count = 0;
|
||||
for worker in state.workers.values() {
|
||||
match worker.status {
|
||||
WorkerStatus::Idle | WorkerStatus::Running | WorkerStatus::Paused => {
|
||||
active_worker_count += 1;
|
||||
}
|
||||
WorkerStatus::Stopped => stopped_worker_count += 1,
|
||||
WorkerStatus::Cancelled => cancelled_worker_count += 1,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -248,7 +247,6 @@ impl Runtime {
|
||||
worker_count: state.workers.len(),
|
||||
active_worker_count,
|
||||
stopped_worker_count,
|
||||
cancelled_worker_count,
|
||||
diagnostic_count: state.diagnostics.len(),
|
||||
os: std::env::consts::OS.to_string(),
|
||||
arch: std::env::consts::ARCH.to_string(),
|
||||
@@ -344,17 +342,6 @@ impl Runtime {
|
||||
return Ok(());
|
||||
}
|
||||
state.status = RuntimeStatus::Stopped;
|
||||
let mut stopped = Vec::new();
|
||||
for (worker_id, worker) in &mut state.workers {
|
||||
if worker.status.is_active() {
|
||||
worker.status = WorkerStatus::Stopped;
|
||||
worker.internal_workers.clear();
|
||||
stopped.push(*worker_id);
|
||||
}
|
||||
}
|
||||
for worker_id in stopped {
|
||||
state.publish_worker_upsert(worker_id)?;
|
||||
}
|
||||
state.persist_runtime_snapshot()?;
|
||||
state.persist_workers()?;
|
||||
Ok(())
|
||||
@@ -716,6 +703,8 @@ impl Runtime {
|
||||
workspace_id: scope.map(|scope| scope.workspace_id.clone()),
|
||||
request: durable_request,
|
||||
run_generation: 1,
|
||||
execution_bound: true,
|
||||
restore_intent: WorkerRestoreIntent::Explicit,
|
||||
working_directory: None,
|
||||
execution_handle: None,
|
||||
internal_workers: BTreeMap::new(),
|
||||
@@ -1052,12 +1041,6 @@ impl Runtime {
|
||||
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
|
||||
)));
|
||||
}
|
||||
(
|
||||
worker.request.clone(),
|
||||
worker.working_directory.clone(),
|
||||
@@ -1070,7 +1053,12 @@ impl Runtime {
|
||||
message: "runtime has no execution backend".to_string(),
|
||||
}
|
||||
})?;
|
||||
state.worker_mut(worker_ref)?.run_generation = run_generation;
|
||||
{
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.run_generation = run_generation;
|
||||
worker.execution_bound = true;
|
||||
worker.restore_intent = WorkerRestoreIntent::Automatic;
|
||||
}
|
||||
state.persist_worker(&worker_ref.worker_id)?;
|
||||
let workspace_scope = worker_request.workspace_api.as_ref().and_then(|api| {
|
||||
state
|
||||
@@ -1124,18 +1112,20 @@ impl Runtime {
|
||||
}
|
||||
|
||||
fn ensure_worker_execution(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> {
|
||||
let has_handle = {
|
||||
let state = self.lock()?;
|
||||
state
|
||||
.worker(worker_ref)?
|
||||
.execution_handle
|
||||
.as_ref()
|
||||
.is_some()
|
||||
};
|
||||
if has_handle {
|
||||
let state = self.lock()?;
|
||||
let worker = state.worker(worker_ref)?;
|
||||
if worker.execution_handle.is_some() {
|
||||
return Ok(());
|
||||
}
|
||||
self.restore_worker(worker_ref).map(|_| ())
|
||||
let message = if worker.status == WorkerStatus::Stopped {
|
||||
"stopped worker requires an explicit restore"
|
||||
} else {
|
||||
"worker has no live execution handle"
|
||||
};
|
||||
Err(RuntimeError::WorkerExecutionUnavailable {
|
||||
worker_id: worker_ref.worker_id,
|
||||
message: message.to_string(),
|
||||
})
|
||||
}
|
||||
|
||||
/// Accept input into a Worker through a workspace-scoped Runtime authorization context.
|
||||
@@ -1488,7 +1478,9 @@ impl Runtime {
|
||||
let detail = {
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.execution_handle = Some(handle);
|
||||
worker.execution_bound = true;
|
||||
worker.status = worker_status_from_run_state(run_state);
|
||||
worker.restore_intent = restore_intent_for_status(worker.status);
|
||||
worker.working_directory = working_directory;
|
||||
worker.detail()
|
||||
};
|
||||
@@ -1517,8 +1509,13 @@ impl Runtime {
|
||||
) -> Result<(), RuntimeError> {
|
||||
let mut state = self.lock()?;
|
||||
if result.is_accepted() {
|
||||
state.worker_mut(worker_ref)?.status = worker_status_from_run_state(result.run_state);
|
||||
let status = worker_status_from_run_state(result.run_state);
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.status = status;
|
||||
worker.restore_intent = restore_intent_for_status(status);
|
||||
state.publish_worker_upsert(worker_ref.worker_id)?;
|
||||
state.persist_runtime_snapshot()?;
|
||||
state.persist_worker(&worker_ref.worker_id)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -1601,15 +1598,26 @@ impl Runtime {
|
||||
self.cancel_worker(worker_ref, reason)
|
||||
}
|
||||
|
||||
/// Cancel a Worker. Repeated cancels are idempotent.
|
||||
/// Cancel the current run while keeping the Worker session available.
|
||||
pub fn cancel_worker(
|
||||
&self,
|
||||
worker_ref: &WorkerRef,
|
||||
reason: Option<String>,
|
||||
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
||||
let current = {
|
||||
let state = self.lock()?;
|
||||
state.ensure_running()?;
|
||||
state.worker(worker_ref)?.status
|
||||
};
|
||||
if matches!(current, WorkerStatus::Idle | WorkerStatus::Stopped) {
|
||||
return Ok(WorkerLifecycleAck {
|
||||
worker_ref: worker_ref.clone(),
|
||||
status: current,
|
||||
});
|
||||
}
|
||||
self.dispatch_lifecycle_to_backend(worker_ref, WorkerExecutionOperation::Cancel)?;
|
||||
let _ = reason;
|
||||
self.transition_worker(worker_ref, WorkerStatus::Cancelled)
|
||||
self.transition_worker_preserving_execution(worker_ref, WorkerStatus::Idle)
|
||||
}
|
||||
|
||||
/// Delete a non-running Worker through a workspace-scoped Runtime authorization context.
|
||||
@@ -1759,6 +1767,10 @@ impl Runtime {
|
||||
if status_changed || activity_changed {
|
||||
state.publish_worker_upsert(worker_ref.worker_id)?;
|
||||
}
|
||||
if status_changed {
|
||||
state.persist_runtime_snapshot()?;
|
||||
state.persist_worker(&worker_ref.worker_id)?;
|
||||
}
|
||||
let event = state.push_worker_observation_event(worker_ref.clone(), payload);
|
||||
Ok(event)
|
||||
}
|
||||
@@ -1791,6 +1803,26 @@ impl Runtime {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn transition_worker_preserving_execution(
|
||||
&self,
|
||||
worker_ref: &WorkerRef,
|
||||
status: WorkerStatus,
|
||||
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
||||
let mut state = self.lock()?;
|
||||
state.ensure_running()?;
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.status = status;
|
||||
worker.restore_intent = restore_intent_for_status(status);
|
||||
let status = worker.status;
|
||||
state.publish_worker_upsert(worker_ref.worker_id)?;
|
||||
state.persist_runtime_snapshot()?;
|
||||
state.persist_worker(&worker_ref.worker_id)?;
|
||||
Ok(WorkerLifecycleAck {
|
||||
worker_ref: worker_ref.clone(),
|
||||
status,
|
||||
})
|
||||
}
|
||||
|
||||
fn transition_worker(
|
||||
&self,
|
||||
worker_ref: &WorkerRef,
|
||||
@@ -1800,18 +1832,9 @@ impl Runtime {
|
||||
state.ensure_running()?;
|
||||
state.ensure_worker_ref(worker_ref)?;
|
||||
|
||||
{
|
||||
let worker = state.worker(worker_ref)?;
|
||||
if !worker.status.is_active() {
|
||||
return Ok(WorkerLifecycleAck {
|
||||
worker_ref: worker_ref.clone(),
|
||||
status: worker.status,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.status = status;
|
||||
worker.restore_intent = restore_intent_for_status(status);
|
||||
worker.execution_handle = None;
|
||||
worker.internal_workers.clear();
|
||||
let status = worker.status;
|
||||
@@ -1867,7 +1890,12 @@ impl Runtime {
|
||||
let worker_ids = state
|
||||
.workers
|
||||
.values()
|
||||
.filter(|worker| worker.execution_handle.is_none())
|
||||
.filter(|worker| {
|
||||
worker.execution_handle.is_none()
|
||||
&& worker.execution_bound
|
||||
&& worker.status.is_active()
|
||||
&& worker.restore_intent == WorkerRestoreIntent::Automatic
|
||||
})
|
||||
.map(|worker| worker.worker_id)
|
||||
.collect::<Vec<_>>();
|
||||
let mut candidates = Vec::with_capacity(worker_ids.len());
|
||||
@@ -1958,7 +1986,9 @@ impl Runtime {
|
||||
{
|
||||
let worker = state.worker_mut(worker_ref)?;
|
||||
worker.execution_handle = Some(handle);
|
||||
worker.execution_bound = true;
|
||||
worker.status = worker_status_from_run_state(run_state);
|
||||
worker.restore_intent = restore_intent_for_status(worker.status);
|
||||
worker.working_directory = working_directory;
|
||||
}
|
||||
state.publish_worker_upsert(worker_ref.worker_id)?;
|
||||
@@ -2227,15 +2257,23 @@ impl RuntimeState {
|
||||
let diagnostics = persisted.diagnostics;
|
||||
let next_diagnostic_id = persisted.next_diagnostic_id;
|
||||
for (worker_id, worker) in persisted.workers {
|
||||
let run_generation = worker
|
||||
.execution
|
||||
.binding
|
||||
.as_ref()
|
||||
.map(|binding| binding.run_generation)
|
||||
.unwrap_or(0);
|
||||
workers.insert(
|
||||
worker_id,
|
||||
WorkerRecord {
|
||||
worker_ref: worker.worker_ref,
|
||||
worker_id: worker.worker_id,
|
||||
status: WorkerStatus::Stopped,
|
||||
status: worker.status,
|
||||
workspace_id: worker.workspace_id,
|
||||
request: worker.request,
|
||||
run_generation: worker.run_generation,
|
||||
run_generation,
|
||||
execution_bound: worker.execution.binding.is_some(),
|
||||
restore_intent: worker.execution.restore_intent,
|
||||
working_directory: worker.working_directory,
|
||||
execution_handle: None,
|
||||
internal_workers: BTreeMap::new(),
|
||||
@@ -2708,9 +2746,11 @@ impl RuntimeState {
|
||||
let worker = self.worker_mut(worker_ref)?;
|
||||
worker.execution_handle = None;
|
||||
worker.status = WorkerStatus::Stopped;
|
||||
worker.restore_intent = WorkerRestoreIntent::Explicit;
|
||||
worker.internal_workers.clear();
|
||||
self.publish_worker_upsert(worker_ref.worker_id)?;
|
||||
self.persist_runtime_snapshot()?;
|
||||
self.persist_worker(&worker_ref.worker_id)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -2926,6 +2966,7 @@ impl RuntimeState {
|
||||
if let Some(next_status) = next_status {
|
||||
let changed = worker.status != next_status;
|
||||
worker.status = next_status;
|
||||
worker.restore_intent = restore_intent_for_status(next_status);
|
||||
changed
|
||||
} else {
|
||||
false
|
||||
@@ -2947,6 +2988,8 @@ struct WorkerRecord {
|
||||
workspace_id: Option<String>,
|
||||
request: CreateWorkerRequest,
|
||||
run_generation: u64,
|
||||
execution_bound: bool,
|
||||
restore_intent: WorkerRestoreIntent,
|
||||
working_directory: Option<CatalogWorkingDirectoryStatus>,
|
||||
execution_handle: Option<WorkerExecutionHandle>,
|
||||
internal_workers: BTreeMap<String, InternalWorkerActivity>,
|
||||
@@ -2991,13 +3034,29 @@ impl WorkerRecord {
|
||||
worker_ref: self.worker_ref.clone(),
|
||||
worker_id: self.worker_id.clone(),
|
||||
request: self.request.clone(),
|
||||
run_generation: self.run_generation,
|
||||
status: self.status,
|
||||
execution: PersistedWorkerExecution {
|
||||
binding: self
|
||||
.execution_bound
|
||||
.then_some(PersistedWorkerExecutionBinding {
|
||||
run_generation: self.run_generation,
|
||||
}),
|
||||
restore_intent: self.restore_intent,
|
||||
},
|
||||
workspace_id: self.workspace_id.clone(),
|
||||
working_directory: self.working_directory.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn restore_intent_for_status(status: WorkerStatus) -> WorkerRestoreIntent {
|
||||
if status.is_active() {
|
||||
WorkerRestoreIntent::Automatic
|
||||
} else {
|
||||
WorkerRestoreIntent::Explicit
|
||||
}
|
||||
}
|
||||
|
||||
fn worker_status_from_run_state(run_state: WorkerExecutionRunState) -> WorkerStatus {
|
||||
match run_state {
|
||||
WorkerExecutionRunState::Idle => WorkerStatus::Idle,
|
||||
@@ -3212,7 +3271,6 @@ fn subscription_worker_state(status: WorkerStatus) -> SubscriptionWorkerState {
|
||||
WorkerStatus::Running => SubscriptionWorkerState::Running,
|
||||
WorkerStatus::Paused => SubscriptionWorkerState::Paused,
|
||||
WorkerStatus::Stopped => SubscriptionWorkerState::Stopped,
|
||||
WorkerStatus::Cancelled => SubscriptionWorkerState::Cancelled,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4917,15 +4975,22 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn input_restores_stopped_worker_without_persisted_connection_state() {
|
||||
fn stopped_worker_rejects_input_until_explicitly_restored() {
|
||||
let (runtime, backend) = runtime_and_backend();
|
||||
let detail = runtime
|
||||
.create_worker(task_request("restore on input"))
|
||||
.create_worker(task_request("restore explicitly"))
|
||||
.unwrap();
|
||||
runtime
|
||||
.send_protocol_method(&detail.worker_ref, Method::Shutdown)
|
||||
.unwrap();
|
||||
|
||||
assert!(matches!(
|
||||
runtime.send_input(&detail.worker_ref, WorkerInput::user("do not wake")),
|
||||
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
||||
));
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
||||
|
||||
runtime.restore_worker(&detail.worker_ref).unwrap();
|
||||
runtime
|
||||
.send_input(&detail.worker_ref, WorkerInput::user("wake up"))
|
||||
.unwrap();
|
||||
@@ -5073,10 +5138,13 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stop_and_cancel_workers_update_summary() {
|
||||
fn stop_and_cancel_workers_keep_four_state_summary() {
|
||||
let runtime = runtime_with_backend();
|
||||
let stopped = runtime.create_worker(task_request("stop me")).unwrap();
|
||||
let cancelled = runtime.create_worker(task_request("cancel me")).unwrap();
|
||||
runtime
|
||||
.send_input(&cancelled.worker_ref, WorkerInput::user("start"))
|
||||
.unwrap();
|
||||
|
||||
let stop_ack = runtime
|
||||
.stop_worker(&stopped.worker_ref, Some("done".to_string()))
|
||||
@@ -5086,13 +5154,13 @@ mod tests {
|
||||
let cancel_ack = runtime
|
||||
.cancel_worker(&cancelled.worker_ref, Some("abort".to_string()))
|
||||
.unwrap();
|
||||
assert_eq!(cancel_ack.status, WorkerStatus::Cancelled);
|
||||
assert_eq!(cancel_ack.status, WorkerStatus::Idle);
|
||||
|
||||
let summary = runtime.summary().unwrap();
|
||||
assert_eq!(summary.worker_count, 2);
|
||||
assert_eq!(summary.active_worker_count, 0);
|
||||
assert_eq!(summary.active_worker_count, 1);
|
||||
assert_eq!(summary.stopped_worker_count, 1);
|
||||
assert_eq!(summary.cancelled_worker_count, 1);
|
||||
assert!(serde_json::from_value::<WorkerStatus>(serde_json::json!("cancelled")).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -5139,14 +5207,16 @@ mod tests {
|
||||
let summary = runtime.summary().unwrap();
|
||||
assert_eq!(summary.active_worker_count, 0);
|
||||
assert_eq!(summary.stopped_worker_count, 1);
|
||||
assert_eq!(summary.cancelled_worker_count, 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cancel_then_stop_preserves_cancelled_terminal_state() {
|
||||
fn cancel_then_stop_transitions_idle_session_to_stopped() {
|
||||
let runtime = runtime_with_backend();
|
||||
let worker = runtime
|
||||
.create_worker(task_request("stable cancelled"))
|
||||
.create_worker(task_request("cancel then stop"))
|
||||
.unwrap();
|
||||
runtime
|
||||
.send_input(&worker.worker_ref, WorkerInput::user("start"))
|
||||
.unwrap();
|
||||
|
||||
let cancel_ack = runtime
|
||||
@@ -5156,17 +5226,16 @@ mod tests {
|
||||
.stop_worker(&worker.worker_ref, Some("late stop".to_string()))
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(cancel_ack.status, WorkerStatus::Cancelled);
|
||||
assert_eq!(stop_ack.status, WorkerStatus::Cancelled);
|
||||
assert_eq!(cancel_ack.status, WorkerStatus::Idle);
|
||||
assert_eq!(stop_ack.status, WorkerStatus::Stopped);
|
||||
assert_eq!(
|
||||
runtime.worker_detail(&worker.worker_ref).unwrap().status,
|
||||
WorkerStatus::Cancelled
|
||||
WorkerStatus::Stopped
|
||||
);
|
||||
|
||||
let summary = runtime.summary().unwrap();
|
||||
assert_eq!(summary.active_worker_count, 0);
|
||||
assert_eq!(summary.stopped_worker_count, 0);
|
||||
assert_eq!(summary.cancelled_worker_count, 1);
|
||||
assert_eq!(summary.stopped_worker_count, 1);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
@@ -5194,239 +5263,38 @@ mod tests {
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_migrates_legacy_numeric_worker_identity_to_workspace_uuid() {
|
||||
let root = fs_store_root("worker-id-v1");
|
||||
let runtime_id = "arcadia";
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: runtime_id.to_string(),
|
||||
display_name: None,
|
||||
},
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
fn fs_store_rejects_schema_older_than_previous_release() {
|
||||
let root = fs_store_root("unsupported-old-schema");
|
||||
let runtime = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
})
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker_scoped(
|
||||
&RuntimeWorkspaceScope::new("workspace-a", "server"),
|
||||
scoped_task_request("legacy", "workspace-a"),
|
||||
)
|
||||
.unwrap();
|
||||
drop(runtime);
|
||||
|
||||
let current_dir = root.join("workers").join(worker.worker_id.to_string());
|
||||
let legacy_dir = root.join("workers").join("7");
|
||||
std::fs::rename(¤t_dir, &legacy_dir).unwrap();
|
||||
let worker_path = legacy_dir.join("worker.json");
|
||||
let mut worker_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
||||
worker_json["schema_version"] = serde_json::json!(1);
|
||||
worker_json["worker_id"] = serde_json::json!(7);
|
||||
worker_json["worker_ref"]["worker_id"] = serde_json::json!(7);
|
||||
let request = worker_json["request"].as_object_mut().unwrap();
|
||||
request.remove("worker_id");
|
||||
request.remove("create_fingerprint");
|
||||
request.insert("idempotency_key".to_string(), serde_json::Value::Null);
|
||||
request.insert(
|
||||
"idempotency_fingerprint".to_string(),
|
||||
serde_json::Value::Null,
|
||||
);
|
||||
std::fs::write(
|
||||
&worker_path,
|
||||
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
let legacy_worker_name = "worker-runtime-7";
|
||||
let legacy_manifest = manifest::WorkerManifest::from_toml(&format!(
|
||||
r#"
|
||||
[worker]
|
||||
name = "{legacy_worker_name}"
|
||||
|
||||
[model]
|
||||
scheme = "anthropic"
|
||||
model_id = "test-model"
|
||||
|
||||
[engine]
|
||||
|
||||
[[scope.allow]]
|
||||
target = "/tmp"
|
||||
permission = "write"
|
||||
"#,
|
||||
))
|
||||
.unwrap();
|
||||
std::fs::write(
|
||||
legacy_dir.join("metadata.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"worker_name": legacy_worker_name,
|
||||
"workspace_id": "workspace-a",
|
||||
"resolved_manifest_snapshot": legacy_manifest
|
||||
}))
|
||||
.unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let runtime_path = root.join("runtime.json");
|
||||
let mut runtime_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
||||
runtime_json["schema_version"] = serde_json::json!(1);
|
||||
runtime_json["workers"] = serde_json::json!({"legacy": "ignored"});
|
||||
runtime_json["next_worker_sequence"] = serde_json::json!(8);
|
||||
runtime_json["next_diagnostic_id"] = serde_json::json!(3);
|
||||
runtime_json["diagnostics"] = serde_json::json!([
|
||||
{
|
||||
"id": 1,
|
||||
"severity": "warning",
|
||||
"code": "mapped_legacy_worker",
|
||||
"message": "mapped diagnostic",
|
||||
"worker_ref": {"worker_id": 7}
|
||||
},
|
||||
{
|
||||
"id": 2,
|
||||
"severity": "warning",
|
||||
"code": "deleted_legacy_worker",
|
||||
"message": "unmapped diagnostic",
|
||||
"worker_ref": {"worker_id": 6}
|
||||
}
|
||||
]);
|
||||
runtime_json["schema_version"] = serde_json::json!(2);
|
||||
std::fs::write(
|
||||
&runtime_path,
|
||||
serde_json::to_vec_pretty(&runtime_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let run_dir = legacy_dir.join("runs").join("6");
|
||||
std::fs::create_dir_all(&run_dir).unwrap();
|
||||
let socket =
|
||||
std::os::unix::net::UnixListener::bind(run_dir.join("worker.sock")).unwrap();
|
||||
drop(socket);
|
||||
}
|
||||
|
||||
let runtime_options = crate::fs_store::FsRuntimeStoreOptions {
|
||||
let error = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: runtime_id.to_string(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
};
|
||||
let runtime_before_dry_run = std::fs::read(&runtime_path).unwrap();
|
||||
let plan = crate::fs_store::FsRuntimeStore::migration_plan(&runtime_options).unwrap();
|
||||
assert!(plan.migration_required);
|
||||
assert_eq!(plan.worker_count, 1);
|
||||
assert_eq!(plan.migrated_worker_aggregate_count, 1);
|
||||
assert_eq!(plan.migrated_diagnostic_worker_ref_count, 1);
|
||||
assert_eq!(plan.cleared_diagnostic_worker_ref_count, 1);
|
||||
assert_eq!(plan.mappings[0].legacy_worker_id, 7);
|
||||
#[cfg(unix)]
|
||||
assert_eq!(
|
||||
plan.excluded_ephemeral_paths,
|
||||
vec!["workers/7/runs/6/worker.sock"]
|
||||
);
|
||||
assert_eq!(
|
||||
std::fs::read(&runtime_path).unwrap(),
|
||||
runtime_before_dry_run
|
||||
);
|
||||
assert!(legacy_dir.exists());
|
||||
|
||||
let restored = Runtime::with_fs_store(runtime_options.clone()).unwrap();
|
||||
let expected = WorkerId::from_legacy_binding("workspace-a", runtime_id, 7);
|
||||
let detail = restored.worker_detail(&WorkerRef::new(expected)).unwrap();
|
||||
assert_eq!(detail.worker_id, expected);
|
||||
assert_eq!(detail.worker_ref.worker_id, expected);
|
||||
let expected_worker_dir = root.join("workers").join(expected.to_string());
|
||||
assert!(expected_worker_dir.exists());
|
||||
#[cfg(unix)]
|
||||
assert!(!expected_worker_dir.join("runs/6/worker.sock").exists());
|
||||
assert!(!legacy_dir.exists());
|
||||
let migrated_runtime: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
||||
assert_eq!(migrated_runtime["schema_version"], serde_json::json!(3));
|
||||
assert!(migrated_runtime.get("workers").is_none());
|
||||
assert!(migrated_runtime.get("next_worker_sequence").is_none());
|
||||
assert_eq!(
|
||||
migrated_runtime["diagnostics"][0]["worker_ref"]["worker_id"],
|
||||
serde_json::json!(expected.to_string())
|
||||
);
|
||||
})
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
migrated_runtime["diagnostics"][1]
|
||||
.get("worker_ref")
|
||||
.is_none()
|
||||
);
|
||||
let diagnostics = restored.diagnostics().unwrap();
|
||||
assert_eq!(
|
||||
diagnostics
|
||||
.iter()
|
||||
.find(|diagnostic| diagnostic.code == "mapped_legacy_worker")
|
||||
.and_then(|diagnostic| diagnostic.worker_ref.as_ref()),
|
||||
Some(&WorkerRef::new(expected))
|
||||
);
|
||||
assert!(
|
||||
diagnostics
|
||||
.iter()
|
||||
.find(|diagnostic| diagnostic.code == "deleted_legacy_worker")
|
||||
.is_some_and(|diagnostic| diagnostic.worker_ref.is_none())
|
||||
);
|
||||
let metadata_path = expected_worker_dir.join("metadata.json");
|
||||
let mut migrated_metadata: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&metadata_path).unwrap()).unwrap();
|
||||
let expected_worker_name = format!("worker-runtime-{expected}");
|
||||
assert_eq!(
|
||||
migrated_metadata["worker_name"],
|
||||
serde_json::json!(expected_worker_name)
|
||||
);
|
||||
assert_eq!(
|
||||
migrated_metadata["resolved_manifest_snapshot"]["worker"]["name"],
|
||||
serde_json::json!(expected_worker_name)
|
||||
error
|
||||
.to_string()
|
||||
.contains("unsupported Runtime store schema version 2; expected 3 or 4")
|
||||
);
|
||||
|
||||
drop(restored);
|
||||
|
||||
let mut schema_v2_runtime: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
||||
schema_v2_runtime["schema_version"] = serde_json::json!(2);
|
||||
std::fs::write(
|
||||
&runtime_path,
|
||||
serde_json::to_vec_pretty(&schema_v2_runtime).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
let migrated_worker_path = expected_worker_dir.join("worker.json");
|
||||
let mut schema_v2_worker: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&migrated_worker_path).unwrap()).unwrap();
|
||||
schema_v2_worker["schema_version"] = serde_json::json!(2);
|
||||
std::fs::write(
|
||||
&migrated_worker_path,
|
||||
serde_json::to_vec_pretty(&schema_v2_worker).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
migrated_metadata["worker_name"] = serde_json::json!(legacy_worker_name);
|
||||
migrated_metadata["resolved_manifest_snapshot"]["worker"]["name"] =
|
||||
serde_json::json!(legacy_worker_name);
|
||||
std::fs::write(
|
||||
&metadata_path,
|
||||
serde_json::to_vec_pretty(&migrated_metadata).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let recovery_plan =
|
||||
crate::fs_store::FsRuntimeStore::migration_plan(&runtime_options).unwrap();
|
||||
assert_eq!(recovery_plan.current_schema_version, 2);
|
||||
assert_eq!(recovery_plan.target_schema_version, 3);
|
||||
assert!(recovery_plan.migration_required);
|
||||
assert_eq!(recovery_plan.worker_count, 1);
|
||||
assert_eq!(recovery_plan.migrated_worker_aggregate_count, 1);
|
||||
assert!(recovery_plan.mappings.is_empty());
|
||||
|
||||
let recovered = Runtime::with_fs_store(runtime_options).unwrap();
|
||||
let recovered_metadata: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(metadata_path).unwrap()).unwrap();
|
||||
assert_eq!(
|
||||
recovered_metadata["worker_name"],
|
||||
serde_json::json!(expected_worker_name)
|
||||
);
|
||||
assert_eq!(
|
||||
recovered_metadata["resolved_manifest_snapshot"]["worker"]["name"],
|
||||
serde_json::json!(expected_worker_name)
|
||||
);
|
||||
drop(recovered);
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
@@ -5466,8 +5334,16 @@ mod tests {
|
||||
let worker_snapshot: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(worker_store_dir.join("worker.json")).unwrap())
|
||||
.unwrap();
|
||||
assert!(worker_snapshot.get("status").is_none());
|
||||
assert!(worker_snapshot.get("execution").is_none());
|
||||
assert_eq!(worker_snapshot["schema_version"], serde_json::json!(4));
|
||||
assert_eq!(worker_snapshot["status"], serde_json::json!("stopped"));
|
||||
assert_eq!(
|
||||
worker_snapshot["execution"]["binding"]["run_generation"],
|
||||
serde_json::json!(1)
|
||||
);
|
||||
assert_eq!(
|
||||
worker_snapshot["execution"]["restore_intent"],
|
||||
serde_json::json!("explicit")
|
||||
);
|
||||
assert!(!root.join("events.jsonl").exists());
|
||||
std::fs::write(
|
||||
worker_store_dir.join("observations.jsonl"),
|
||||
@@ -5624,8 +5500,8 @@ mod tests {
|
||||
display_name: None,
|
||||
})
|
||||
.unwrap();
|
||||
let stopped_worker = backendless.worker_detail(&worker.worker_ref).unwrap();
|
||||
assert_eq!(stopped_worker.status, WorkerStatus::Stopped);
|
||||
let persisted_worker = backendless.worker_detail(&worker.worker_ref).unwrap();
|
||||
assert_eq!(persisted_worker.status, WorkerStatus::Idle);
|
||||
drop(backendless);
|
||||
|
||||
let restoring_backend = Arc::new(TestExecutionBackend::default());
|
||||
@@ -5649,6 +5525,272 @@ mod tests {
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_automatically_restores_every_active_lifecycle_state() {
|
||||
for status in ["idle", "running", "paused"] {
|
||||
let root = fs_store_root(&format!("automatic-{status}"));
|
||||
let options = crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
};
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
options.clone(),
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker(task_request(&format!("restore {status}")))
|
||||
.unwrap();
|
||||
drop(runtime);
|
||||
|
||||
let worker_path = root
|
||||
.join("workers")
|
||||
.join(worker.worker_id.to_string())
|
||||
.join("worker.json");
|
||||
let mut worker_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
||||
worker_json["status"] = serde_json::json!(status);
|
||||
std::fs::write(
|
||||
&worker_path,
|
||||
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let backend = Arc::new(TestExecutionBackend::default());
|
||||
let restored =
|
||||
Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap();
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 1, "status={status}");
|
||||
assert_eq!(
|
||||
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
||||
WorkerStatus::Idle,
|
||||
"status={status}"
|
||||
);
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_current_schema_requires_lifecycle_authority() {
|
||||
let root = fs_store_root("current-schema-requires-lifecycle");
|
||||
let options = crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
};
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
options.clone(),
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker(task_request("missing lifecycle authority"))
|
||||
.unwrap();
|
||||
drop(runtime);
|
||||
|
||||
let worker_path = root
|
||||
.join("workers")
|
||||
.join(worker.worker_id.to_string())
|
||||
.join("worker.json");
|
||||
let mut worker_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
||||
worker_json.as_object_mut().unwrap().remove("status");
|
||||
std::fs::write(
|
||||
&worker_path,
|
||||
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let restored = Runtime::with_fs_store_and_execution_backend(
|
||||
options,
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
assert!(restored.list_workers().unwrap().is_empty());
|
||||
assert!(
|
||||
restored
|
||||
.diagnostics()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|diagnostic| diagnostic.code == "worker_snapshot_ignored")
|
||||
);
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_shutdown_preserves_active_worker_restore_intent() {
|
||||
let root = fs_store_root("shutdown-preserves-worker");
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
},
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker(task_request("preserve active worker"))
|
||||
.unwrap();
|
||||
|
||||
runtime.stop_runtime().unwrap();
|
||||
|
||||
let snapshot: serde_json::Value = serde_json::from_slice(
|
||||
&std::fs::read(
|
||||
root.join("workers")
|
||||
.join(worker.worker_id.to_string())
|
||||
.join("worker.json"),
|
||||
)
|
||||
.unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(snapshot["status"], serde_json::json!("idle"));
|
||||
assert_eq!(
|
||||
snapshot["execution"]["restore_intent"],
|
||||
serde_json::json!("automatic")
|
||||
);
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_stopped_worker_requires_explicit_restore() {
|
||||
let root = fs_store_root("stopped-explicit-restore");
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
},
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker(task_request("explicit restore only"))
|
||||
.unwrap();
|
||||
runtime
|
||||
.stop_worker(&worker.worker_ref, Some("operator stop".to_string()))
|
||||
.unwrap();
|
||||
drop(runtime);
|
||||
|
||||
let backend = Arc::new(TestExecutionBackend::default());
|
||||
let restored = Runtime::with_fs_store_and_execution_backend(
|
||||
crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
},
|
||||
backend.clone(),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
||||
assert_eq!(
|
||||
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
||||
WorkerStatus::Stopped
|
||||
);
|
||||
assert!(matches!(
|
||||
restored.send_input(&worker.worker_ref, WorkerInput::user("implicit restore")),
|
||||
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
||||
));
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
||||
|
||||
restored.restore_worker(&worker.worker_ref).unwrap();
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
|
||||
assert_eq!(
|
||||
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
||||
WorkerStatus::Idle
|
||||
);
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_migrates_schema_v3_workers_to_stopped_explicit_restore() {
|
||||
let root = fs_store_root("schema-v3-restore-intent");
|
||||
let options = crate::fs_store::FsRuntimeStoreOptions {
|
||||
root: root.clone(),
|
||||
runtime_id: "test-runtime".to_string(),
|
||||
display_name: None,
|
||||
};
|
||||
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||
options.clone(),
|
||||
Arc::new(TestExecutionBackend::default()),
|
||||
)
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let worker = runtime
|
||||
.create_worker(task_request("schema v3 worker"))
|
||||
.unwrap();
|
||||
drop(runtime);
|
||||
|
||||
let runtime_path = root.join("runtime.json");
|
||||
let worker_path = root
|
||||
.join("workers")
|
||||
.join(worker.worker_id.to_string())
|
||||
.join("worker.json");
|
||||
let mut runtime_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
||||
runtime_json["schema_version"] = serde_json::json!(3);
|
||||
std::fs::write(
|
||||
&runtime_path,
|
||||
serde_json::to_vec_pretty(&runtime_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
let mut worker_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
||||
worker_json["schema_version"] = serde_json::json!(3);
|
||||
worker_json.as_object_mut().unwrap().remove("status");
|
||||
worker_json.as_object_mut().unwrap().remove("execution");
|
||||
worker_json["run_generation"] = serde_json::json!(7);
|
||||
std::fs::write(
|
||||
&worker_path,
|
||||
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let backend = Arc::new(TestExecutionBackend::default());
|
||||
let migrated =
|
||||
Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap();
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
||||
assert_eq!(
|
||||
migrated.worker_detail(&worker.worker_ref).unwrap().status,
|
||||
WorkerStatus::Stopped
|
||||
);
|
||||
let migrated_json: serde_json::Value =
|
||||
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
||||
assert_eq!(migrated_json["schema_version"], serde_json::json!(4));
|
||||
assert_eq!(migrated_json["status"], serde_json::json!("stopped"));
|
||||
assert_eq!(
|
||||
migrated_json["execution"]["binding"]["run_generation"],
|
||||
serde_json::json!(7)
|
||||
);
|
||||
assert_eq!(
|
||||
migrated_json["execution"]["restore_intent"],
|
||||
serde_json::json!("explicit")
|
||||
);
|
||||
assert!(matches!(
|
||||
migrated.send_input(&worker.worker_ref, WorkerInput::notify("do not restore")),
|
||||
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
||||
));
|
||||
|
||||
migrated.restore_worker(&worker.worker_ref).unwrap();
|
||||
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
|
||||
assert_eq!(backend.run_generations.lock().unwrap().as_slice(), &[8]);
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
#[cfg(feature = "fs-store")]
|
||||
#[test]
|
||||
fn fs_store_stops_worker_and_reports_when_execution_restore_fails() {
|
||||
@@ -5702,7 +5844,10 @@ mod tests {
|
||||
WorkerInput::user("after failed restore"),
|
||||
)
|
||||
.unwrap_err();
|
||||
assert!(matches!(err, RuntimeError::WorkerExecutionRejected { .. }));
|
||||
assert!(matches!(
|
||||
err,
|
||||
RuntimeError::WorkerExecutionUnavailable { .. }
|
||||
));
|
||||
|
||||
let _ = std::fs::remove_dir_all(root);
|
||||
}
|
||||
|
||||
@@ -3679,6 +3679,60 @@ mod tests {
|
||||
assert_eq!(call_count.load(Ordering::SeqCst), 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stopped_runtime_worker_can_restore_and_accept_input() {
|
||||
let client = MockClient::new(simple_text_events());
|
||||
let runtime_base = tempfile::tempdir().unwrap();
|
||||
let cwd = tempfile::tempdir().unwrap();
|
||||
let store = tempfile::tempdir().unwrap();
|
||||
let factory = MockFactory {
|
||||
client,
|
||||
runtime_base: runtime_base.path().to_path_buf(),
|
||||
cwd: cwd.path().to_path_buf(),
|
||||
store_dir: store.path().join("sessions"),
|
||||
worker_metadata_dir: store.path().join("workers"),
|
||||
observed_cwds: Arc::new(Mutex::new(Vec::new())),
|
||||
observed_workspace_clients: Arc::new(Mutex::new(Vec::new())),
|
||||
};
|
||||
let backend = Arc::new(WorkerRuntimeExecutionBackend::new(factory).unwrap());
|
||||
let runtime =
|
||||
EmbeddedRuntime::with_execution_backend(RuntimeOptions::default(), backend.clone())
|
||||
.unwrap();
|
||||
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||
let detail = runtime
|
||||
.create_worker(create_request("restore-after-stop"))
|
||||
.unwrap();
|
||||
|
||||
runtime.stop_worker(&detail.worker_ref, None).unwrap();
|
||||
assert_eq!(
|
||||
runtime.worker_detail(&detail.worker_ref).unwrap().status,
|
||||
crate::catalog::WorkerStatus::Stopped
|
||||
);
|
||||
assert!(
|
||||
!backend
|
||||
.workers
|
||||
.lock()
|
||||
.unwrap()
|
||||
.contains_key(&detail.worker_ref)
|
||||
);
|
||||
|
||||
runtime.restore_worker(&detail.worker_ref).unwrap();
|
||||
assert_eq!(
|
||||
runtime.worker_detail(&detail.worker_ref).unwrap().status,
|
||||
crate::catalog::WorkerStatus::Idle
|
||||
);
|
||||
assert!(
|
||||
backend
|
||||
.workers
|
||||
.lock()
|
||||
.unwrap()
|
||||
.contains_key(&detail.worker_ref)
|
||||
);
|
||||
runtime
|
||||
.send_input(&detail.worker_ref, WorkerInput::user("continue"))
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stopping_and_deleting_worker_preserves_bound_working_directory() {
|
||||
let client = MockClient::new(simple_text_events());
|
||||
|
||||
+24
-25
@@ -10,8 +10,8 @@ use agen::llm_client::client::LlmClient;
|
||||
use agen::llm_client::types::Role;
|
||||
use agen::state::Mutable;
|
||||
use agen::{
|
||||
Engine, EngineError, EngineResult, EngineRunExit, History, HistoryEntry, Item, StopReason,
|
||||
ToolExecutionPolicy, ToolOutputLimits, UsageRecord,
|
||||
Engine, EngineError, EngineResult, EngineRunExit, History, HistoryEntry, Item,
|
||||
RunInterruptionReason, ToolExecutionPolicy, ToolOutputLimits, UsageRecord,
|
||||
};
|
||||
use arc_swap::ArcSwap;
|
||||
use session_store::{
|
||||
@@ -2605,7 +2605,7 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
||||
) -> bool {
|
||||
if !matches!(
|
||||
result,
|
||||
EngineRunExit::Paused | EngineRunExit::Interrupted(StopReason::Cancelled)
|
||||
EngineRunExit::Paused | EngineRunExit::Interrupted(RunInterruptionReason::Cancelled)
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
@@ -3403,15 +3403,15 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
||||
self.last_run_interrupted = true;
|
||||
Ok(WorkerRunResult::Paused)
|
||||
}
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached) => {
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached) => {
|
||||
self.last_run_interrupted = false;
|
||||
Ok(WorkerRunResult::LimitReached)
|
||||
}
|
||||
EngineRunExit::Interrupted(reason) => {
|
||||
self.last_run_interrupted = true;
|
||||
Ok(WorkerRunResult::Interrupted {
|
||||
code: stop_reason_error_code(&reason),
|
||||
message: stop_reason_message(&reason),
|
||||
code: run_interruption_reason_error_code(&reason),
|
||||
message: run_interruption_reason_message(&reason),
|
||||
})
|
||||
}
|
||||
EngineRunExit::Yielded => unreachable!("yielded handled above"),
|
||||
@@ -3731,9 +3731,9 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
||||
result,
|
||||
EngineRunExit::Paused
|
||||
| EngineRunExit::Yielded
|
||||
| EngineRunExit::Interrupted(StopReason::Cancelled)
|
||||
| EngineRunExit::Interrupted(StopReason::ContextWindowExceeded)
|
||||
| EngineRunExit::Interrupted(StopReason::Unexpected(_))
|
||||
| EngineRunExit::Interrupted(RunInterruptionReason::Cancelled)
|
||||
| EngineRunExit::Interrupted(RunInterruptionReason::ContextWindowExceeded)
|
||||
| EngineRunExit::Interrupted(RunInterruptionReason::Unexpected(_))
|
||||
);
|
||||
let active_run_turn_count = self.engine.as_ref().unwrap().active_run_turn_count();
|
||||
match result {
|
||||
@@ -3751,7 +3751,7 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
||||
active_run_turn_count,
|
||||
})?;
|
||||
}
|
||||
EngineRunExit::Interrupted(StopReason::LimitReached) => {
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::LimitReached) => {
|
||||
self.commit_entry(LogEntry::RunCompleted {
|
||||
ts: segment_log::now_millis(),
|
||||
interrupted: false,
|
||||
@@ -3763,7 +3763,7 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
|
||||
self.commit_entry(LogEntry::RunErrored {
|
||||
ts: segment_log::now_millis(),
|
||||
interrupted,
|
||||
message: stop_reason_message(reason),
|
||||
message: run_interruption_reason_message(reason),
|
||||
})?;
|
||||
}
|
||||
}
|
||||
@@ -6122,15 +6122,14 @@ fn restore_manifest_from_worker_metadata_snapshot(
|
||||
}
|
||||
}
|
||||
|
||||
fn stop_reason_error_code(reason: &StopReason) -> ErrorCode {
|
||||
fn run_interruption_reason_error_code(reason: &RunInterruptionReason) -> ErrorCode {
|
||||
match reason {
|
||||
StopReason::ContextWindowExceeded | StopReason::Unexpected(EngineError::Client(_)) => {
|
||||
ErrorCode::ProviderError
|
||||
}
|
||||
StopReason::Unexpected(EngineError::Tool(_)) => ErrorCode::ToolError,
|
||||
StopReason::LimitReached
|
||||
| StopReason::Cancelled
|
||||
| StopReason::Unexpected(
|
||||
RunInterruptionReason::ContextWindowExceeded
|
||||
| RunInterruptionReason::Unexpected(EngineError::Client(_)) => ErrorCode::ProviderError,
|
||||
RunInterruptionReason::Unexpected(EngineError::Tool(_)) => ErrorCode::ToolError,
|
||||
RunInterruptionReason::LimitReached
|
||||
| RunInterruptionReason::Cancelled
|
||||
| RunInterruptionReason::Unexpected(
|
||||
EngineError::Aborted(_)
|
||||
| EngineError::Cancelled
|
||||
| EngineError::PauseRequested
|
||||
@@ -6141,12 +6140,12 @@ fn stop_reason_error_code(reason: &StopReason) -> ErrorCode {
|
||||
}
|
||||
}
|
||||
|
||||
fn stop_reason_message(reason: &StopReason) -> String {
|
||||
fn run_interruption_reason_message(reason: &RunInterruptionReason) -> String {
|
||||
match reason {
|
||||
StopReason::LimitReached => "engine turn limit reached".to_string(),
|
||||
StopReason::ContextWindowExceeded => "model context window reached".to_string(),
|
||||
StopReason::Cancelled => "engine run cancelled".to_string(),
|
||||
StopReason::Unexpected(error) => format!("unexpected engine failure: {error}"),
|
||||
RunInterruptionReason::LimitReached => "engine turn limit reached".to_string(),
|
||||
RunInterruptionReason::ContextWindowExceeded => "model context window reached".to_string(),
|
||||
RunInterruptionReason::Cancelled => "engine run cancelled".to_string(),
|
||||
RunInterruptionReason::Unexpected(error) => format!("unexpected engine failure: {error}"),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8695,7 +8694,7 @@ mod build_summary_prompt_tests {
|
||||
]);
|
||||
let _ = worker
|
||||
.handle_worker_result(
|
||||
EngineRunExit::Interrupted(StopReason::Cancelled),
|
||||
EngineRunExit::Interrupted(RunInterruptionReason::Cancelled),
|
||||
worker.history().len(),
|
||||
)
|
||||
.await
|
||||
|
||||
@@ -7,9 +7,10 @@ publish = false
|
||||
|
||||
[features]
|
||||
default = []
|
||||
typescript = ["dep:ts-rs"]
|
||||
typescript = ["dep:ts-rs", "protocol/typescript"]
|
||||
|
||||
[dependencies]
|
||||
protocol.workspace = true
|
||||
serde = { workspace = true, features = ["derive"] }
|
||||
ts-rs = { version = "12.0.1", optional = true }
|
||||
|
||||
@@ -24,6 +25,10 @@ serde_json.workspace = true
|
||||
name = "generate_workdir_api_types"
|
||||
required-features = ["typescript"]
|
||||
|
||||
[[example]]
|
||||
name = "generate_worker_launch_api_types"
|
||||
required-features = ["typescript"]
|
||||
|
||||
[[example]]
|
||||
name = "generate_companion_api_types"
|
||||
required-features = ["typescript"]
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
fn main() {
|
||||
print!("{}", workspace_api::worker_launch_api_typescript());
|
||||
}
|
||||
@@ -579,6 +579,8 @@ pub struct WorkingDirectoryOccupancy {
|
||||
/// retains the Backend-generated Repository id and is never a Workspace public
|
||||
/// projection.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[cfg_attr(feature = "typescript", ts(optional_fields = nullable))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct RuntimeWorkingDirectoryCleanupTarget {
|
||||
pub kind: String,
|
||||
@@ -590,6 +592,8 @@ pub struct RuntimeWorkingDirectoryCleanupTarget {
|
||||
/// surfaces must project this through [`WorkingDirectorySummary`] so the UUID is
|
||||
/// replaced with `repository_key`.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[cfg_attr(feature = "typescript", ts(optional_fields = nullable))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct RuntimeWorkingDirectorySummary {
|
||||
pub working_directory_id: String,
|
||||
@@ -607,6 +611,7 @@ pub struct RuntimeWorkingDirectorySummary {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub current_tree: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
#[cfg_attr(feature = "typescript", ts(type = "number | null"))]
|
||||
pub observed_at_epoch_seconds: Option<u64>,
|
||||
pub materializer_kind: WorkingDirectoryMaterializerKind,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
@@ -908,6 +913,7 @@ pub struct RuntimeConnectionTestResponse {
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
pub struct WorkerWorkspaceSummary {
|
||||
pub visibility: String,
|
||||
pub identity: String,
|
||||
@@ -916,12 +922,14 @@ pub struct WorkerWorkspaceSummary {
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
pub struct WorkerImplementationSummary {
|
||||
pub kind: String,
|
||||
pub display_hint: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
pub struct WorkerCapabilitySummary {
|
||||
pub can_stop: bool,
|
||||
pub can_spawn_followup: bool,
|
||||
@@ -1095,6 +1103,7 @@ pub struct WorkspaceWorkerDiscoveryPage {
|
||||
/// do not carry one. The Workspace Server must resolve it from Workspace
|
||||
/// authority before constructing this response.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
pub struct WorkerSummary {
|
||||
pub runtime_id: String,
|
||||
pub worker_id: String,
|
||||
@@ -1122,6 +1131,142 @@ pub struct WorkerSummary {
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
/// Runtime-owned Worker summary embedded in Worker launch responses.
|
||||
///
|
||||
/// This preserves the existing launch wire shape. Workspace-owned Worker list
|
||||
/// and detail responses use [`WorkerSummary`] instead.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct WorkerLaunchWorkerSummary {
|
||||
pub runtime_id: String,
|
||||
pub worker_id: String,
|
||||
pub host_id: String,
|
||||
pub display_name: String,
|
||||
pub label: String,
|
||||
pub profile: Option<String>,
|
||||
pub singleton_key: Option<String>,
|
||||
pub tags: Vec<String>,
|
||||
pub workspace: WorkerWorkspaceSummary,
|
||||
pub state: String,
|
||||
pub last_seen_at: Option<String>,
|
||||
pub pinned: bool,
|
||||
pub retention_state: String,
|
||||
pub implementation: WorkerImplementationSummary,
|
||||
pub capabilities: WorkerCapabilitySummary,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
#[cfg_attr(feature = "typescript", ts(optional = nullable))]
|
||||
pub working_directory: Option<RuntimeWorkingDirectorySummary>,
|
||||
#[serde(default)]
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct WorkerLaunchOptionsResponse {
|
||||
pub workspace_id: String,
|
||||
pub runtimes: Vec<WorkerLaunchRuntimeOption>,
|
||||
pub default_profile: Option<String>,
|
||||
pub profiles: Vec<WorkerLaunchProfileCandidate>,
|
||||
pub repositories: Vec<WorkingDirectoryRepositoryOption>,
|
||||
pub working_directories: Vec<WorkingDirectorySummary>,
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct WorkerLaunchRuntimeOption {
|
||||
pub runtime_id: String,
|
||||
pub display_name: String,
|
||||
pub built_in: bool,
|
||||
pub worker_creation_available: bool,
|
||||
pub working_directory_required: bool,
|
||||
pub status: String,
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct WorkerLaunchProfileCandidate {
|
||||
pub id: String,
|
||||
pub label: String,
|
||||
pub description: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct WorkingDirectoryRepositoryOption {
|
||||
pub repository_key: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
#[cfg_attr(feature = "typescript", ts(optional = nullable))]
|
||||
pub default_selector: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct BrowserWorkerWorkingDirectorySelection {
|
||||
pub working_directory_id: String,
|
||||
#[serde(default)]
|
||||
pub relative_cwd: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct CreateWorkspaceWorkerTicketAssignmentRequest {
|
||||
pub ticket_id: String,
|
||||
pub operation_id: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct CreateWorkspaceWorkerRequest {
|
||||
pub runtime_id: String,
|
||||
pub display_name: String,
|
||||
#[serde(default)]
|
||||
pub profile: Option<String>,
|
||||
#[serde(default)]
|
||||
pub ticket_assignment: Option<CreateWorkspaceWorkerTicketAssignmentRequest>,
|
||||
#[serde(default)]
|
||||
pub initial_submit: Vec<protocol::Segment>,
|
||||
#[serde(default)]
|
||||
pub working_directory: Option<BrowserWorkerWorkingDirectorySelection>,
|
||||
/// Backend idempotency key used only for authenticated Worker-owned spawn/control.
|
||||
#[serde(default)]
|
||||
pub control_operation_id: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct BrowserCreateWorkerResponse {
|
||||
pub workspace_id: String,
|
||||
pub runtime_id: String,
|
||||
pub worker_id: String,
|
||||
pub console_href: String,
|
||||
pub worker: WorkerLaunchWorkerSummary,
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[cfg_attr(feature = "typescript", derive(ts_rs::TS))]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct BrowserWorkspaceOrchestratorResponse {
|
||||
pub workspace_id: String,
|
||||
pub online: bool,
|
||||
pub disposition: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
#[cfg_attr(feature = "typescript", ts(optional = nullable))]
|
||||
pub worker: Option<WorkerLaunchWorkerSummary>,
|
||||
pub diagnostics: Vec<Diagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum WorkerOperationState {
|
||||
@@ -1390,6 +1535,74 @@ pub fn workdir_api_typescript() -> String {
|
||||
)
|
||||
}
|
||||
|
||||
#[cfg(feature = "typescript")]
|
||||
pub fn worker_launch_api_typescript() -> String {
|
||||
use ts_rs::TS;
|
||||
|
||||
let config = ts_rs::Config::default();
|
||||
let declarations = [
|
||||
DiagnosticSeverity::decl(&config),
|
||||
Diagnostic::decl(&config),
|
||||
WorkingDirectoryMaterializerKind::decl(&config),
|
||||
WorkingDirectoryStatusKind::decl(&config),
|
||||
WorkingDirectoryCleanupTarget::decl(&config),
|
||||
RuntimeWorkingDirectoryCleanupTarget::decl(&config),
|
||||
RuntimeWorkingDirectorySummary::decl(&config),
|
||||
WorkingDirectoryOccupancy::decl(&config),
|
||||
WorkingDirectorySummary::decl(&config),
|
||||
WorkerWorkspaceSummary::decl(&config),
|
||||
WorkerImplementationSummary::decl(&config),
|
||||
WorkerCapabilitySummary::decl(&config),
|
||||
WorkerLaunchWorkerSummary::decl(&config),
|
||||
WorkerLaunchRuntimeOption::decl(&config),
|
||||
WorkerLaunchProfileCandidate::decl(&config),
|
||||
WorkingDirectoryRepositoryOption::decl(&config),
|
||||
WorkerLaunchOptionsResponse::decl(&config),
|
||||
BrowserWorkerWorkingDirectorySelection::decl(&config),
|
||||
CreateWorkspaceWorkerTicketAssignmentRequest::decl(&config),
|
||||
CreateWorkspaceWorkerRequest::decl(&config),
|
||||
BrowserCreateWorkerResponse::decl(&config),
|
||||
BrowserWorkspaceOrchestratorResponse::decl(&config),
|
||||
];
|
||||
format!(
|
||||
"// Generated from workspace-api. Do not edit by hand.\n// Regenerate: cargo run -q -p workspace-api --features typescript --example generate_worker_launch_api_types > web/workspace/src/lib/generated/worker-launch-api.ts\n\nimport type {{ Segment }} from \"./protocol\";\n\n{}\n",
|
||||
declarations
|
||||
.into_iter()
|
||||
.map(|declaration| format!("export {declaration}"))
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n\n")
|
||||
)
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "typescript"))]
|
||||
mod worker_launch_typescript_tests {
|
||||
#[test]
|
||||
fn generated_worker_launch_api_contract_is_current() {
|
||||
let expected = super::worker_launch_api_typescript();
|
||||
let path = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
|
||||
.join("../../web/workspace/src/lib/generated/worker-launch-api.ts");
|
||||
let actual = std::fs::read_to_string(&path)
|
||||
.unwrap_or_else(|error| panic!("failed to read {}: {error}", path.display()));
|
||||
assert_eq!(
|
||||
normalize(&actual),
|
||||
normalize(&expected),
|
||||
"regenerate Worker launch API TypeScript types with `cargo run -q -p workspace-api --features typescript --example generate_worker_launch_api_types > web/workspace/src/lib/generated/worker-launch-api.ts` and format the generated file",
|
||||
);
|
||||
}
|
||||
|
||||
fn normalize(value: &str) -> String {
|
||||
value
|
||||
.chars()
|
||||
.filter_map(|character| match character {
|
||||
character if character.is_whitespace() => None,
|
||||
',' => Some(';'),
|
||||
character => Some(character),
|
||||
})
|
||||
.collect::<String>()
|
||||
.replace("=|", "=")
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "typescript"))]
|
||||
mod workdir_typescript_tests {
|
||||
#[test]
|
||||
@@ -1423,6 +1636,113 @@ mod workdir_typescript_tests {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn worker_launch_summary() -> WorkerLaunchWorkerSummary {
|
||||
WorkerLaunchWorkerSummary {
|
||||
runtime_id: "runtime-a".to_string(),
|
||||
worker_id: "worker-a".to_string(),
|
||||
host_id: "host-a".to_string(),
|
||||
display_name: "Worker A".to_string(),
|
||||
label: "worker-a".to_string(),
|
||||
profile: None,
|
||||
singleton_key: None,
|
||||
tags: Vec::new(),
|
||||
workspace: WorkerWorkspaceSummary {
|
||||
visibility: "workspace".to_string(),
|
||||
identity: "workspace-a".to_string(),
|
||||
workspace_id: Some("workspace-a".to_string()),
|
||||
},
|
||||
state: "idle".to_string(),
|
||||
last_seen_at: None,
|
||||
pinned: false,
|
||||
retention_state: "active".to_string(),
|
||||
implementation: WorkerImplementationSummary {
|
||||
kind: "runtime".to_string(),
|
||||
display_hint: "Runtime Worker".to_string(),
|
||||
},
|
||||
capabilities: WorkerCapabilitySummary {
|
||||
can_stop: true,
|
||||
can_spawn_followup: false,
|
||||
},
|
||||
working_directory: None,
|
||||
diagnostics: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn worker_launch_optional_omission_and_request_shape_are_stable() {
|
||||
assert_eq!(
|
||||
serde_json::to_value(WorkingDirectoryRepositoryOption {
|
||||
repository_key: "main".to_string(),
|
||||
default_selector: None,
|
||||
})
|
||||
.unwrap(),
|
||||
serde_json::json!({ "repository_key": "main" })
|
||||
);
|
||||
|
||||
let orchestrator = serde_json::to_value(BrowserWorkspaceOrchestratorResponse {
|
||||
workspace_id: "workspace-a".to_string(),
|
||||
online: false,
|
||||
disposition: "unavailable".to_string(),
|
||||
worker: None,
|
||||
diagnostics: Vec::new(),
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
orchestrator,
|
||||
serde_json::json!({
|
||||
"workspace_id": "workspace-a",
|
||||
"online": false,
|
||||
"disposition": "unavailable",
|
||||
"diagnostics": [],
|
||||
})
|
||||
);
|
||||
|
||||
let worker = serde_json::to_value(worker_launch_summary()).unwrap();
|
||||
assert!(
|
||||
!worker
|
||||
.as_object()
|
||||
.unwrap()
|
||||
.contains_key("working_directory")
|
||||
);
|
||||
assert_eq!(worker["profile"], serde_json::Value::Null);
|
||||
assert_eq!(worker["singleton_key"], serde_json::Value::Null);
|
||||
assert_eq!(worker["last_seen_at"], serde_json::Value::Null);
|
||||
|
||||
let request = serde_json::to_value(CreateWorkspaceWorkerRequest {
|
||||
runtime_id: "runtime-a".to_string(),
|
||||
display_name: "Worker A".to_string(),
|
||||
profile: None,
|
||||
ticket_assignment: None,
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
request,
|
||||
serde_json::json!({
|
||||
"runtime_id": "runtime-a",
|
||||
"display_name": "Worker A",
|
||||
"profile": null,
|
||||
"ticket_assignment": null,
|
||||
"initial_submit": [],
|
||||
"working_directory": null,
|
||||
"control_operation_id": null,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn worker_launch_request_rejects_unknown_fields() {
|
||||
let error = serde_json::from_value::<CreateWorkspaceWorkerRequest>(serde_json::json!({
|
||||
"runtime_id": "runtime-a",
|
||||
"display_name": "Worker A",
|
||||
"unexpected": true,
|
||||
}))
|
||||
.unwrap_err();
|
||||
assert!(error.to_string().contains("unknown field"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repository_key_validation_is_canonical_and_bounded() {
|
||||
let max = "a".repeat(64);
|
||||
|
||||
@@ -3842,7 +3842,6 @@ fn embedded_worker_status_label(status: EmbeddedWorkerStatus) -> &'static str {
|
||||
EmbeddedWorkerStatus::Running => "running",
|
||||
EmbeddedWorkerStatus::Paused => "paused",
|
||||
EmbeddedWorkerStatus::Stopped => "stopped",
|
||||
EmbeddedWorkerStatus::Cancelled => "cancelled",
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5827,7 +5826,7 @@ mod tests {
|
||||
json!({
|
||||
"workers": [
|
||||
worker_json_with_status("remote:primary", &worker_ids[0], "stopped"),
|
||||
worker_json_with_status("remote:primary", &worker_ids[1], "cancelled"),
|
||||
worker_json_with_status("remote:primary", &worker_ids[1], "running"),
|
||||
worker_json_with_status("remote:primary", &worker_ids[2], "paused"),
|
||||
worker_json_with_status("remote:primary", &worker_ids[3], "idle")
|
||||
]
|
||||
@@ -5866,11 +5865,11 @@ mod tests {
|
||||
let workers = registry.list_workers(10);
|
||||
assert_eq!(workers.items.len(), 4);
|
||||
assert!(!workers.items[0].capabilities.can_stop);
|
||||
assert!(!workers.items[1].capabilities.can_stop);
|
||||
assert!(workers.items[1].capabilities.can_stop);
|
||||
assert!(workers.items[2].capabilities.can_stop);
|
||||
assert!(workers.items[3].capabilities.can_stop);
|
||||
assert_eq!(workers.items[0].state, "stopped");
|
||||
assert_eq!(workers.items[1].state, "cancelled");
|
||||
assert_eq!(workers.items[1].state, "running");
|
||||
assert_eq!(workers.items[2].state, "paused");
|
||||
assert_eq!(workers.items[3].state, "idle");
|
||||
|
||||
|
||||
@@ -58,25 +58,28 @@ use worker::feature::builtin::{WorkerObservationSubject, WorkerObservationSubjec
|
||||
use worker_runtime::resource::{BackendResourceError, BackendResourceFetchRequest};
|
||||
use worker_runtime::worker_backend::{ProfileRuntimeWorkerFactory, WorkerRuntimeExecutionBackend};
|
||||
use workspace_api::{
|
||||
CreateRemoteRuntimeRequest, CreateRepositorySshCredentialRequest,
|
||||
CreateWorkspaceRepositoryRequest, CreateWorkspaceRepositoryResponse,
|
||||
DeleteRepositorySshCredentialRequest, DeleteRepositorySshHostTrustRequest,
|
||||
ObjectiveCreateRequest, ObjectiveEditRequest, ObjectiveLinkTicketRequest,
|
||||
ObjectiveStateRequest, ProfileSettingsResponse, PutRepositorySshHostTrustRequest,
|
||||
RepositoryAccessProjection, RepositoryDetailResponse, RepositoryListResponse,
|
||||
RepositoryLogResponse, RepositorySshCredential, RepositorySshHostTrust,
|
||||
BrowserCreateWorkerResponse, BrowserWorkspaceOrchestratorResponse, CreateRemoteRuntimeRequest,
|
||||
CreateRepositorySshCredentialRequest, CreateWorkspaceRepositoryRequest,
|
||||
CreateWorkspaceRepositoryResponse, CreateWorkspaceWorkerRequest,
|
||||
CreateWorkspaceWorkerTicketAssignmentRequest, DeleteRepositorySshCredentialRequest,
|
||||
DeleteRepositorySshHostTrustRequest, ObjectiveCreateRequest, ObjectiveEditRequest,
|
||||
ObjectiveLinkTicketRequest, ObjectiveStateRequest, ProfileSettingsResponse,
|
||||
PutRepositorySshHostTrustRequest, RepositoryAccessProjection, RepositoryDetailResponse,
|
||||
RepositoryListResponse, RepositoryLogResponse, RepositorySshCredential, RepositorySshHostTrust,
|
||||
RotateRepositorySshCredentialRequest, RuntimeConnectionTestResponse, RuntimeManagementSummary,
|
||||
TICKET_ORCHESTRATION_PLANS_QUERY_PATH, TICKET_RELATIONS_QUERY_PATH,
|
||||
UpdateWorkspaceMetadataRequest,
|
||||
UpdateWorkspaceMetadataRequest, WorkerLaunchOptionsResponse, WorkerLaunchProfileCandidate,
|
||||
WorkerLaunchRuntimeOption, WorkerLaunchWorkerSummary,
|
||||
WorkingDirectoryCreateRequest as BrowserWorkingDirectoryCreateRequest,
|
||||
WorkingDirectoryCreateResponse as BrowserWorkingDirectoryCreateResponse,
|
||||
WorkingDirectoryDetailResponse as BrowserWorkingDirectoryDetailResponse,
|
||||
WorkingDirectoryListResponse as BrowserWorkingDirectoryListResponse,
|
||||
WorkingDirectoryRemovalDisposition, WorkingDirectoryRemovalRequest,
|
||||
WorkingDirectoryRemovalResponse, WorkspaceCatalogListResponse, WorkspaceCreateResponse,
|
||||
WorkspaceExtensionPointState, WorkspaceExtensionPoints, WorkspaceMetadataMutationResponse,
|
||||
WorkspaceMetadataSettingsResponse, WorkspacePermissionSummary, WorkspaceRepositoryRecord,
|
||||
WorkspaceResponse, WorkspaceRuntimeResource, WorkspaceSummary, WorkspaceWorkerDiscoveryItem,
|
||||
WorkingDirectoryRemovalResponse, WorkingDirectoryRepositoryOption,
|
||||
WorkspaceCatalogListResponse, WorkspaceCreateResponse, WorkspaceExtensionPointState,
|
||||
WorkspaceExtensionPoints, WorkspaceMetadataMutationResponse, WorkspaceMetadataSettingsResponse,
|
||||
WorkspacePermissionSummary, WorkspaceRepositoryRecord, WorkspaceResponse,
|
||||
WorkspaceRuntimeResource, WorkspaceSummary, WorkspaceWorkerDiscoveryItem,
|
||||
WorkspaceWorkerDiscoveryPage, WorkspaceWorkerSubject,
|
||||
};
|
||||
|
||||
@@ -3097,98 +3100,6 @@ pub struct WorkerRetentionResponse {
|
||||
pub retention_state: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
pub struct WorkerLaunchOptionsResponse {
|
||||
pub workspace_id: String,
|
||||
pub runtimes: Vec<WorkerLaunchRuntimeOption>,
|
||||
pub default_profile: Option<String>,
|
||||
pub profiles: Vec<WorkerLaunchProfileCandidate>,
|
||||
pub repositories: Vec<WorkingDirectoryRepositoryOption>,
|
||||
pub working_directories: Vec<WorkingDirectorySummary>,
|
||||
pub diagnostics: Vec<RuntimeDiagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
pub struct WorkerLaunchRuntimeOption {
|
||||
pub runtime_id: String,
|
||||
pub display_name: String,
|
||||
pub built_in: bool,
|
||||
pub worker_creation_available: bool,
|
||||
pub working_directory_required: bool,
|
||||
pub status: String,
|
||||
pub diagnostics: Vec<RuntimeDiagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct WorkerLaunchProfileCandidate {
|
||||
pub id: String,
|
||||
pub label: String,
|
||||
pub description: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct WorkingDirectoryRepositoryOption {
|
||||
pub repository_key: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub default_selector: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct BrowserWorkerWorkingDirectorySelection {
|
||||
pub working_directory_id: String,
|
||||
#[serde(default)]
|
||||
pub relative_cwd: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
pub struct BrowserWorkspaceOrchestratorResponse {
|
||||
pub workspace_id: String,
|
||||
pub online: bool,
|
||||
pub disposition: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub worker: Option<WorkerSummary>,
|
||||
pub diagnostics: Vec<RuntimeDiagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct CreateWorkspaceWorkerTicketAssignmentRequest {
|
||||
pub ticket_id: String,
|
||||
pub operation_id: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct CreateWorkspaceWorkerRequest {
|
||||
pub runtime_id: String,
|
||||
pub display_name: String,
|
||||
#[serde(default)]
|
||||
pub profile: Option<String>,
|
||||
#[serde(default)]
|
||||
pub ticket_assignment: Option<CreateWorkspaceWorkerTicketAssignmentRequest>,
|
||||
#[serde(default)]
|
||||
pub initial_submit: Vec<Segment>,
|
||||
#[serde(default)]
|
||||
pub working_directory: Option<BrowserWorkerWorkingDirectorySelection>,
|
||||
/// Backend idempotency key used only for authenticated Worker-owned spawn/control.
|
||||
#[serde(default)]
|
||||
pub control_operation_id: Option<String>,
|
||||
/// Trusted resolution populated only by the authenticated worker-control handler.
|
||||
#[serde(skip, default)]
|
||||
pub resolved_control_operation: Option<WorkerControlOperation>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
pub struct BrowserCreateWorkerResponse {
|
||||
pub workspace_id: String,
|
||||
#[serde(flatten)]
|
||||
pub worker_ref: RuntimeWorkerRef,
|
||||
pub console_href: String,
|
||||
pub worker: WorkerSummary,
|
||||
pub diagnostics: Vec<RuntimeDiagnostic>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct LogQuery {
|
||||
limit: Option<usize>,
|
||||
@@ -8135,9 +8046,9 @@ fn start_memory_staging_consolidation(
|
||||
resolved_config_bundle,
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
},
|
||||
)?;
|
||||
if result.state != WorkerOperationState::Accepted {
|
||||
@@ -9088,18 +8999,20 @@ async fn spawn_known_worker(
|
||||
.map(|byte| format!("{byte:02x}"))
|
||||
.collect::<String>()
|
||||
);
|
||||
request.resolved_control_operation = Some(WorkerControlOperation {
|
||||
let resolved_control_operation = Some(WorkerControlOperation {
|
||||
operation_id: scoped_worker_control_operation_id(&controller, &operation_id),
|
||||
input_fingerprint,
|
||||
});
|
||||
let response = create_workspace_worker(State(api.clone()), headers, Json(request)).await?;
|
||||
let response =
|
||||
create_workspace_worker_inner(api.clone(), headers, request, resolved_control_operation)
|
||||
.await?;
|
||||
if let Err(error) = api
|
||||
.store
|
||||
.create_worker_control_grant(&WorkerControlGrantRecord {
|
||||
workspace_id: path.workspace_id.clone(),
|
||||
grant_id: new_id("wcg"),
|
||||
controller,
|
||||
subject: response.0.worker_ref.clone(),
|
||||
subject: RuntimeWorkerRef::new(&response.0.runtime_id, &response.0.worker_id),
|
||||
relation: relation.to_string(),
|
||||
origin: "worker_spawn".to_string(),
|
||||
permissions: vec![
|
||||
@@ -9432,9 +9345,9 @@ async fn scoped_start_workspace_orchestrator(
|
||||
resolved_config_bundle: None,
|
||||
resolved_worker_observation_enabled: true,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
},
|
||||
)?;
|
||||
if result.state != WorkerOperationState::Accepted || result.worker.is_none() {
|
||||
@@ -9462,6 +9375,42 @@ async fn scoped_start_workspace_orchestrator(
|
||||
Ok(Json(workspace_orchestrator_response(&api, "created")))
|
||||
}
|
||||
|
||||
fn worker_launch_worker_summary(worker: WorkerSummary) -> WorkerLaunchWorkerSummary {
|
||||
WorkerLaunchWorkerSummary {
|
||||
runtime_id: worker.worker.runtime_id,
|
||||
worker_id: worker.worker.worker_id,
|
||||
host_id: worker.host_id,
|
||||
display_name: worker.display_name,
|
||||
label: worker.label,
|
||||
profile: worker.profile,
|
||||
singleton_key: worker.singleton_key,
|
||||
tags: worker.tags,
|
||||
workspace: workspace_api::WorkerWorkspaceSummary {
|
||||
visibility: worker.workspace.visibility,
|
||||
identity: worker.workspace.identity,
|
||||
workspace_id: worker.workspace.workspace_id,
|
||||
},
|
||||
state: worker.state,
|
||||
last_seen_at: worker.last_seen_at,
|
||||
pinned: worker.pinned,
|
||||
retention_state: worker.retention_state,
|
||||
implementation: workspace_api::WorkerImplementationSummary {
|
||||
kind: worker.implementation.kind,
|
||||
display_hint: worker.implementation.display_hint,
|
||||
},
|
||||
capabilities: workspace_api::WorkerCapabilitySummary {
|
||||
can_stop: worker.capabilities.can_stop,
|
||||
can_spawn_followup: worker.capabilities.can_spawn_followup,
|
||||
},
|
||||
working_directory: worker.working_directory,
|
||||
diagnostics: worker
|
||||
.diagnostics
|
||||
.into_iter()
|
||||
.map(workspace_api::Diagnostic::from)
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
|
||||
fn workspace_orchestrator_response(
|
||||
api: &WorkspaceApi,
|
||||
disposition: &str,
|
||||
@@ -9472,13 +9421,20 @@ fn workspace_orchestrator_response(
|
||||
.is_some_and(workspace_orchestrator_is_online);
|
||||
let diagnostics = worker
|
||||
.as_ref()
|
||||
.map(|worker| worker.diagnostics.clone())
|
||||
.map(|worker| {
|
||||
worker
|
||||
.diagnostics
|
||||
.iter()
|
||||
.cloned()
|
||||
.map(workspace_api::Diagnostic::from)
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default();
|
||||
BrowserWorkspaceOrchestratorResponse {
|
||||
workspace_id: api.config.workspace_id.clone(),
|
||||
online,
|
||||
disposition: disposition.to_string(),
|
||||
worker,
|
||||
worker: worker.map(worker_launch_worker_summary),
|
||||
diagnostics,
|
||||
}
|
||||
}
|
||||
@@ -12835,6 +12791,15 @@ async fn create_workspace_worker(
|
||||
State(api): State<WorkspaceApi>,
|
||||
headers: HeaderMap,
|
||||
Json(request): Json<CreateWorkspaceWorkerRequest>,
|
||||
) -> ApiResult<Json<BrowserCreateWorkerResponse>> {
|
||||
create_workspace_worker_inner(api, headers, request, None).await
|
||||
}
|
||||
|
||||
async fn create_workspace_worker_inner(
|
||||
api: WorkspaceApi,
|
||||
headers: HeaderMap,
|
||||
request: CreateWorkspaceWorkerRequest,
|
||||
resolved_control_operation: Option<WorkerControlOperation>,
|
||||
) -> ApiResult<Json<BrowserCreateWorkerResponse>> {
|
||||
let CreateWorkspaceWorkerRequest {
|
||||
runtime_id,
|
||||
@@ -12844,7 +12809,6 @@ async fn create_workspace_worker(
|
||||
initial_submit,
|
||||
working_directory,
|
||||
control_operation_id: _,
|
||||
resolved_control_operation,
|
||||
} = request;
|
||||
let config_state = api
|
||||
.config_store
|
||||
@@ -13142,10 +13106,14 @@ fn browser_worker_response_from_summary(
|
||||
);
|
||||
Ok(BrowserCreateWorkerResponse {
|
||||
workspace_id,
|
||||
worker_ref: RuntimeWorkerRef::new(&runtime_id, &worker_id),
|
||||
runtime_id,
|
||||
worker_id,
|
||||
console_href,
|
||||
worker,
|
||||
diagnostics,
|
||||
worker: worker_launch_worker_summary(worker),
|
||||
diagnostics: diagnostics
|
||||
.into_iter()
|
||||
.map(workspace_api::Diagnostic::from)
|
||||
.collect(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -13576,22 +13544,13 @@ fn compensate_failed_worker_spawn(
|
||||
let cancellation = api
|
||||
.runtime
|
||||
.cancel_worker(&worker.worker, lifecycle_request.clone());
|
||||
let cancellation_accepted = cancellation
|
||||
let stop = api.runtime.stop_worker(&worker.worker, lifecycle_request);
|
||||
let stop_accepted = stop
|
||||
.as_ref()
|
||||
.is_ok_and(|result| result.state == WorkerOperationState::Accepted);
|
||||
let stop = (!cancellation_accepted)
|
||||
.then(|| api.runtime.stop_worker(&worker.worker, lifecycle_request));
|
||||
let stop_accepted = stop.as_ref().is_some_and(|result| {
|
||||
result
|
||||
.as_ref()
|
||||
.is_ok_and(|result| result.state == WorkerOperationState::Accepted)
|
||||
});
|
||||
let termination_detail = (!cancellation_accepted && !stop_accepted).then(|| {
|
||||
let termination_detail = (!stop_accepted).then(|| {
|
||||
let cancellation = lifecycle_failure_detail("cancel", &cancellation);
|
||||
let stop = stop
|
||||
.as_ref()
|
||||
.map(|result| lifecycle_failure_detail("stop", result))
|
||||
.unwrap_or_else(|| "stop was not attempted".to_string());
|
||||
let stop = lifecycle_failure_detail("stop", &stop);
|
||||
format!("{cancellation}; {stop}")
|
||||
});
|
||||
|
||||
@@ -15187,7 +15146,11 @@ fn worker_launch_options_response(api: &WorkspaceApi) -> ApiResult<WorkerLaunchO
|
||||
worker_creation_available: runtime.worker_creation_available,
|
||||
working_directory_required: !built_in,
|
||||
status: runtime.status,
|
||||
diagnostics: runtime.diagnostics,
|
||||
diagnostics: runtime
|
||||
.diagnostics
|
||||
.into_iter()
|
||||
.map(workspace_api::Diagnostic::from)
|
||||
.collect(),
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
@@ -17707,9 +17670,9 @@ mod tests {
|
||||
resolved_config_bundle: None,
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
assert!(
|
||||
api.validate_worker_spawn_repository_scope(&workdir_flow_launch)
|
||||
@@ -17953,9 +17916,9 @@ mod tests {
|
||||
resolved_config_bundle: None,
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
|
||||
assert!(
|
||||
@@ -18002,7 +17965,6 @@ mod tests {
|
||||
}],
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18018,7 +17980,8 @@ mod tests {
|
||||
.get_current_ticket_coder_assignment(&api.config.workspace_id, &ticket.id)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(current.worker, response.worker_ref);
|
||||
let response_ref = RuntimeWorkerRef::new(&response.runtime_id, &response.worker_id);
|
||||
assert_eq!(current.worker, response_ref);
|
||||
let operation = api
|
||||
.store
|
||||
.get_ticket_assignment_operation(
|
||||
@@ -18028,7 +17991,7 @@ mod tests {
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(operation.assignment_id, Some(current.assignment_id));
|
||||
assert_eq!(operation.worker, Some(response.worker_ref));
|
||||
assert_eq!(operation.worker, Some(response_ref));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -18058,7 +18021,6 @@ mod tests {
|
||||
}],
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
@@ -18092,7 +18054,6 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18100,11 +18061,11 @@ mod tests {
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert(
|
||||
"x-yoi-runtime-id",
|
||||
axum::http::HeaderValue::from_str(&created.worker_ref.runtime_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&created.runtime_id).unwrap(),
|
||||
);
|
||||
headers.insert(
|
||||
"x-yoi-worker-id",
|
||||
axum::http::HeaderValue::from_str(&created.worker_ref.worker_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&created.worker_id).unwrap(),
|
||||
);
|
||||
let error =
|
||||
authenticate_worker_mutation_source(&api, "other-workspace", &headers).unwrap_err();
|
||||
@@ -18128,14 +18089,14 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let generic_ref = RuntimeWorkerRef::new(&generic.runtime_id, &generic.worker_id);
|
||||
assert!(matches!(
|
||||
require_online_workspace_orchestrator_source(&api, &generic.worker_ref),
|
||||
require_online_workspace_orchestrator_source(&api, &generic_ref),
|
||||
Err(Error::TicketAssignmentConflict(_))
|
||||
));
|
||||
|
||||
@@ -18147,10 +18108,11 @@ mod tests {
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let orchestrator = started.worker.unwrap().worker;
|
||||
let orchestrator = started.worker.unwrap();
|
||||
let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id);
|
||||
require_online_workspace_orchestrator_source(&api, &orchestrator).unwrap();
|
||||
assert!(matches!(
|
||||
require_online_workspace_orchestrator_source(&api, &generic.worker_ref),
|
||||
require_online_workspace_orchestrator_source(&api, &generic_ref),
|
||||
Err(Error::TicketAssignmentConflict(_))
|
||||
));
|
||||
|
||||
@@ -18219,8 +18181,8 @@ mod tests {
|
||||
assert!(started.online);
|
||||
let worker = started
|
||||
.worker
|
||||
.expect("production Workspace Orchestrator Worker")
|
||||
.worker;
|
||||
.expect("production Workspace Orchestrator Worker");
|
||||
let worker = RuntimeWorkerRef::new(&worker.runtime_id, &worker.worker_id);
|
||||
|
||||
let stopped = api
|
||||
.runtime
|
||||
@@ -18320,12 +18282,12 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let controller = controller_worker.worker_ref;
|
||||
let controller =
|
||||
RuntimeWorkerRef::new(&controller_worker.runtime_id, &controller_worker.worker_id);
|
||||
assert_ne!(
|
||||
scoped_worker_control_operation_id(&controller, "same-operation"),
|
||||
scoped_worker_control_operation_id(
|
||||
@@ -18351,7 +18313,6 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: Some("control-spawn-retry".to_string()),
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
|
||||
let Json(first) = spawn_known_worker(
|
||||
@@ -18375,7 +18336,8 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(retried.worker_ref, first.worker_ref);
|
||||
assert_eq!(retried.runtime_id, first.runtime_id);
|
||||
assert_eq!(retried.worker_id, first.worker_id);
|
||||
let mut conflicting_request = request();
|
||||
conflicting_request.display_name = "Different controlled child".to_string();
|
||||
let conflict = spawn_known_worker(
|
||||
@@ -18405,7 +18367,10 @@ mod tests {
|
||||
.list_active_worker_control_grants(&workspace_id, &controller, 10)
|
||||
.unwrap();
|
||||
assert_eq!(grants.len(), 1);
|
||||
assert_eq!(grants[0].subject, first.worker_ref);
|
||||
assert_eq!(
|
||||
grants[0].subject,
|
||||
RuntimeWorkerRef::new(&first.runtime_id, &first.worker_id)
|
||||
);
|
||||
assert_eq!(grants[0].operation_id, "control-spawn-retry");
|
||||
}
|
||||
|
||||
@@ -18427,7 +18392,6 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18445,7 +18409,6 @@ mod tests {
|
||||
initial_submit: Vec::new(),
|
||||
working_directory: None,
|
||||
control_operation_id: None,
|
||||
resolved_control_operation: None,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18466,14 +18429,16 @@ mod tests {
|
||||
dedicated.singleton_key.as_deref(),
|
||||
Some(crate::hosts::WORKSPACE_ORCHESTRATOR_SINGLETON_KEY)
|
||||
);
|
||||
assert_ne!(dedicated.worker.worker_id, generic.worker_ref.worker_id);
|
||||
let dedicated_ref = RuntimeWorkerRef::new(&dedicated.runtime_id, &dedicated.worker_id);
|
||||
let generic_ref = RuntimeWorkerRef::new(&generic.runtime_id, &generic.worker_id);
|
||||
assert_ne!(dedicated.worker_id, generic.worker_id);
|
||||
|
||||
api.store
|
||||
.create_worker_control_grant(&WorkerControlGrantRecord {
|
||||
workspace_id: workspace_id.clone(),
|
||||
grant_id: "orchestrator-controls-generic".to_string(),
|
||||
controller: dedicated.worker.clone(),
|
||||
subject: generic.worker_ref.clone(),
|
||||
controller: dedicated_ref.clone(),
|
||||
subject: generic_ref.clone(),
|
||||
relation: "spawned".to_string(),
|
||||
origin: "test".to_string(),
|
||||
permissions: vec!["observe".to_string()],
|
||||
@@ -18486,11 +18451,11 @@ mod tests {
|
||||
let mut observation_headers = HeaderMap::new();
|
||||
observation_headers.insert(
|
||||
"x-yoi-runtime-id",
|
||||
axum::http::HeaderValue::from_str(&dedicated.worker.runtime_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&dedicated.runtime_id).unwrap(),
|
||||
);
|
||||
observation_headers.insert(
|
||||
"x-yoi-worker-id",
|
||||
axum::http::HeaderValue::from_str(&dedicated.worker.worker_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&dedicated.worker_id).unwrap(),
|
||||
);
|
||||
let Json(known) = list_known_workers(
|
||||
State(api.clone()),
|
||||
@@ -18502,7 +18467,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(known.items.len(), 1);
|
||||
assert_eq!(known.items[0].subject, generic.worker_ref);
|
||||
assert_eq!(known.items[0].subject, generic_ref);
|
||||
assert_eq!(known.items[0].permissions, ["observe"]);
|
||||
|
||||
let Json(sessions) = scoped_list_worker_observation_sessions(
|
||||
@@ -18521,8 +18486,8 @@ mod tests {
|
||||
.iter()
|
||||
.any(|session| {
|
||||
session["subject"]["kind"] == "runtime_worker"
|
||||
&& session["subject"]["runtime_id"] == generic.worker_ref.runtime_id
|
||||
&& session["subject"]["worker_id"] == generic.worker_ref.worker_id
|
||||
&& session["subject"]["runtime_id"] == generic.runtime_id
|
||||
&& session["subject"]["worker_id"] == generic.worker_id
|
||||
})
|
||||
);
|
||||
let Json(capture) = scoped_capture_worker_observation_session(
|
||||
@@ -18532,8 +18497,8 @@ mod tests {
|
||||
}),
|
||||
observation_headers.clone(),
|
||||
Json(WorkerObservationSubjectRef::RuntimeWorker {
|
||||
runtime_id: generic.worker_ref.runtime_id.clone(),
|
||||
worker_id: generic.worker_ref.worker_id.clone(),
|
||||
runtime_id: generic.runtime_id.clone(),
|
||||
worker_id: generic.worker_id.clone(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18556,8 +18521,8 @@ mod tests {
|
||||
}),
|
||||
observation_headers.clone(),
|
||||
Json(WorkerObservationSubjectRef::RuntimeWorker {
|
||||
runtime_id: generic.worker_ref.runtime_id.clone(),
|
||||
worker_id: generic.worker_ref.worker_id.clone(),
|
||||
runtime_id: generic.runtime_id.clone(),
|
||||
worker_id: generic.worker_id.clone(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -18570,11 +18535,11 @@ mod tests {
|
||||
let mut unauthorized_headers = HeaderMap::new();
|
||||
unauthorized_headers.insert(
|
||||
"x-yoi-runtime-id",
|
||||
axum::http::HeaderValue::from_str(&generic.worker_ref.runtime_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&generic.runtime_id).unwrap(),
|
||||
);
|
||||
unauthorized_headers.insert(
|
||||
"x-yoi-worker-id",
|
||||
axum::http::HeaderValue::from_str(&generic.worker_ref.worker_id).unwrap(),
|
||||
axum::http::HeaderValue::from_str(&generic.worker_id).unwrap(),
|
||||
);
|
||||
let Json(unauthorized) = scoped_list_worker_observation_sessions(
|
||||
State(api.clone()),
|
||||
@@ -18596,10 +18561,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(existing.disposition, "existing");
|
||||
assert_eq!(
|
||||
existing.worker.unwrap().worker.worker_id,
|
||||
dedicated.worker.worker_id
|
||||
);
|
||||
assert_eq!(existing.worker.unwrap().worker_id, dedicated.worker_id);
|
||||
|
||||
let Json(status) = scoped_workspace_orchestrator_status(
|
||||
State(api),
|
||||
@@ -18607,10 +18569,7 @@ mod tests {
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
status.worker.unwrap().worker.worker_id,
|
||||
dedicated.worker.worker_id
|
||||
);
|
||||
assert_eq!(status.worker.unwrap().worker_id, dedicated.worker_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -20414,7 +20373,8 @@ mod tests {
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let orchestrator = started.worker.unwrap().worker;
|
||||
let orchestrator = started.worker.unwrap();
|
||||
let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id);
|
||||
execution.take_inputs();
|
||||
|
||||
let mut input = ticket::NewTicket::new("Bounded notification");
|
||||
@@ -20784,7 +20744,7 @@ mod tests {
|
||||
.iter()
|
||||
.any(|assignment| assignment.role == "coder")
|
||||
);
|
||||
assert_eq!(api.runtime.worker(&worker).unwrap().state, "cancelled");
|
||||
assert_eq!(api.runtime.worker(&worker).unwrap().state, "idle");
|
||||
|
||||
let Json(replayed) = scoped_cancel_ticket_implementation(State(api), path(), request())
|
||||
.await
|
||||
@@ -21525,8 +21485,8 @@ mod tests {
|
||||
.unwrap()
|
||||
.0
|
||||
.worker
|
||||
.expect("Workspace Orchestrator should be available")
|
||||
.worker;
|
||||
.expect("Workspace Orchestrator should be available");
|
||||
let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id);
|
||||
let _ = execution.take_inputs();
|
||||
let backend = browser_ticket_backend(&api).unwrap();
|
||||
let mut input = ticket::NewTicket::new("Queued notification");
|
||||
@@ -21722,7 +21682,7 @@ mod tests {
|
||||
Some(ticket_ref.id.as_str())
|
||||
);
|
||||
*api.orchestrator_attention_fingerprint.lock().unwrap() = None;
|
||||
let worker_id = started.worker.as_ref().unwrap().worker.worker_id.clone();
|
||||
let worker_id = started.worker.as_ref().unwrap().worker_id.clone();
|
||||
maybe_dispatch_orchestrator_turn_end(
|
||||
&api,
|
||||
&worker_id,
|
||||
@@ -21853,9 +21813,9 @@ mod tests {
|
||||
resolved_config_bundle: None,
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
let Json(first) = scoped_create_runtime_worker(
|
||||
State(api.clone()),
|
||||
@@ -22018,7 +21978,6 @@ mod tests {
|
||||
ticket_id: second_ticket.id.clone(),
|
||||
operation_id: "pending-spawn-operation".to_string(),
|
||||
}),
|
||||
resolved_control_operation: None,
|
||||
..request
|
||||
};
|
||||
pending_request.resolved_workspace_api =
|
||||
@@ -22136,9 +22095,9 @@ mod tests {
|
||||
resolved_config_bundle: None,
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
let Json(created) = scoped_create_runtime_worker(
|
||||
State(api.clone()),
|
||||
@@ -22846,7 +22805,8 @@ mod tests {
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let source = orchestrator.worker.unwrap().worker;
|
||||
let source = orchestrator.worker.unwrap();
|
||||
let source = RuntimeWorkerRef::new(&source.runtime_id, &source.worker_id);
|
||||
let verified_source = || crate::worker_source::VerifiedWorkerMutationSource {
|
||||
runtime_id: source.runtime_id.clone(),
|
||||
worker_id: source.worker_id.clone(),
|
||||
@@ -22922,7 +22882,8 @@ mod tests {
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let source = orchestrator.worker.unwrap().worker;
|
||||
let source = orchestrator.worker.unwrap();
|
||||
let source = RuntimeWorkerRef::new(&source.runtime_id, &source.worker_id);
|
||||
|
||||
let spawned = api
|
||||
.spawn_workspace_worker(
|
||||
@@ -27445,9 +27406,9 @@ mod tests {
|
||||
resolved_config_bundle: Some(runtime_test_bundle()),
|
||||
resolved_worker_observation_enabled: false,
|
||||
resolved_worker_observation_grants: Vec::new(),
|
||||
resolved_control_operation: None,
|
||||
resolved_workspace_api: None,
|
||||
resolved_memory_settings: None,
|
||||
resolved_control_operation: None,
|
||||
};
|
||||
let spawned = api
|
||||
.spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request)
|
||||
|
||||
Reference in New Issue
Block a user