refactor: share Worker launch REST DTOs

This commit is contained in:
2026-09-03 14:23:33 +09:00
parent 9fb1b90856
commit 1f32c693df
18 changed files with 1437 additions and 313 deletions
+23 -45
View File
@@ -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,
+6 -1
View File
@@ -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());
}
+320
View File
@@ -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);
+140 -170
View File
@@ -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>,
@@ -7758,9 +7669,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 {
@@ -8711,18 +8622,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![
@@ -9055,9 +8968,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() {
@@ -9085,6 +8998,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,
@@ -9095,13 +9044,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,
}
}
@@ -12458,6 +12414,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,
@@ -12467,7 +12432,6 @@ async fn create_workspace_worker(
initial_submit,
working_directory,
control_operation_id: _,
resolved_control_operation,
} = request;
let config_state = api
.config_store
@@ -12765,10 +12729,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(),
})
}
@@ -14810,7 +14778,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();
@@ -17210,9 +17182,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)
@@ -17456,9 +17428,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!(
@@ -17505,7 +17477,6 @@ mod tests {
}],
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await
@@ -17521,7 +17492,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(
@@ -17531,7 +17503,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]
@@ -17561,7 +17533,6 @@ mod tests {
}],
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await;
@@ -17595,7 +17566,6 @@ mod tests {
initial_submit: Vec::new(),
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await
@@ -17603,11 +17573,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();
@@ -17631,14 +17601,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(_))
));
@@ -17650,10 +17620,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(_))
));
@@ -17722,8 +17693,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
@@ -17823,12 +17794,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(
@@ -17854,7 +17825,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(
@@ -17878,7 +17848,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(
@@ -17908,7 +17879,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");
}
@@ -17930,7 +17904,6 @@ mod tests {
initial_submit: Vec::new(),
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await
@@ -17948,7 +17921,6 @@ mod tests {
initial_submit: Vec::new(),
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await
@@ -17969,14 +17941,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()],
@@ -17989,11 +17963,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()),
@@ -18005,7 +17979,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(
@@ -18024,8 +17998,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(
@@ -18035,8 +18009,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
@@ -18059,8 +18033,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
@@ -18073,11 +18047,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()),
@@ -18099,10 +18073,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),
@@ -18110,10 +18081,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]
@@ -19917,7 +19885,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");
@@ -21028,8 +20997,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");
@@ -21225,7 +21194,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,
@@ -21356,9 +21325,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()),
@@ -21521,7 +21490,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 =
@@ -21639,9 +21607,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()),
@@ -22349,7 +22317,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(),
@@ -22425,7 +22394,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(
@@ -26947,9 +26917,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)