From a2781f57e91c843277467810cd89765c7965e9f3 Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 5 Aug 2026 21:26:55 +0900 Subject: [PATCH] worker: structure runtime worker identities --- crates/worker-runtime/src/catalog.rs | 7 +- crates/worker-runtime/src/identity.rs | 52 ++ crates/workspace-server/src/hosts.rs | 314 +++---- crates/workspace-server/src/lib.rs | 11 +- crates/workspace-server/src/observation.rs | 102 +-- .../workspace-server/src/resource_broker.rs | 64 +- crates/workspace-server/src/server.rs | 807 ++++++++---------- crates/workspace-server/src/store.rs | 181 ++-- .../src/workspace_subscription.rs | 53 +- 9 files changed, 775 insertions(+), 816 deletions(-) diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index d1a513e8..b299ee6e 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -1,4 +1,4 @@ -use crate::identity::{WorkerId, WorkerRef}; +use crate::identity::{RuntimeWorkerRef, WorkerId, WorkerRef}; use crate::interaction::WorkerInput; use crate::profile_archive::{ProfileSourceArchive, ProfileSourceArchiveRef}; use serde::{Deserialize, Serialize}; @@ -132,9 +132,8 @@ pub struct WorkingDirectoryCleanupTarget { #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct WorkingDirectoryOccupancy { - pub runtime_id: String, - pub runtime_worker_id: u64, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub display_name: String, pub linked_at: String, } diff --git a/crates/worker-runtime/src/identity.rs b/crates/worker-runtime/src/identity.rs index b5cd80f2..ce3969a6 100644 --- a/crates/worker-runtime/src/identity.rs +++ b/crates/worker-runtime/src/identity.rs @@ -30,6 +30,32 @@ impl fmt::Display for WorkerId { } } +/// Backend-visible Worker identity, namespaced by the Runtime that owns the Worker record. +/// +/// This is intentionally distinct from [`WorkerRef`], which is meaningful only inside one +/// Runtime. Do not flatten this reference into a concatenated string for authority decisions. +#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub struct RuntimeWorkerRef { + pub runtime_id: String, + pub worker_id: String, +} + +impl RuntimeWorkerRef { + pub fn new(runtime_id: impl Into, worker_id: impl Into) -> Self { + Self { + runtime_id: runtime_id.into(), + worker_id: worker_id.into(), + } + } + + pub fn local_worker_ref(&self) -> Result { + self.worker_id + .parse::() + .map(WorkerId::new) + .map(WorkerRef::new) + } +} + /// Runtime-local authority reference for Worker operations. #[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] pub struct WorkerRef { @@ -41,3 +67,29 @@ impl WorkerRef { Self { worker_id } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn runtime_worker_ref_preserves_structured_identity_and_json_fields() { + let worker = RuntimeWorkerRef::new("arcadia", "30"); + assert_eq!(worker.runtime_id, "arcadia"); + assert_eq!(worker.worker_id, "30"); + assert_eq!( + worker.local_worker_ref().unwrap(), + WorkerRef::new(WorkerId::new(30)) + ); + assert_eq!( + serde_json::to_value(&worker).unwrap(), + serde_json::json!({"runtime_id": "arcadia", "worker_id": "30"}) + ); + } + + #[test] + fn runtime_worker_ref_does_not_treat_composite_text_as_local_worker_id() { + let worker = RuntimeWorkerRef::new("arcadia", "embedded-worker-runtime-5"); + assert!(worker.local_worker_ref().is_err()); + } +} diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 7eb67d2f..582d65dd 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -1,5 +1,5 @@ use crate::Error; -use crate::resource_broker::BackendResourceBroker; +use crate::resource_broker::{BackendResourceBroker, BackendResourceTarget}; use chrono::Utc; use reqwest::blocking::{Client as BlockingHttpClient, RequestBuilder}; use reqwest::header::{AUTHORIZATION, CONTENT_TYPE}; @@ -45,7 +45,9 @@ use worker_runtime::http_server::{ RuntimeHttpWorkerWorkspaceApiRequest, RuntimeHttpWorkersResponse, RuntimeHttpWorkingDirectoriesResponse, RuntimeHttpWorkingDirectoryResponse, }; -use worker_runtime::identity::{WorkerId as EmbeddedWorkerId, WorkerRef as EmbeddedWorkerRef}; +use worker_runtime::identity::{ + RuntimeWorkerRef, WorkerId as EmbeddedWorkerId, WorkerRef as EmbeddedWorkerRef, +}; use worker_runtime::interaction::{ WorkerInput as EmbeddedWorkerInput, WorkerInputKind as EmbeddedWorkerInputKind, }; @@ -235,8 +237,8 @@ pub struct WorkerCapabilitySummary { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerSummary { - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub host_id: String, /// Human-readable display name. This is not identity and may be duplicated. pub display_name: String, @@ -487,8 +489,8 @@ pub struct WorkerLifecycleRequest { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerLifecycleResult { pub state: WorkerOperationState, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub diagnostics: Vec, } @@ -505,8 +507,8 @@ pub enum WorkerInputKind { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerDeleteResult { pub state: WorkerOperationState, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub deleted: bool, pub diagnostics: Vec, } @@ -529,8 +531,8 @@ pub struct WorkerCompletionsRequest { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerCompletionsResult { - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub kind: protocol::CompletionKind, pub prefix: String, pub entries: Vec, @@ -540,8 +542,8 @@ pub struct WorkerCompletionsResult { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerInputResult { pub state: WorkerOperationState, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub diagnostics: Vec, } @@ -561,8 +563,7 @@ pub enum RuntimeRegistryError { UnknownRuntime(String), UnknownHost(String), UnknownWorker { - runtime_id: String, - worker_id: String, + worker: RuntimeWorkerRef, }, RuntimeOperationFailed { runtime_id: String, @@ -579,10 +580,10 @@ impl RuntimeRegistryError { } Self::UnknownRuntime(runtime_id) => format!("unknown runtime `{runtime_id}`"), Self::UnknownHost(host_id) => format!("unknown host `{host_id}`"), - Self::UnknownWorker { - runtime_id, - worker_id, - } => format!("unknown worker `{worker_id}` in runtime `{runtime_id}`"), + Self::UnknownWorker { worker } => format!( + "unknown worker `{}` in runtime `{}`", + worker.worker_id, worker.runtime_id + ), Self::RuntimeOperationFailed { message, .. } => message.clone(), } } @@ -595,13 +596,7 @@ impl RuntimeRegistryError { }, Self::UnknownRuntime(runtime_id) => Error::UnknownRuntime(runtime_id), Self::UnknownHost(host_id) => Error::UnknownHost(host_id), - Self::UnknownWorker { - runtime_id, - worker_id, - } => Error::UnknownWorker { - runtime_id, - worker_id, - }, + Self::UnknownWorker { worker } => Error::UnknownWorker { worker }, Self::RuntimeOperationFailed { runtime_id, code, @@ -798,8 +793,7 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { ) -> WorkerLifecycleResult { WorkerLifecycleResult { state: WorkerOperationState::Unsupported, - runtime_id: self.runtime_id().to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id().to_string(), worker_id.to_string()), diagnostics: vec![diagnostic( "worker_stop_pending", DiagnosticSeverity::Info, @@ -817,8 +811,7 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { ) -> WorkerLifecycleResult { WorkerLifecycleResult { state: WorkerOperationState::Unsupported, - runtime_id: self.runtime_id().to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id().to_string(), worker_id.to_string()), diagnostics: vec![diagnostic( "worker_cancel_pending", DiagnosticSeverity::Info, @@ -832,8 +825,7 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { fn delete_worker(&self, worker_id: &str) -> WorkerDeleteResult { WorkerDeleteResult { state: WorkerOperationState::Unsupported, - runtime_id: self.runtime_id().to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id().to_string(), worker_id.to_string()), deleted: false, diagnostics: vec![diagnostic( "worker_delete_unsupported", @@ -853,8 +845,7 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { fn send_input(&self, worker_id: &str, _request: WorkerInputRequest) -> WorkerInputResult { WorkerInputResult { state: WorkerOperationState::Unsupported, - runtime_id: self.runtime_id().to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id().to_string(), worker_id.to_string()), diagnostics: vec![diagnostic( "worker_input_pending", DiagnosticSeverity::Info, @@ -871,8 +862,7 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { request: WorkerCompletionsRequest, ) -> WorkerCompletionsResult { WorkerCompletionsResult { - runtime_id: self.runtime_id().to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id().to_string(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries: Vec::new(), @@ -1099,11 +1089,9 @@ impl RuntimeRegistry { } } - pub fn worker( - &self, - runtime_id: &str, - worker_id: &str, - ) -> Result { + pub fn worker(&self, worker: &RuntimeWorkerRef) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1116,9 +1104,10 @@ impl RuntimeRegistry { pub fn restore_worker( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1127,10 +1116,11 @@ impl RuntimeRegistry { pub fn replace_worker_workspace_api( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, workspace_api: WorkspaceApiRef, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1241,10 +1231,11 @@ impl RuntimeRegistry { pub fn send_protocol_method( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, method: protocol::Method, ) -> Result, RuntimeRegistryError> { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1261,10 +1252,11 @@ impl RuntimeRegistry { pub fn send_input( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, request: WorkerInputRequest, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1281,10 +1273,11 @@ impl RuntimeRegistry { pub fn worker_completions( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, request: WorkerCompletionsRequest, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1301,10 +1294,11 @@ impl RuntimeRegistry { pub fn stop_worker( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, request: WorkerLifecycleRequest, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1321,10 +1315,11 @@ impl RuntimeRegistry { pub fn cancel_worker( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, request: WorkerLifecycleRequest, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1341,9 +1336,10 @@ impl RuntimeRegistry { pub fn delete_worker( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; @@ -1360,17 +1356,17 @@ impl RuntimeRegistry { pub fn observation_source( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> Result { + let runtime_id = worker.runtime_id.as_str(); + let worker_id = worker.worker_id.as_str(); validate_backend_identifier("runtime_id", runtime_id)?; validate_backend_identifier("worker_id", worker_id)?; let runtime = self.runtime(runtime_id)?; runtime .observation_source(worker_id) .ok_or_else(|| RuntimeRegistryError::UnknownWorker { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: worker.clone(), }) } @@ -1483,8 +1479,7 @@ impl EmbeddedWorkerRuntime { true, ); WorkerSummary { - runtime_id: self.runtime_id.clone(), - worker_id, + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id.clone()), host_id: self.host_id.clone(), display_name: display.display_name.clone(), label: display.display_name, @@ -1522,8 +1517,7 @@ impl EmbeddedWorkerRuntime { true, ); WorkerSummary { - runtime_id: self.runtime_id.clone(), - worker_id, + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id.clone()), host_id: self.host_id.clone(), display_name: display.display_name.clone(), label: display.display_name, @@ -1974,8 +1968,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { match self.runtime.stop_worker(&worker_ref, request.reason) { Ok(_) => WorkerLifecycleResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), diagnostics: Vec::new(), }, Err(error) => embedded_lifecycle_rejected( @@ -2018,8 +2011,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { match self.runtime.cancel_worker(&worker_ref, request.reason) { Ok(_) => WorkerLifecycleResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), diagnostics: Vec::new(), }, Err(error) => embedded_lifecycle_rejected( @@ -2034,8 +2026,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { let Some(worker_ref) = self.worker_ref(worker_id) else { return WorkerDeleteResult { state: WorkerOperationState::Rejected, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), deleted: false, diagnostics: vec![diagnostic( "embedded_worker_id_invalid", @@ -2047,15 +2038,16 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { match self.runtime.delete_worker(&worker_ref) { Ok(result) => WorkerDeleteResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: result.worker_id.to_string(), + worker: RuntimeWorkerRef::new( + self.runtime_id.clone(), + result.worker_id.to_string(), + ), deleted: result.deleted, diagnostics: Vec::new(), }, Err(error) => WorkerDeleteResult { state: WorkerOperationState::Rejected, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), deleted: false, diagnostics: vec![embedded_runtime_diagnostic(&error)], }, @@ -2072,8 +2064,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { } Some(crate::observation::RuntimeObservationSource::embedded( crate::observation::EmbeddedRuntimeObservationSource { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id), runtime: self.runtime.clone(), worker_ref, }, @@ -2096,8 +2087,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { } let Some(worker_ref) = self.worker_ref(worker_id) else { return Err(RuntimeRegistryError::UnknownWorker { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id), }); }; self.runtime @@ -2148,8 +2138,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { match self.runtime.send_input(&worker_ref, input) { Ok(_) => WorkerInputResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), diagnostics: Vec::new(), }, Err(error) => embedded_input_rejected( @@ -2167,8 +2156,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { ) -> WorkerCompletionsResult { if !self.execution_enabled { return WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries: Vec::new(), @@ -2183,8 +2171,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { } let Some(worker_ref) = self.worker_ref(worker_id) else { return WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries: Vec::new(), @@ -2200,16 +2187,14 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { .worker_completions(&worker_ref, request.kind, &request.prefix) { Ok(entries) => WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries, diagnostics: Vec::new(), }, Err(error) => WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries: Vec::new(), @@ -2572,8 +2557,7 @@ impl RemoteWorkerRuntime { false, ); WorkerSummary { - runtime_id: self.runtime_id.clone(), - worker_id, + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id.clone()), host_id: self.host_id.clone(), display_name: display.display_name.clone(), label: display.display_name, @@ -2615,8 +2599,7 @@ impl RemoteWorkerRuntime { false, ); WorkerSummary { - runtime_id: self.runtime_id.clone(), - worker_id, + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id.clone()), host_id: self.host_id.clone(), display_name: display.display_name.clone(), label: display.display_name, @@ -2655,8 +2638,7 @@ impl RemoteWorkerRuntime { ) -> WorkerLifecycleResult { WorkerLifecycleResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), diagnostics: vec![diagnostic( "remote_runtime_lifecycle_accepted", DiagnosticSeverity::Info, @@ -3073,15 +3055,16 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { { Ok(response) => WorkerDeleteResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: response.worker.worker_id.to_string(), + worker: RuntimeWorkerRef::new( + self.runtime_id.clone(), + response.worker.worker_id.to_string(), + ), deleted: response.worker.deleted, diagnostics: Vec::new(), }, Err(diagnostic) => WorkerDeleteResult { state: WorkerOperationState::Rejected, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), deleted: false, diagnostics: vec![diagnostic], }, @@ -3094,8 +3077,7 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { ) -> Option { Some(crate::observation::RuntimeObservationSource::remote_ws( crate::observation::RuntimeObservationSourceConfig { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(&self.runtime_id, worker_id), endpoint: self.ws_endpoint(worker_id), bearer_token: self .runtime_capability_token(&format!("/v1/workers/{worker_id}/protocol")) @@ -3122,8 +3104,7 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { ) { Ok(_) => WorkerInputResult { state: WorkerOperationState::Accepted, - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), diagnostics: Vec::new(), }, Err(diagnostic) => remote_input_rejected(&self.runtime_id, worker_id, diagnostic), @@ -3144,16 +3125,14 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { &request, ) { Ok(response) => WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: response.kind, prefix: response.prefix, entries: response.entries, diagnostics: Vec::new(), }, Err(diagnostic) => WorkerCompletionsResult { - runtime_id: self.runtime_id.clone(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(self.runtime_id.clone(), worker_id.to_string()), kind: request.kind, prefix: request.prefix, entries: Vec::new(), @@ -3247,10 +3226,12 @@ fn profile_source_archive_http_source( backend_base_url: &str, ) -> Result { let archive = profile_source_archive_for_request(request, profile)?; + let target = runtime_id + .map(BackendResourceTarget::Runtime) + .unwrap_or(BackendResourceTarget::Workspace); let _handle = resource_broker.issue_profile_source_archive_handle( workspace_id.to_string(), - runtime_id, - None, + target, archive.clone(), ); let etag = format!("\"profile-source:{}\"", archive.reference.digest); @@ -3294,10 +3275,12 @@ fn builtin_profile_config_bundle( let (profile_source_archive, profile_source_archive_handle) = match archive_transport { ProfileSourceArchiveTransport::Inline => (Some(archive), None), ProfileSourceArchiveTransport::BackendResourceHandle => { + let target = runtime_id + .map(BackendResourceTarget::Runtime) + .unwrap_or(BackendResourceTarget::Workspace); let handle = resource_broker.issue_profile_source_archive_handle( workspace_id.to_string(), - runtime_id, - None, + target, archive, ); (None, Some(handle)) @@ -3515,8 +3498,7 @@ fn embedded_input_rejected( ) -> WorkerInputResult { WorkerInputResult { state: WorkerOperationState::Rejected, - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id.to_string(), worker_id.to_string()), diagnostics: vec![diagnostic], } } @@ -3528,8 +3510,7 @@ fn remote_input_rejected( ) -> WorkerInputResult { WorkerInputResult { state: WorkerOperationState::Rejected, - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id.to_string(), worker_id.to_string()), diagnostics: vec![diagnostic], } } @@ -3541,8 +3522,7 @@ fn embedded_lifecycle_rejected( ) -> WorkerLifecycleResult { WorkerLifecycleResult { state: WorkerOperationState::Rejected, - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id.to_string(), worker_id.to_string()), diagnostics: vec![diagnostic], } } @@ -3554,8 +3534,7 @@ fn remote_lifecycle_rejected( ) -> WorkerLifecycleResult { WorkerLifecycleResult { state: WorkerOperationState::Rejected, - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id.to_string(), worker_id.to_string()), diagnostics: vec![diagnostic], } } @@ -3779,8 +3758,7 @@ fn operation_failed_or_unknown_worker( message: diagnostic.message, }) .unwrap_or_else(|| RuntimeRegistryError::UnknownWorker { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id, worker_id), }) } @@ -3887,8 +3865,7 @@ fn worker_spawn_intent_label(intent: &WorkerSpawnIntent) -> &'static str { pub fn placeholder_worker(host_id: impl Into) -> WorkerSummary { let host_id = host_id.into(); WorkerSummary { - runtime_id: "placeholder".to_string(), - worker_id: "worker-placeholder".to_string(), + worker: RuntimeWorkerRef::new("placeholder", "worker-placeholder"), host_id, display_name: "Worker runtime actions are not implemented".to_string(), label: "Worker runtime actions are not implemented".to_string(), @@ -3951,6 +3928,29 @@ mod tests { } } + #[test] + fn worker_summary_keeps_flat_wire_identity_while_using_structured_internal_identity() { + let summary = placeholder_worker("placeholder"); + assert_eq!( + summary.worker, + RuntimeWorkerRef::new("placeholder", "worker-placeholder") + ); + let value = serde_json::to_value(summary).unwrap(); + assert_eq!(value["runtime_id"], "placeholder"); + assert_eq!(value["worker_id"], "worker-placeholder"); + assert!(value.get("worker").is_none()); + + let lifecycle = WorkerLifecycleResult { + state: WorkerOperationState::Accepted, + worker: RuntimeWorkerRef::new("arcadia", "30"), + diagnostics: Vec::new(), + }; + let value = serde_json::to_value(lifecycle).unwrap(); + assert_eq!(value["runtime_id"], "arcadia"); + assert_eq!(value["worker_id"], "30"); + assert!(value.get("worker").is_none()); + } + #[test] fn embedded_orchestrator_profile_enables_workdir_and_worker_authority() { let root = tempfile::tempdir().unwrap(); @@ -4289,8 +4289,7 @@ mod tests { runtime_id: runtime_id.to_string(), host_id: host_id.to_string(), workers: vec![WorkerSummary { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id, worker_id), host_id: host_id.to_string(), display_name: label.to_string(), label: label.to_string(), @@ -4382,7 +4381,7 @@ mod tests { worker: self .workers .iter() - .find(|worker| worker.worker_id == worker_id) + .find(|worker| worker.worker.worker_id == worker_id) .cloned(), diagnostics: Vec::new(), } @@ -4406,13 +4405,17 @@ mod tests { )), ]); - let from_runtime_b = registry.worker("runtime-b", "shared-worker").unwrap(); - assert_eq!(from_runtime_b.runtime_id, "runtime-b"); + let from_runtime_b = registry + .worker(&RuntimeWorkerRef::new("runtime-b", "shared-worker")) + .unwrap(); + assert_eq!(from_runtime_b.worker.runtime_id, "runtime-b"); assert_eq!(from_runtime_b.host_id, "host-b"); assert_eq!(from_runtime_b.label, "worker from runtime b"); - let from_runtime_a = registry.worker("runtime-a", "shared-worker").unwrap(); - assert_eq!(from_runtime_a.runtime_id, "runtime-a"); + let from_runtime_a = registry + .worker(&RuntimeWorkerRef::new("runtime-a", "shared-worker")) + .unwrap(); + assert_eq!(from_runtime_a.worker.runtime_id, "runtime-a"); assert_eq!(from_runtime_a.host_id, "host-a"); assert_eq!(from_runtime_a.label, "worker from runtime a"); } @@ -4436,7 +4439,7 @@ mod tests { let listed = registry.list_workers_for_runtime("runtime-b", 10).unwrap(); assert_eq!(listed.items.len(), 1); - assert_eq!(listed.items[0].runtime_id, "runtime-b"); + assert_eq!(listed.items[0].worker.runtime_id, "runtime-b"); assert_eq!(listed.items[0].host_id, "host-b"); assert_eq!(listed.items[0].label, "worker from runtime b"); } @@ -4455,7 +4458,9 @@ mod tests { Some("builtin:companion") ); - let worker = registry.worker("runtime-a", "worker-a").unwrap(); + let worker = registry + .worker(&RuntimeWorkerRef::new("runtime-a", "worker-a")) + .unwrap(); assert_eq!(worker.profile.as_deref(), Some("builtin:companion")); } @@ -4468,7 +4473,9 @@ mod tests { "worker from runtime a", ))]); - let unknown_runtime = registry.worker("runtime-missing", "worker-a").unwrap_err(); + let unknown_runtime = registry + .worker(&RuntimeWorkerRef::new("runtime-missing", "worker-a")) + .unwrap_err(); assert_eq!( unknown_runtime, RuntimeRegistryError::UnknownRuntime("runtime-missing".to_string()) @@ -4478,18 +4485,19 @@ mod tests { Error::UnknownRuntime(runtime_id) if runtime_id == "runtime-missing" )); - let unknown_worker = registry.worker("runtime-a", "999").unwrap_err(); + let unknown_worker = registry + .worker(&RuntimeWorkerRef::new("runtime-a", "999")) + .unwrap_err(); assert_eq!( unknown_worker, RuntimeRegistryError::UnknownWorker { - runtime_id: "runtime-a".to_string(), - worker_id: "999".to_string(), + worker: RuntimeWorkerRef::new("runtime-a", "999"), } ); assert!(matches!( unknown_worker.into_error(), - Error::UnknownWorker { runtime_id, worker_id } - if runtime_id == "runtime-a" && worker_id == "999" + Error::UnknownWorker { worker } + if worker == RuntimeWorkerRef::new("runtime-a", "999") )); } @@ -4591,7 +4599,7 @@ mod tests { assert!(worker.capabilities.can_stop); let input = runtime.send_input( - &worker.worker_id, + &worker.worker.worker_id, WorkerInputRequest { kind: WorkerInputKind::User, content: "hello".to_string(), @@ -4603,7 +4611,7 @@ mod tests { let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2); loop { let detail = runtime - .worker(&worker.worker_id) + .worker(&worker.worker.worker_id) .worker .expect("worker detail"); if detail.state == "idle" { @@ -4671,15 +4679,14 @@ mod tests { .any(|evidence| evidence.kind == "embedded_runtime_backend_internal_projection") ); let worker = spawned.worker.expect("created embedded worker"); - assert_eq!(worker.runtime_id, EMBEDDED_RUNTIME_ID); + assert_eq!(worker.worker.runtime_id, EMBEDDED_RUNTIME_ID); assert_eq!(worker.workspace.visibility, "backend_internal"); assert_eq!(worker.workspace.identity, "runtime_registry_worker"); assert_eq!(worker.implementation.kind, "embedded_worker_runtime"); assert_eq!(worker.profile.as_deref(), Some("builtin:coder")); let input = registry .send_input( - EMBEDDED_RUNTIME_ID, - &worker.worker_id, + &worker.worker, WorkerInputRequest { kind: WorkerInputKind::User, content: "hello embedded runtime".to_string(), @@ -4688,12 +4695,10 @@ mod tests { ) .unwrap(); assert_eq!(input.state, WorkerOperationState::Accepted); - assert_eq!(input.runtime_id, EMBEDDED_RUNTIME_ID); - assert_eq!(input.worker_id, worker.worker_id); + assert_eq!(input.worker.runtime_id, EMBEDDED_RUNTIME_ID); + assert_eq!(input.worker.worker_id, worker.worker.worker_id); - let detail = registry - .worker(EMBEDDED_RUNTIME_ID, &worker.worker_id) - .unwrap(); + let detail = registry.worker(&worker.worker).unwrap(); let json = serde_json::to_string(&(embedded_summary, worker, input, detail)).unwrap(); for forbidden in [ @@ -4871,7 +4876,7 @@ mod tests { ); let observation = registry - .observation_source("remote:primary", "1") + .observation_source(&RuntimeWorkerRef::new("remote:primary", "1")) .expect("remote runtime exposes backend-owned WS observation source"); let crate::observation::RuntimeObservationSource::RemoteWs(observation) = observation else { @@ -4883,8 +4888,8 @@ mod tests { let workers = registry.list_workers(10); assert_eq!(workers.items.len(), 1); - assert_eq!(workers.items[0].runtime_id, "remote:primary"); - assert_eq!(workers.items[0].worker_id, "1"); + assert_eq!(workers.items[0].worker.runtime_id, "remote:primary"); + assert_eq!(workers.items[0].worker.worker_id, "1"); assert_eq!( workers.items[0].implementation.kind, "remote_worker_runtime" @@ -4897,8 +4902,7 @@ mod tests { let input = registry .send_input( - "remote:primary", - "1", + &RuntimeWorkerRef::new("remote:primary", "1"), WorkerInputRequest { kind: WorkerInputKind::User, content: "hello remote".to_string(), @@ -4975,7 +4979,9 @@ mod tests { assert_eq!(workers.items[2].state, "paused"); assert_eq!(workers.items[3].state, "idle"); - let stopped_detail = registry.worker("remote:primary", "1").unwrap(); + let stopped_detail = registry + .worker(&RuntimeWorkerRef::new("remote:primary", "1")) + .unwrap(); assert!(!stopped_detail.capabilities.can_stop); assert_eq!(stopped_detail.state, "stopped"); @@ -5154,7 +5160,7 @@ mod tests { ); let error = registry - .worker("remote:primary", "999") + .worker(&RuntimeWorkerRef::new("remote:primary", "999")) .expect_err("auth failure is a backend operation error"); assert!(matches!( error, diff --git a/crates/workspace-server/src/lib.rs b/crates/workspace-server/src/lib.rs index ff111ff6..36a3af02 100644 --- a/crates/workspace-server/src/lib.rs +++ b/crates/workspace-server/src/lib.rs @@ -43,6 +43,8 @@ pub use repositories::{ pub use server::{AuthConfig, ServerConfig, WorkspaceApi, build_router, serve}; pub use store::{ControlPlaneStore, SqliteWorkspaceStore, WorkspaceRecord}; +use worker_runtime::identity::RuntimeWorkerRef; + pub type Result = std::result::Result; #[derive(Debug, thiserror::Error)] @@ -65,11 +67,8 @@ pub enum Error { UnknownHost(String), #[error("unknown runtime `{0}`")] UnknownRuntime(String), - #[error("unknown worker `{worker_id}` in runtime `{runtime_id}`")] - UnknownWorker { - runtime_id: String, - worker_id: String, - }, + #[error("unknown worker `{}` in runtime `{}`", worker.worker_id, worker.runtime_id)] + UnknownWorker { worker: RuntimeWorkerRef }, #[error("invalid runtime {kind} `{value}`")] InvalidRuntimeIdentifier { kind: String, value: String }, #[error("worker name is reserved for a dedicated Workspace service: {0}")] @@ -93,6 +92,8 @@ pub enum Error { TicketAssignmentConflict(String), #[error("Workdir attachment conflict: {0}")] WorkdirAttachmentConflict(String), + #[error("Registry inconsistency: {0}")] + RegistryInconsistency(String), #[error("Worker source identity is invalid: {0}")] WorkerSourceIdentity(String), #[error("workspace identity error: {0}")] diff --git a/crates/workspace-server/src/observation.rs b/crates/workspace-server/src/observation.rs index 756108a6..7bba5712 100644 --- a/crates/workspace-server/src/observation.rs +++ b/crates/workspace-server/src/observation.rs @@ -1,7 +1,7 @@ use std::collections::{BTreeMap, VecDeque}; use std::sync::Arc; -use worker_runtime::identity::WorkerRef; +use worker_runtime::identity::{RuntimeWorkerRef, WorkerRef}; use worker_runtime::observation::{WorkerObservationCursor, WorkerObservationEvent}; use axum::http::StatusCode; @@ -15,8 +15,7 @@ use tokio_tungstenite::tungstenite::{Error as TungsteniteError, Message as Tungs /// Backend-private source for a runtime worker observation stream. #[derive(Clone, PartialEq, Eq)] pub struct RuntimeObservationSourceConfig { - pub runtime_id: String, - pub worker_id: String, + pub worker: RuntimeWorkerRef, pub endpoint: String, pub bearer_token: Option, } @@ -24,8 +23,8 @@ pub struct RuntimeObservationSourceConfig { impl std::fmt::Debug for RuntimeObservationSourceConfig { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("RuntimeObservationSourceConfig") - .field("runtime_id", &self.runtime_id) - .field("worker_id", &self.worker_id) + .field("runtime_id", &self.worker.runtime_id) + .field("worker_id", &self.worker.worker_id) .field("endpoint", &"") .field( "bearer_token", @@ -37,8 +36,7 @@ impl std::fmt::Debug for RuntimeObservationSourceConfig { #[derive(Clone)] pub struct EmbeddedRuntimeObservationSource { - pub runtime_id: String, - pub worker_id: String, + pub worker: RuntimeWorkerRef, pub runtime: worker_runtime::Runtime, pub worker_ref: WorkerRef, } @@ -60,15 +58,15 @@ impl RuntimeObservationSource { pub fn runtime_id(&self) -> &str { match self { - Self::RemoteWs(config) => &config.runtime_id, - Self::Embedded(source) => &source.runtime_id, + Self::RemoteWs(config) => &config.worker.runtime_id, + Self::Embedded(source) => &source.worker.runtime_id, } } pub fn worker_id(&self) -> &str { match self { - Self::RemoteWs(config) => &config.worker_id, - Self::Embedded(source) => &source.worker_id, + Self::RemoteWs(config) => &config.worker.worker_id, + Self::Embedded(source) => &source.worker.worker_id, } } } @@ -78,8 +76,8 @@ impl std::fmt::Debug for RuntimeObservationSource { match self { Self::RemoteWs(config) => formatter .debug_struct("RemoteRuntimeObservationSource") - .field("runtime_id", &config.runtime_id) - .field("worker_id", &config.worker_id) + .field("runtime_id", &config.worker.runtime_id) + .field("worker_id", &config.worker.worker_id) .field("endpoint", &"") .field( "bearer_token", @@ -88,8 +86,8 @@ impl std::fmt::Debug for RuntimeObservationSource { .finish(), Self::Embedded(source) => formatter .debug_struct("EmbeddedRuntimeObservationSource") - .field("runtime_id", &source.runtime_id) - .field("worker_id", &source.worker_id) + .field("runtime_id", &source.worker.runtime_id) + .field("worker_id", &source.worker.worker_id) .finish(), } } @@ -98,8 +96,8 @@ impl std::fmt::Debug for RuntimeObservationSource { /// Event consumed from a Runtime-owned worker observation WebSocket. #[derive(Clone, Debug, Serialize, Deserialize)] pub struct RuntimeObservationUpstreamEvent { - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker: RuntimeWorkerRef, pub runtime_event_id: String, pub payload: protocol::Event, } @@ -132,11 +130,7 @@ impl ObservationProxyError { } } -#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)] -struct ObservationKey { - runtime_id: String, - worker_id: String, -} +type ObservationKey = RuntimeWorkerRef; /// Backend-owned in-memory v0 observation proxy state. #[derive(Clone)] @@ -156,15 +150,7 @@ impl BackendObservationProxy { pub fn new(sources: Vec) -> Self { let sources = sources .into_iter() - .map(|source| { - ( - ObservationKey { - runtime_id: source.runtime_id.clone(), - worker_id: source.worker_id.clone(), - }, - source, - ) - }) + .map(|source| (source.worker.clone(), source)) .collect(); Self { sources: Arc::new(sources), @@ -173,19 +159,16 @@ impl BackendObservationProxy { pub fn source( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> Result { self.sources - .get(&ObservationKey { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), - }) + .get(worker) .cloned() .map(RuntimeObservationSource::remote_ws) .ok_or_else(|| { ObservationProxyError::WorkerNotFound(format!( - "worker {worker_id} is not registered for runtime {runtime_id}" + "worker {} is not registered for runtime {}", + worker.worker_id, worker.runtime_id )) }) } @@ -209,8 +192,7 @@ fn map_runtime_connect_error(error: TungsteniteError) -> ObservationProxyError { } pub struct RuntimeWsObservationClient { - runtime_id: String, - worker_id: String, + worker: RuntimeWorkerRef, stream: tokio_tungstenite::WebSocketStream< tokio_tungstenite::MaybeTlsStream, >, @@ -240,8 +222,7 @@ impl RuntimeWsObservationClient { .await .map_err(map_runtime_connect_error)?; Ok(Self { - runtime_id: source.runtime_id.clone(), - worker_id: source.worker_id.clone(), + worker: source.worker.clone(), stream, }) } @@ -291,8 +272,7 @@ impl RuntimeWsObservationClient { )) })?; return Ok(RuntimeObservationUpstreamEvent { - runtime_id: self.runtime_id.clone(), - worker_id: self.worker_id.clone(), + worker: self.worker.clone(), runtime_event_id: "protocol".to_string(), payload, }); @@ -330,8 +310,7 @@ impl RuntimeObservationClient { } pub struct EmbeddedObservationClient { - runtime_id: String, - worker_id: String, + worker: RuntimeWorkerRef, worker_ref: WorkerRef, cursor: WorkerObservationCursor, receiver: tokio::sync::broadcast::Receiver, @@ -346,7 +325,7 @@ impl EmbeddedObservationClient { .map_err(|err| { ObservationProxyError::WorkerNotFound(format!( "embedded Worker '{}' is not observable: {err}", - source.worker_id + source.worker.worker_id )) })?; let receiver = source @@ -355,7 +334,7 @@ impl EmbeddedObservationClient { .map_err(|err| { ObservationProxyError::WorkerNotFound(format!( "embedded Worker '{}' observation subscription is unavailable: {err}", - source.worker_id + source.worker.worker_id )) })?; let mut queued = VecDeque::new(); @@ -365,12 +344,11 @@ impl EmbeddedObservationClient { .map_err(|err| { ObservationProxyError::WorkerNotFound(format!( "embedded Worker '{}' snapshot is unavailable: {err}", - source.worker_id + source.worker.worker_id )) })?; queued.push_back(RuntimeObservationUpstreamEvent { - runtime_id: source.runtime_id.clone(), - worker_id: source.worker_id.clone(), + worker: source.worker.clone(), runtime_event_id: "snapshot".to_string(), payload: snapshot, }); @@ -380,19 +358,14 @@ impl EmbeddedObservationClient { .map_err(|err| { ObservationProxyError::RuntimeUnavailable(format!( "embedded Worker '{}' observation cursor is unavailable: {err}", - source.worker_id + source.worker.worker_id )) })? { - queued.push_back(Self::map_event( - &source.runtime_id, - &source.worker_id, - event, - )); + queued.push_back(Self::map_event(&source.worker, event)); } Ok(Self { - runtime_id: source.runtime_id.clone(), - worker_id: source.worker_id.clone(), + worker: source.worker.clone(), worker_ref: source.worker_ref.clone(), cursor, receiver, @@ -418,7 +391,7 @@ impl EmbeddedObservationClient { "embedded runtime emitted a malformed cursor".into(), ) })?; - return Ok(Self::map_event(&self.runtime_id, &self.worker_id, event)); + return Ok(Self::map_event(&self.worker, event)); } Ok(_) => continue, Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => { @@ -436,13 +409,11 @@ impl EmbeddedObservationClient { } fn map_event( - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, event: WorkerObservationEvent, ) -> RuntimeObservationUpstreamEvent { RuntimeObservationUpstreamEvent { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: worker.clone(), runtime_event_id: event.cursor.clone(), payload: event.payload, } @@ -455,8 +426,7 @@ mod tests { fn sensitive_source() -> RuntimeObservationSourceConfig { RuntimeObservationSourceConfig { - runtime_id: "remote-runtime".to_string(), - worker_id: "worker-1".to_string(), + worker: RuntimeWorkerRef::new("remote-runtime", "worker-1"), endpoint: "wss://remote.example.invalid/private/workers/worker-1/protocol/ws" .to_string(), bearer_token: Some("top-secret-bearer-token".to_string()), diff --git a/crates/workspace-server/src/resource_broker.rs b/crates/workspace-server/src/resource_broker.rs index 94ad6b93..2549b52e 100644 --- a/crates/workspace-server/src/resource_broker.rs +++ b/crates/workspace-server/src/resource_broker.rs @@ -3,7 +3,7 @@ use chrono::{Duration, Utc}; use std::collections::HashMap; use std::sync::{Arc, Mutex}; use uuid::Uuid; -use worker_runtime::identity::WorkerId; +use worker_runtime::identity::RuntimeWorkerRef; use worker_runtime::profile_archive::ProfileSourceArchive; use worker_runtime::resource::{ BackendResourceClient, BackendResourceError, BackendResourceFetchRequest, @@ -17,10 +17,17 @@ pub struct BackendResourceBroker { resources: Arc>>, } +#[derive(Clone, Copy, Debug)] +pub enum BackendResourceTarget<'a> { + Workspace, + Runtime(&'a str), + Worker(&'a RuntimeWorkerRef), +} + #[derive(Clone)] struct StoredResource { runtime_id: Option, - worker_id: Option, + worker: Option, handle: BackendResourceHandle, archive: ProfileSourceArchive, } @@ -29,11 +36,17 @@ impl BackendResourceBroker { pub fn issue_profile_source_archive_handle( &self, workspace_id: impl Into, - runtime_id: Option<&str>, - worker_id: Option<&WorkerId>, + target: BackendResourceTarget<'_>, archive: ProfileSourceArchive, ) -> BackendResourceHandle { let workspace_id = workspace_id.into(); + let (runtime_id, worker) = match target { + BackendResourceTarget::Workspace => (None, None), + BackendResourceTarget::Runtime(runtime_id) => (Some(runtime_id.to_string()), None), + BackendResourceTarget::Worker(worker) => { + (Some(worker.runtime_id.clone()), Some(worker.clone())) + } + }; let nonce = Uuid::now_v7().to_string(); let audit_correlation_id = format!("resource-fetch-{nonce}"); let expires_at = Utc::now() + Duration::minutes(15); @@ -41,8 +54,8 @@ impl BackendResourceBroker { kind: BackendResourceKind::ProfileSourceArchive, workspace_id: workspace_id.clone(), scope_id: Some("workspace-profile-source".to_string()), - runtime_id: runtime_id.map(|id| id.to_string()), - worker_id: worker_id.map(|id| id.to_string()), + runtime_id: runtime_id.clone(), + worker_id: worker.as_ref().map(|worker| worker.worker_id.clone()), resource_id: archive.reference.id.clone(), digest: archive.reference.digest.clone(), operation: BackendResourceOperation::FetchArchive, @@ -57,8 +70,8 @@ impl BackendResourceBroker { profile_source_graph: Some(archive.reference.source_graph.clone()), }; let stored = StoredResource { - runtime_id: runtime_id.map(|id| id.to_string()), - worker_id: worker_id.map(|id| id.to_string()), + runtime_id, + worker, handle: handle.clone(), archive, }; @@ -117,8 +130,10 @@ impl BackendResourceBroker { }); } } - if let Some(expected_worker_id) = stored.worker_id.as_deref() { - if Some(expected_worker_id) != request.worker_id.as_deref() { + if let Some(expected_worker) = stored.worker.as_ref() { + if expected_worker.runtime_id != request.runtime_id + || Some(expected_worker.worker_id.as_str()) != request.worker_id.as_deref() + { return Err(BackendResourceError::Unauthorized { message: "worker id does not match resource handle".to_string(), }); @@ -167,7 +182,6 @@ fn verify_handle_shape(handle: &BackendResourceHandle) -> Result<(), BackendReso mod tests { use super::*; use std::collections::BTreeMap; - use worker_runtime::identity::WorkerId; use worker_runtime::profile_archive::{ ProfileSourceArchive, ProfileSourceArchiveRef, ProfileSourceGraphSummary, sha256_hex, }; @@ -215,13 +229,13 @@ mod tests { fn request( handle: BackendResourceHandle, runtime_id: &str, - worker_id: Option<&WorkerId>, + worker_id: Option<&str>, ) -> BackendResourceFetchRequest { BackendResourceFetchRequest { audit_correlation_id: handle.audit_correlation_id.clone(), handle, runtime_id: runtime_id.to_string(), - worker_id: worker_id.map(|id| id.to_string()), + worker_id: worker_id.map(str::to_string), } } @@ -231,8 +245,7 @@ mod tests { let runtime_id = "runtime-test"; let handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - None, + BackendResourceTarget::Runtime(runtime_id), archive(), ); let response = broker @@ -253,8 +266,7 @@ mod tests { let runtime_a = "runtime-a"; let handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_a), - None, + BackendResourceTarget::Runtime(runtime_a), archive(), ); let err = broker @@ -267,16 +279,15 @@ mod tests { fn broker_rejects_worker_mismatch() { let broker = BackendResourceBroker::default(); let runtime_id = "runtime-test"; - let worker_a = WorkerId::new(1); - let worker_b = WorkerId::new(2); + let worker_a = RuntimeWorkerRef::new(runtime_id, "1"); + let worker_b = RuntimeWorkerRef::new(runtime_id, "2"); let handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - Some(&worker_a), + BackendResourceTarget::Worker(&worker_a), archive(), ); let err = broker - .fetch_profile_source_archive(request(handle, &runtime_id, Some(&worker_b))) + .fetch_profile_source_archive(request(handle, runtime_id, Some(&worker_b.worker_id))) .unwrap_err(); assert!(matches!(err, BackendResourceError::Unauthorized { .. })); } @@ -287,8 +298,7 @@ mod tests { let runtime_id = "runtime-test"; let handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - None, + BackendResourceTarget::Runtime(runtime_id), archive(), ); broker @@ -313,8 +323,7 @@ mod tests { let runtime_id = "runtime-test"; let mut handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - None, + BackendResourceTarget::Runtime(runtime_id), archive(), ); handle.scope_id = Some("tampered-scope".to_string()); @@ -331,8 +340,7 @@ mod tests { let archive = archive_with_len((DEFAULT_PROFILE_SOURCE_ARCHIVE_MAX_BYTES + 1) as usize); let mut handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - None, + BackendResourceTarget::Runtime(runtime_id), archive, ); handle.max_bytes = DEFAULT_PROFILE_SOURCE_ARCHIVE_MAX_BYTES + 1024; diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 1bc35343..542522c8 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -110,7 +110,7 @@ use worker_runtime::http_server::{ RuntimeHttpConfigBundleAvailabilityResponse, RuntimeHttpConfigBundlesResponse, RuntimeHttpSummaryResponse, RuntimeHttpWorkerResponse, RuntimeHttpWorkersResponse, }; -use worker_runtime::identity::WorkerId; +use worker_runtime::identity::RuntimeWorkerRef; use worker_runtime::interaction::{ WorkerInput as EmbeddedWorkerInput, WorkerInputKind as EmbeddedWorkerInputKind, }; @@ -259,8 +259,8 @@ pub struct WorkspaceApi { observation_proxy: BackendObservationProxy, runtime_subscription_broker: RuntimeSubscriptionBroker, resource_broker: BackendResourceBroker, - workdir_sessions: Arc>>, - workdir_session_locks: Arc>>>>, + workdir_sessions: Arc>>, + workdir_session_locks: Arc>>>>, } impl WorkspaceApi { @@ -437,14 +437,14 @@ impl WorkspaceApi { } return Ok(result); }; - let replacement = match self.runtime.replace_worker_workspace_api( - runtime_id, - &worker.worker_id, - workspace_api, - ) { + let worker_ref = worker.worker.clone(); + let replacement = match self + .runtime + .replace_worker_workspace_api(&worker_ref, workspace_api) + { Ok(replacement) => replacement, Err(error) => { - let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); + let _ = self.runtime.delete_worker(&worker_ref); if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { let _ = self.store.release_worker_workdir_attachment_reservation( &self.config.workspace_id, @@ -456,7 +456,7 @@ impl WorkspaceApi { } }; if replacement.state != WorkerOperationState::Accepted { - let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); + let _ = self.runtime.delete_worker(&worker_ref); if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { let _ = self.store.release_worker_workdir_attachment_reservation( &self.config.workspace_id, @@ -478,22 +478,18 @@ impl WorkspaceApi { .into()); } if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - let runtime_worker_id = match parse_runtime_worker_id_for_registry(&worker.worker_id) { - Ok(worker_id) => worker_id, - Err(error) => { - let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - ); - return Err(error); - } - }; + if let Err(error) = parse_runtime_worker_id_for_registry(&worker.worker.worker_id) { + let _ = self.runtime.delete_worker(&worker_ref); + let _ = self.store.release_worker_workdir_attachment_reservation( + &self.config.workspace_id, + workdir_id, + reservation_id, + ); + return Err(error); + } let attachment = WorkerWorkdirLinkRecord { workspace_id: self.config.workspace_id.clone(), - runtime_id: runtime_id.to_string(), - runtime_worker_id, + worker: worker_ref.clone(), workdir_id: workdir_id.clone(), role: "attachment".to_string(), linked_at: now_registry_timestamp(), @@ -503,7 +499,7 @@ impl WorkspaceApi { .store .finalize_reserved_worker_workdir_attachment(&attachment, reservation_id) { - let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); + let _ = self.runtime.delete_worker(&worker_ref); let _ = self.store.release_worker_workdir_attachment_reservation( &self.config.workspace_id, workdir_id, @@ -517,12 +513,11 @@ impl WorkspaceApi { fn restore_workspace_worker( &self, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> ApiResult { let binding = self .runtime - .replace_worker_workspace_api(runtime_id, worker_id, self.workspace_api_ref(runtime_id)) + .replace_worker_workspace_api(worker, self.workspace_api_ref(&worker.runtime_id)) .map_err(|error| error.into_error())?; if binding.state != WorkerOperationState::Accepted { return Ok(WorkerRestoreResult { @@ -533,7 +528,7 @@ impl WorkspaceApi { } Ok(self .runtime - .restore_worker(runtime_id, worker_id) + .restore_worker(worker) .map_err(|error| error.into_error())?) } @@ -1132,8 +1127,8 @@ enum RuntimeWorkersStatusFilter { #[derive(Debug, Serialize, Deserialize)] pub struct WorkerRestoreResponse { pub workspace_id: String, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker_ref: RuntimeWorkerRef, pub result: WorkerRestoreResult, } @@ -1300,8 +1295,8 @@ pub struct RuntimeCleanupExecutionResponse { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerRetentionResponse { pub workspace_id: String, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker_ref: RuntimeWorkerRef, pub pinned: bool, pub retention_state: String, } @@ -1460,8 +1455,8 @@ pub struct BrowserCreateWorkerRequest { #[derive(Debug, Serialize, Deserialize)] pub struct BrowserCreateWorkerResponse { pub workspace_id: String, - pub runtime_id: String, - pub worker_id: String, + #[serde(flatten)] + pub worker_ref: RuntimeWorkerRef, pub console_href: String, pub worker: WorkerSummary, pub diagnostics: Vec, @@ -1643,8 +1638,8 @@ struct ScopedConfigBundlePath { #[derive(Debug, Deserialize)] struct ScopedRuntimeWorkerPath { workspace_id: String, - runtime_id: String, - worker_id: String, + #[serde(flatten)] + worker: RuntimeWorkerRef, } fn validate_workspace_scope(api: &WorkspaceApi, workspace_id: &str) -> ApiResult<()> { @@ -1864,8 +1859,8 @@ struct TicketWorkerAssignmentMutationResponse { #[serde(deny_unknown_fields)] struct SetTicketWorkerAssignmentRequest { operation_id: String, - runtime_id: String, - worker_id: String, + #[serde(flatten)] + worker: RuntimeWorkerRef, expected_assignment_id: Option, assigned_by: Option, } @@ -1887,11 +1882,9 @@ async fn scoped_get_ticket_worker_assignment( let assignment = api .store .get_current_ticket_worker_assignment(&path.workspace_id, &ticket.id)?; - let worker = assignment.as_ref().and_then(|assignment| { - api.runtime - .worker(&assignment.runtime_id, &assignment.worker_id) - .ok() - }); + let worker = assignment + .as_ref() + .and_then(|assignment| api.runtime.worker(&assignment.worker).ok()); Ok(Json(TicketWorkerAssignmentResponse { workspace_id: path.workspace_id, ticket_id: ticket.id, @@ -1925,8 +1918,8 @@ async fn set_ticket_worker_assignment( validate_workspace_scope(&api, &path.workspace_id)?; let ticket = api.authority.ticket(&path.id)?; let operation_id = require_ticket_assignment_value("operation_id", request.operation_id)?; - let runtime_id = require_ticket_assignment_value("runtime_id", request.runtime_id)?; - let worker_id = require_ticket_assignment_value("worker_id", request.worker_id)?; + let runtime_id = require_ticket_assignment_value("runtime_id", request.worker.runtime_id)?; + let worker_id = require_ticket_assignment_value("worker_id", request.worker.worker_id)?; let expected_assignment_id = request .expected_assignment_id .map(|value| require_ticket_assignment_value("expected_assignment_id", value)) @@ -1936,17 +1929,17 @@ async fn set_ticket_worker_assignment( .map(|value| require_ticket_assignment_value("assigned_by", value)) .transpose()? .unwrap_or_else(|| "workspace-api".to_string()); + let requested_worker = RuntimeWorkerRef::new(runtime_id, worker_id); let worker = api .runtime - .worker(&runtime_id, &worker_id) + .worker(&requested_worker) .map_err(|err| err.into_error())?; let assigned_at = Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true); let record = TicketWorkerAssignmentRecord { workspace_id: path.workspace_id.clone(), ticket_id: ticket.id.clone(), assignment_id: new_id("tasg"), - runtime_id: worker.runtime_id, - worker_id: worker.worker_id, + worker: worker.worker.clone(), assigned_by, assigned_at, }; @@ -2017,8 +2010,7 @@ fn assign_ticket_worker_from_lifecycle( workspace_id: api.config.workspace_id.clone(), ticket_id: ticket.id, assignment_id: new_id("tasg"), - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), + worker: RuntimeWorkerRef::new(runtime_id, worker_id), assigned_by: "worker-lifecycle".to_string(), assigned_at, }; @@ -2054,12 +2046,12 @@ fn existing_lifecycle_assignment_worker( assignment.operation_id ))); } - let Some(worker_id) = operation.worker_id else { + let Some(worker_ref) = operation.worker else { return Ok(None); }; let worker = api .runtime - .worker(runtime_id, &worker_id) + .worker(&worker_ref) .map_err(|error| error.into_error())?; if operation.assignment_id.is_none() && worker.state == "stopped" { return Ok(None); @@ -2358,7 +2350,7 @@ async fn execute_worker_ticket_rest_operation( operation_kind, &previous_state, ticket.meta.workflow_state.as_str(), - Some((source.runtime_id, source.worker_id)), + Some(source), ); } Ok(result) @@ -2854,11 +2846,7 @@ async fn scoped_ticket_doctor( }) } -#[derive(Debug, Clone, PartialEq, Eq)] -struct WorkerMutationSource { - runtime_id: String, - worker_id: String, -} +type WorkerMutationSource = RuntimeWorkerRef; fn ticket_mutation_target(operation: &TicketBackendOperation) -> Option<&TicketIdOrSlug> { match operation { @@ -2941,8 +2929,7 @@ fn ticket_operation_initial_state(operation: &TicketBackendOperation) -> String #[derive(Debug, Clone)] struct WorkerTicketSourceContext { - runtime_id: String, - worker_id: String, + worker: RuntimeWorkerRef, actor_role: String, assignment_id: Option, } @@ -2950,8 +2937,14 @@ struct WorkerTicketSourceContext { impl WorkerTicketSourceContext { fn attributes(&self, operation_kind: &str) -> BTreeMap { let mut attributes = BTreeMap::from([ - ("source_runtime_id".to_string(), self.runtime_id.clone()), - ("source_worker_id".to_string(), self.worker_id.clone()), + ( + "source_runtime_id".to_string(), + self.worker.runtime_id.clone(), + ), + ( + "source_worker_id".to_string(), + self.worker.worker_id.clone(), + ), ("source_actor_role".to_string(), self.actor_role.clone()), ( "source_operation_kind".to_string(), @@ -2988,20 +2981,18 @@ fn worker_ticket_source_context( .flatten() }); let orchestrator = find_workspace_orchestrator(api); - let is_current_assignment = assignment.as_ref().is_some_and(|assignment| { - assignment.runtime_id == source.runtime_id && assignment.worker_id == source.worker_id - }); - let is_orchestrator = orchestrator.as_ref().is_some_and(|worker| { - worker.runtime_id == source.runtime_id && worker.worker_id == source.worker_id - }); + let is_current_assignment = assignment + .as_ref() + .is_some_and(|assignment| &assignment.worker == source); + let is_orchestrator = orchestrator + .as_ref() + .is_some_and(|worker| worker.worker == *source); let actor_role = worker_source_actor_role(is_current_assignment, is_orchestrator); WorkerTicketSourceContext { - runtime_id: source.runtime_id.clone(), - worker_id: source.worker_id.clone(), + worker: source.clone(), actor_role: actor_role.to_string(), assignment_id: assignment.and_then(|assignment| { - (assignment.runtime_id == source.runtime_id && assignment.worker_id == source.worker_id) - .then_some(assignment.assignment_id) + (assignment.worker == *source).then_some(assignment.assignment_id) }), } } @@ -3015,7 +3006,7 @@ fn notify_ticket_recipients( source_operation_kind: &str, previous_state: &str, current_state: &str, - source: Option<(String, String)>, + source: Option, ) { let mut recipients = Vec::new(); if let Some(assignment) = api @@ -3024,33 +3015,32 @@ fn notify_ticket_recipients( .ok() .flatten() { - recipients.push((assignment.runtime_id, assignment.worker_id)); + recipients.push(assignment.worker.clone()); } if (matches!(previous_state, "queued" | "inprogress") || matches!(current_state, "queued" | "inprogress")) && let Some(orchestrator) = find_workspace_orchestrator(api) { - recipients.push((orchestrator.runtime_id, orchestrator.worker_id)); + recipients.push(orchestrator.worker.clone()); } recipients.sort(); recipients.dedup(); let source_fields = source .as_ref() - .map(|(runtime_id, worker_id)| { - format!(" source_runtime_id={runtime_id} source_worker_id={worker_id}") + .map(|source| { + format!( + " source_runtime_id={} source_worker_id={}", + source.runtime_id, source.worker_id + ) }) .unwrap_or_default(); - for (runtime_id, worker_id) in recipients { - if source - .as_ref() - .is_some_and(|source| source.0 == runtime_id && source.1 == worker_id) - { + for recipient in recipients { + if source.as_ref().is_some_and(|source| source == &recipient) { continue; } let _ = api.runtime.send_input( - &runtime_id, - &worker_id, + &recipient, WorkerInputRequest { kind: WorkerInputKind::Notify, content: format!( @@ -3079,51 +3069,41 @@ fn authenticate_worker_mutation_source( .ok_or_else(|| { Error::WorkerSourceIdentity("missing Runtime-bound Worker id".to_string()) })?; - api.runtime.worker(runtime_id, worker_id).map_err(|_| { + let worker = RuntimeWorkerRef::new(runtime_id, worker_id); + api.runtime.worker(&worker).map_err(|_| { Error::WorkerSourceIdentity("Runtime-bound Worker identity does not exist".to_string()) })?; - Ok(WorkerMutationSource { - runtime_id: runtime_id.to_string(), - worker_id: worker_id.to_string(), - }) + Ok(worker) } fn current_worker_identity( api: &WorkspaceApi, workspace_id: &str, headers: &HeaderMap, -) -> Result<(String, u64)> { - let source = authenticate_worker_mutation_source(api, workspace_id, headers)?; - let worker_id = source.worker_id.parse::().map_err(|_| { - Error::WorkerSourceIdentity(format!( - "Runtime-bound Worker id must be numeric, got `{}`", - source.worker_id - )) - })?; - Ok((source.runtime_id, worker_id)) +) -> Result { + authenticate_worker_mutation_source(api, workspace_id, headers) } fn current_worker_active_attachment( api: &WorkspaceApi, - runtime_id: &str, - worker_id: u64, + worker: &RuntimeWorkerRef, ) -> ApiResult { if let Some(link) = api .store - .list_worker_workdir_links(&api.config.workspace_id, runtime_id, worker_id)? + .list_worker_workdir_links(&api.config.workspace_id, worker)? .into_iter() .next() { return Ok(link); } - if api.store.worker_workdir_link_history_exists( - &api.config.workspace_id, - runtime_id, - worker_id, - )? { + if api + .store + .worker_workdir_link_history_exists(&api.config.workspace_id, worker)? + { return Err(Error::WorkdirAttachmentConflict(format!( - "Worker {runtime_id}:{worker_id} has no active Workdir attachment" + "Worker {}:{} has no active Workdir attachment", + worker.runtime_id, worker.worker_id )) .into()); } @@ -3132,15 +3112,15 @@ fn current_worker_active_attachment( // outer spawn handler has projected that binding into the Backend registry. Import the same // binding transactionally on the first identity-bound operation so initial input cannot race // attachment authority. - let worker = api + let observed_worker = api .runtime - .worker(runtime_id, &worker_id.to_string()) + .worker(worker) .map_err(|error| error.into_error())?; - if worker.working_directory.is_some() { - sync_worker_observation(api, &worker)?; + if observed_worker.working_directory.is_some() { + sync_worker_observation(api, &observed_worker)?; if let Some(link) = api .store - .list_worker_workdir_links(&api.config.workspace_id, runtime_id, worker_id)? + .list_worker_workdir_links(&api.config.workspace_id, worker)? .into_iter() .next() { @@ -3148,43 +3128,41 @@ fn current_worker_active_attachment( } } Err(Error::WorkdirAttachmentConflict(format!( - "Worker {runtime_id}:{worker_id} has no active Workdir attachment" + "Worker {}:{} has no active Workdir attachment", + worker.runtime_id, worker.worker_id )) .into()) } fn current_worker_session_lock( api: &WorkspaceApi, - runtime_id: &str, - worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Arc> { api.workdir_session_locks .lock() .expect("Workdir session lock registry poisoned") - .entry((runtime_id.to_string(), worker_id)) + .entry(worker.clone()) .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(()))) .clone() } -fn runtime_local_owner_worker_id( - caller_runtime_id: &str, +fn runtime_local_owner_worker_id<'a>( + caller: &'a RuntimeWorkerRef, target_runtime_id: &str, - caller_worker_id: u64, -) -> Option { - (caller_runtime_id == target_runtime_id).then(|| caller_worker_id.to_string()) +) -> Option<&'a str> { + (caller.runtime_id == target_runtime_id).then_some(caller.worker_id.as_str()) } async fn open_current_worker_workdir_session_locked( api: &WorkspaceApi, - runtime_id: &str, - worker_id: u64, + worker: &RuntimeWorkerRef, link: &WorkerWorkdirLinkRecord, ) -> Result { if let Some(session) = api .workdir_sessions .lock() .expect("Workdir session registry lock poisoned") - .get(&(runtime_id.to_string(), worker_id)) + .get(worker) .cloned() { return Ok(session); @@ -3195,24 +3173,20 @@ async fn open_current_worker_workdir_session_locked( .into_iter() .find(|workdir| workdir.workdir_id == link.workdir_id) .ok_or_else(|| Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), + runtime_id: worker.runtime_id.clone(), code: "working_directory_not_found".to_string(), message: format!( "attached Workdir {} is not registered in this Workspace", link.workdir_id ), })?; - let owner_worker_id = runtime_local_owner_worker_id(runtime_id, &workdir.runtime_id, worker_id); + let owner_worker_id = runtime_local_owner_worker_id(worker, &workdir.runtime_id); let session = api .runtime - .open_workdir_session( - &workdir.runtime_id, - &workdir.workdir_id, - owner_worker_id.as_deref(), - ) + .open_workdir_session(&workdir.runtime_id, &workdir.workdir_id, owner_worker_id) .await .map_err(|error| error.into_error())?; - let key = (runtime_id.to_string(), worker_id); + let key = worker.clone(); let (selected, unused) = { let mut sessions = api .workdir_sessions @@ -3233,7 +3207,7 @@ async fn open_current_worker_workdir_session_locked( .close() .await .map_err(|error| Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), + runtime_id: worker.runtime_id.clone(), code: "duplicate_workdir_session_close_failed".to_string(), message: error.to_string(), })?; @@ -3243,10 +3217,9 @@ async fn open_current_worker_workdir_session_locked( async fn close_current_worker_session_locked( api: &WorkspaceApi, - runtime_id: &str, - worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result<()> { - let key = (runtime_id.to_string(), worker_id); + let key = worker.clone(); let session = api .workdir_sessions .lock() @@ -3258,7 +3231,7 @@ async fn close_current_worker_session_locked( .close() .await .map_err(|error| Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), + runtime_id: worker.runtime_id.clone(), code: "workdir_session_close_failed".to_string(), message: error.to_string(), })?; @@ -3277,7 +3250,7 @@ async fn scoped_attach_current_worker_workdir( Json(request): Json, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - let (runtime_id, worker_id) = current_worker_identity(&api, &path.workspace_id, &headers)?; + let worker = current_worker_identity(&api, &path.workspace_id, &headers)?; let workdir_id = request.workdir_id.trim(); if workdir_id.is_empty() || workdir_id.chars().any(char::is_control) { return Err(Error::InvalidRecordId(request.workdir_id).into()); @@ -3289,30 +3262,26 @@ async fn scoped_attach_current_worker_workdir( .any(|workdir| workdir.workdir_id == workdir_id) { return Err(Error::RuntimeOperationFailed { - runtime_id: runtime_id.clone(), + runtime_id: worker.runtime_id.clone(), code: "working_directory_not_found".to_string(), message: format!("unknown Workdir `{workdir_id}`"), } .into()); } - let session_lock = current_worker_session_lock(&api, &runtime_id, worker_id); + let session_lock = current_worker_session_lock(&api, &worker); let _session_guard = session_lock.lock().await; let link = api.store.attach_worker_workdir(&WorkerWorkdirLinkRecord { workspace_id: api.config.workspace_id.clone(), - runtime_id: runtime_id.clone(), - runtime_worker_id: worker_id, + worker: worker.clone(), workdir_id: workdir_id.to_string(), role: "attachment".to_string(), linked_at: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), unlinked_at: None, })?; - if let Err(error) = - open_current_worker_workdir_session_locked(&api, &runtime_id, worker_id, &link).await - { + if let Err(error) = open_current_worker_workdir_session_locked(&api, &worker, &link).await { let _ = api.store.detach_worker_workdir( &api.config.workspace_id, - &runtime_id, - worker_id, + &worker, Some(workdir_id), &Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), ); @@ -3331,15 +3300,14 @@ async fn scoped_detach_current_worker_workdir( headers: HeaderMap, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - let (runtime_id, worker_id) = current_worker_identity(&api, &path.workspace_id, &headers)?; - let session_lock = current_worker_session_lock(&api, &runtime_id, worker_id); + let worker = current_worker_identity(&api, &path.workspace_id, &headers)?; + let session_lock = current_worker_session_lock(&api, &worker); let _session_guard = session_lock.lock().await; - let link = current_worker_active_attachment(&api, &runtime_id, worker_id)?; - close_current_worker_session_locked(&api, &runtime_id, worker_id).await?; + let link = current_worker_active_attachment(&api, &worker)?; + close_current_worker_session_locked(&api, &worker).await?; api.store.detach_worker_workdir( &api.config.workspace_id, - &runtime_id, - worker_id, + &worker, Some(&link.workdir_id), &Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), )?; @@ -3357,16 +3325,15 @@ async fn scoped_execute_current_worker_workdir_operation( Json(operation): Json, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - let (runtime_id, worker_id) = current_worker_identity(&api, &path.workspace_id, &headers)?; - let session_lock = current_worker_session_lock(&api, &runtime_id, worker_id); + let worker = current_worker_identity(&api, &path.workspace_id, &headers)?; + let session_lock = current_worker_session_lock(&api, &worker); let _session_guard = session_lock.lock().await; - let link = current_worker_active_attachment(&api, &runtime_id, worker_id)?; - let session = - open_current_worker_workdir_session_locked(&api, &runtime_id, worker_id, &link).await?; + let link = current_worker_active_attachment(&api, &worker)?; + let session = open_current_worker_workdir_session_locked(&api, &worker, &link).await?; let result = execute_workdir_session_operation(&session, operation) .await .map_err(|error| Error::RuntimeOperationFailed { - runtime_id, + runtime_id: worker.runtime_id.clone(), code: "workdir_session_operation_failed".to_string(), message: error.to_string(), })?; @@ -3497,7 +3464,8 @@ fn maybe_dispatch_orchestrator_turn_end( let Some(orchestrator) = find_workspace_orchestrator(api) else { return; }; - if orchestrator.runtime_id == EMBEDDED_WORKER_RUNTIME_ID && orchestrator.worker_id == worker_id + if orchestrator.worker.runtime_id == EMBEDDED_WORKER_RUNTIME_ID + && orchestrator.worker.worker_id == worker_id { dispatch_orchestrator_queue_attention(api); } @@ -3573,8 +3541,7 @@ fn dispatch_orchestrator_queue_attention(api: &WorkspaceApi) { let accepted = api .runtime .send_input( - &orchestrator.runtime_id, - &orchestrator.worker_id, + &orchestrator.worker, WorkerInputRequest { kind: WorkerInputKind::Notify, content, @@ -3792,7 +3759,7 @@ fn start_memory_staging_consolidation( status: "started".to_string(), summary: format!( "Started Memory consolidater '{}' for {candidate_count} staging candidate(s).", - worker.worker_id + worker.worker.worker_id ), candidate_count, total_bytes, @@ -3818,13 +3785,13 @@ fn try_reuse_memory_consolidation_worker( if consolidaters.is_empty() { return Ok(None); } - consolidaters.sort_by(|a, b| a.worker_id.cmp(&b.worker_id)); + consolidaters.sort_by(|a, b| a.worker.worker_id.cmp(&b.worker.worker_id)); if let Some(worker) = consolidaters.iter().find(|worker| worker.state != "idle") { return Ok(Some(MemoryConsolidationOutput { status: "skipped_existing_not_idle".to_string(), summary: format!( "Existing Memory consolidater '{}' is '{}', not confirmed idle.", - worker.worker_id, worker.state + worker.worker.worker_id, worker.state ), candidate_count, total_bytes, @@ -3837,8 +3804,7 @@ fn try_reuse_memory_consolidation_worker( let input = api .runtime .send_input( - &worker.runtime_id, - &worker.worker_id, + &worker.worker, WorkerInputRequest { kind: WorkerInputKind::User, content: input_content.to_string(), @@ -3851,7 +3817,7 @@ fn try_reuse_memory_consolidation_worker( status: "skipped_existing_input_rejected".to_string(), summary: format!( "Existing idle Memory consolidater '{}' rejected the new consolidation input.", - worker.worker_id + worker.worker.worker_id ), candidate_count, total_bytes, @@ -3861,7 +3827,7 @@ fn try_reuse_memory_consolidation_worker( status: "reused".to_string(), summary: format!( "Reused Memory consolidater '{}' for {candidate_count} staging candidate(s).", - worker.worker_id + worker.worker.worker_id ), candidate_count, total_bytes, @@ -4162,12 +4128,12 @@ async fn scoped_start_workspace_orchestrator( } let restored = api .runtime - .restore_worker(&existing.runtime_id, &existing.worker_id) + .restore_worker(&existing.worker) .map_err(|error| error.into_error())?; if restored.state != WorkerOperationState::Accepted { return Err(ApiError::with_diagnostics( Error::RuntimeOperationFailed { - runtime_id: existing.runtime_id, + runtime_id: existing.worker.runtime_id.clone(), code: "workspace_orchestrator_restore_rejected".to_string(), message: "Runtime rejected Workspace Orchestrator restore".to_string(), }, @@ -4521,28 +4487,21 @@ async fn set_worker_retention( runtime_worker_id: String, pinned: bool, ) -> ApiResult> { - let runtime_worker_registry_id = parse_runtime_worker_id_for_registry(&runtime_worker_id)?; + parse_runtime_worker_id_for_registry(&runtime_worker_id)?; + let worker_ref = RuntimeWorkerRef::new(runtime_id.clone(), runtime_worker_id.clone()); if api .store - .get_worker_registry( - &api.config.workspace_id, - runtime_id.as_str(), - runtime_worker_registry_id, - )? + .get_worker_registry(&api.config.workspace_id, &worker_ref)? .is_none() { - if let Ok(worker) = api - .runtime - .worker(runtime_id.as_str(), runtime_worker_id.as_str()) - { + if let Ok(worker) = api.runtime.worker(&worker_ref) { let _ = sync_worker_observation(&api, &worker); } } let retention_state = if pinned { "pinned" } else { "normal" }; let changed = api.store.update_worker_retention( &api.config.workspace_id, - runtime_id.as_str(), - runtime_worker_registry_id, + &worker_ref, retention_state, now_registry_timestamp().as_str(), )?; @@ -4555,8 +4514,7 @@ async fn set_worker_retention( } Ok(Json(WorkerRetentionResponse { workspace_id: api.config.workspace_id, - runtime_id, - worker_id: runtime_worker_id, + worker_ref, pinned, retention_state: retention_state.to_string(), })) @@ -4567,15 +4525,11 @@ fn build_runtime_cleanup_plan( runtime_id: &str, ) -> ApiResult { let workers = workers_response(api.clone())?; - let live_running_worker_ids: HashSet<(String, u64)> = workers + let live_running_worker_ids: HashSet = workers .items .iter() .filter(|worker| worker.state == "running") - .filter_map(|worker| { - parse_runtime_worker_id_for_registry(worker.worker_id.as_str()) - .ok() - .map(|worker_id| (worker.runtime_id.clone(), worker_id)) - }) + .map(|worker| worker.worker.clone()) .collect(); let (workdir_summaries, mut diagnostics) = match runtime_working_directory_summaries(api, runtime_id) { @@ -4600,12 +4554,7 @@ fn build_runtime_cleanup_plan( .list_worker_registry(&api.config.workspace_id, 500)?; let worker_by_id: HashMap<_, _> = worker_records .iter() - .map(|record| { - ( - (record.runtime_id.clone(), record.runtime_worker_id.clone()), - record.clone(), - ) - }) + .map(|record| (record.worker.clone(), record.clone())) .collect(); let observed_workdirs: HashMap<_, _> = workdir_summaries .into_iter() @@ -4615,15 +4564,12 @@ fn build_runtime_cleanup_plan( let mut worker_candidates = Vec::new(); for record in worker_records .iter() - .filter(|record| record.runtime_id == runtime_id) + .filter(|record| record.worker.runtime_id == runtime_id) { - let links = api.store.list_worker_workdir_links( - &api.config.workspace_id, - record.runtime_id.as_str(), - record.runtime_worker_id, - )?; - let is_running = live_running_worker_ids - .contains(&(record.runtime_id.clone(), record.runtime_worker_id.clone())); + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, &record.worker)?; + let is_running = live_running_worker_ids.contains(&record.worker); let pinned = record.retention_state == "pinned"; let blocking_reason = if pinned { Some("worker is pinned".to_string()) @@ -4635,13 +4581,13 @@ fn build_runtime_cleanup_plan( worker_candidates.push(CleanupWorkerCandidate { target_id: format!( "worker:{}:{}", - encode_path_segment(record.runtime_id.as_str()), - encode_path_segment(&record.runtime_worker_id.to_string()) + encode_path_segment(record.worker.runtime_id.as_str()), + encode_path_segment(&record.worker.worker_id) ), action: CleanupTargetKind::WorkerDelete, - worker_id: record.runtime_worker_id.to_string(), - runtime_worker_id: record.runtime_worker_id.to_string(), - runtime_id: record.runtime_id.clone(), + worker_id: record.worker.worker_id.clone(), + runtime_worker_id: record.worker.worker_id.clone(), + runtime_id: record.worker.runtime_id.clone(), reason: if blocking_reason.is_some() { "Worker cannot be deleted until blocking conditions are cleared".to_string() } else { @@ -4666,21 +4612,16 @@ fn build_runtime_cleanup_plan( .list_workdir_worker_links(&api.config.workspace_id, record.workdir_id.as_str())?; let linked_workers = links .iter() - .filter_map(|link| { - worker_by_id.get(&(link.runtime_id.clone(), link.runtime_worker_id.clone())) - }) + .filter_map(|link| worker_by_id.get(&link.worker)) .collect::>(); let linked_worker_ids = links .iter() - .map(|link| link.runtime_worker_id.to_string()) + .map(|link| link.worker.worker_id.clone()) .collect::>(); let linked_running_worker_ids = linked_workers .iter() - .filter(|worker| { - live_running_worker_ids - .contains(&(worker.runtime_id.clone(), worker.runtime_worker_id.clone())) - }) - .map(|worker| worker.runtime_worker_id.to_string()) + .filter(|worker| live_running_worker_ids.contains(&worker.worker)) + .map(|worker| worker.worker.worker_id.clone()) .collect::>(); let pinned_linked = linked_workers .iter() @@ -4822,21 +4763,22 @@ async fn execute_runtime_cleanup( "pinned Worker/history cannot be deleted", )); } - let runtime_worker_id = parse_runtime_worker_id_for_registry(&candidate.runtime_worker_id)?; - let session_lock = current_worker_session_lock(api, runtime_id, runtime_worker_id); + parse_runtime_worker_id_for_registry(&candidate.runtime_worker_id)?; + let worker = RuntimeWorkerRef::new( + candidate.runtime_id.clone(), + candidate.runtime_worker_id.clone(), + ); + let session_lock = current_worker_session_lock(api, &worker); let _session_guard = session_lock.lock().await; - close_current_worker_session_locked(api, runtime_id, runtime_worker_id).await?; + close_current_worker_session_locked(api, &worker).await?; cleanup_runtime_worker_for_execution(api, runtime_id, candidate)?; - api.store.delete_worker_registry( - &api.config.workspace_id, - candidate.runtime_id.as_str(), - runtime_worker_id, - )?; + api.store + .delete_worker_registry(&api.config.workspace_id, &worker)?; drop(_session_guard); api.workdir_session_locks .lock() .expect("Workdir session lock registry poisoned") - .remove(&(runtime_id.to_string(), runtime_worker_id)); + .remove(&worker); results.push(RuntimeCleanupExecutionResult { target_id: candidate.target_id.clone(), action: candidate.action.clone(), @@ -4963,9 +4905,9 @@ fn cleanup_runtime_worker_for_execution( runtime_id: &str, candidate: &CleanupWorkerCandidate, ) -> ApiResult<()> { + let worker = RuntimeWorkerRef::new(runtime_id, &candidate.runtime_worker_id); match api.runtime.stop_worker( - runtime_id, - candidate.runtime_worker_id.as_str(), + &worker, WorkerLifecycleRequest { reason: Some("cleanup worker before deletion".to_string()), ticket_assignment: None, @@ -4986,10 +4928,7 @@ fn cleanup_runtime_worker_for_execution( Err(error) => return Err(error.into_error().into()), } - match api - .runtime - .delete_worker(runtime_id, candidate.runtime_worker_id.as_str()) - { + match api.runtime.delete_worker(&worker) { Ok(result) if result.deleted && result.state == WorkerOperationState::Accepted => Ok(()), Ok(result) => Err(ApiError::with_diagnostics( Error::RuntimeOperationFailed { @@ -5149,7 +5088,11 @@ async fn scoped_get_runtime_worker( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - get_runtime_worker(State(api), AxumPath((path.runtime_id, path.worker_id))).await + get_runtime_worker( + State(api), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), + ) + .await } #[derive(Debug, Default, Deserialize)] @@ -5165,8 +5108,8 @@ async fn scoped_restore_runtime_worker( ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; let workspace_id = path.workspace_id.clone(); - let runtime_id = path.runtime_id.clone(); - let worker_id = path.worker_id.clone(); + let runtime_id = path.worker.runtime_id.clone(); + let worker_id = path.worker.worker_id.clone(); let assignment_request = match ( query.ticket_id.clone(), query.assignment_operation_id.clone(), @@ -5207,18 +5150,17 @@ async fn scoped_restore_runtime_worker( &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), )?; if let Some(worker) = existing_lifecycle_assignment_worker(&api, assignment, &runtime_id)? { - if worker.worker_id != worker_id { + if worker.worker.worker_id != worker_id { return Err(Error::TicketAssignmentConflict(format!( "assignment operation {} belongs to worker {}, not {}", - assignment.operation_id, worker.worker_id, worker_id + assignment.operation_id, worker.worker.worker_id, worker_id )) .into()); } assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker_id)?; return Ok(Json(WorkerRestoreResponse { workspace_id, - runtime_id, - worker_id, + worker_ref: RuntimeWorkerRef::new(&runtime_id, &worker_id), result: crate::hosts::WorkerRestoreResult { state: WorkerOperationState::Accepted, worker: Some(worker), @@ -5243,7 +5185,7 @@ async fn scoped_pin_runtime_worker( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - set_worker_retention(api, path.runtime_id, path.worker_id, true).await + set_worker_retention(api, path.worker.runtime_id, path.worker.worker_id, true).await } async fn scoped_unpin_runtime_worker( @@ -5251,7 +5193,7 @@ async fn scoped_unpin_runtime_worker( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - set_worker_retention(api, path.runtime_id, path.worker_id, false).await + set_worker_retention(api, path.worker.runtime_id, path.worker.worker_id, false).await } async fn scoped_runtime_cleanup_plan( @@ -5281,7 +5223,7 @@ async fn scoped_send_runtime_worker_input( validate_workspace_scope(&api, &path.workspace_id)?; send_runtime_worker_input( State(api), - AxumPath((path.runtime_id, path.worker_id)), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), Json(request), ) .await @@ -5295,7 +5237,7 @@ async fn scoped_runtime_worker_completions( validate_workspace_scope(&api, &path.workspace_id)?; runtime_worker_completions( State(api), - AxumPath((path.runtime_id, path.worker_id)), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), Json(request), ) .await @@ -5309,7 +5251,7 @@ async fn scoped_stop_runtime_worker( validate_workspace_scope(&api, &path.workspace_id)?; stop_runtime_worker( State(api), - AxumPath((path.runtime_id, path.worker_id)), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), Json(request), ) .await @@ -5323,7 +5265,7 @@ async fn scoped_cancel_runtime_worker( validate_workspace_scope(&api, &path.workspace_id)?; cancel_runtime_worker( State(api), - AxumPath((path.runtime_id, path.worker_id)), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), Json(request), ) .await @@ -5337,9 +5279,13 @@ async fn scoped_worker_protocol_ws( if let Err(err) = validate_workspace_scope(&api, &path.workspace_id) { return err.into_response(); } - worker_protocol_ws(State(api), AxumPath((path.runtime_id, path.worker_id)), ws) - .await - .into_response() + worker_protocol_ws( + State(api), + AxumPath((path.worker.runtime_id, path.worker.worker_id)), + ws, + ) + .await + .into_response() } async fn scoped_list_host_workers( @@ -6654,7 +6600,7 @@ fn record_browser_worker_spawn( )?; if let Some(working_directory) = worker.working_directory.as_ref() { let workdir_record = - workdir_record_from_summary(api, worker.runtime_id.as_str(), working_directory); + workdir_record_from_summary(api, worker.worker.runtime_id.as_str(), working_directory); api.store.upsert_workdir_registry(&workdir_record)?; link_worker_to_workdir( api, @@ -6671,13 +6617,13 @@ fn record_browser_worker_spawn( { if let Ok(result) = api .runtime - .working_directory(worker.runtime_id.as_str(), workdir_id) + .working_directory(worker.worker.runtime_id.as_str(), workdir_id) .map_err(|err| err.into_error()) { if let Some(status) = result.working_directory { let record = workdir_record_from_summary( api, - worker.runtime_id.as_str(), + worker.worker.runtime_id.as_str(), &status.summary, ); api.store.upsert_workdir_registry(&record)?; @@ -6692,8 +6638,8 @@ fn record_browser_worker_spawn( link_worker_to_workdir(api, &worker_record, workdir_id, None)?; } } - let runtime_id = worker.runtime_id.clone(); - let worker_id = worker.worker_id.clone(); + let runtime_id = worker.worker.runtime_id.clone(); + let worker_id = worker.worker.worker_id.clone(); let workspace_id = api.workspace_id().to_string(); let console_href = format!( "/w/{}/runtimes/{}/workers/{}/console", @@ -6703,8 +6649,7 @@ fn record_browser_worker_spawn( ); Ok(BrowserCreateWorkerResponse { workspace_id, - runtime_id, - worker_id, + worker_ref: RuntimeWorkerRef::new(&runtime_id, &worker_id), console_href, worker, diagnostics: result.diagnostics, @@ -6771,16 +6716,15 @@ async fn get_runtime_worker( State(api): State, AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, ) -> ApiResult> { + let worker_ref = RuntimeWorkerRef::new(runtime_id, worker_id); let worker = api .runtime - .worker(&runtime_id, &worker_id) + .worker(&worker_ref) .map_err(|err| err.into_error())?; let record = sync_worker_observation(&api, &worker)?; - let links = api.store.list_worker_workdir_links( - &api.config.workspace_id, - record.runtime_id.as_str(), - record.runtime_worker_id, - )?; + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, &record.worker)?; let workdirs = api .store .list_workdir_registry(&api.config.workspace_id, 500)?; @@ -6796,14 +6740,13 @@ async fn restore_runtime_worker( State(api): State, AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, ) -> ApiResult> { - let mut result = api.restore_workspace_worker(&runtime_id, &worker_id)?; + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); + let mut result = api.restore_workspace_worker(&worker)?; if let Some(worker) = result.worker.as_ref() { let record = sync_worker_observation(&api, worker)?; - let links = api.store.list_worker_workdir_links( - &api.config.workspace_id, - record.runtime_id.as_str(), - record.runtime_worker_id, - )?; + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, &record.worker)?; let workdirs = api .store .list_workdir_registry(&api.config.workspace_id, 500)?; @@ -6816,8 +6759,7 @@ async fn restore_runtime_worker( } Ok(Json(WorkerRestoreResponse { workspace_id: api.workspace_id().to_string(), - runtime_id, - worker_id, + worker_ref: RuntimeWorkerRef::new(&runtime_id, &worker_id), result, })) } @@ -6911,7 +6853,12 @@ async fn create_runtime_worker( ) -> ApiResult> { if let Some(assignment) = request.ticket_assignment.as_ref() { if let Some(worker) = existing_lifecycle_assignment_worker(&api, assignment, &runtime_id)? { - assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker.worker_id)?; + assign_ticket_worker_from_lifecycle( + &api, + assignment, + &runtime_id, + &worker.worker.worker_id, + )?; return Ok(Json(WorkerSpawnResult { state: WorkerOperationState::Accepted, worker: Some(worker), @@ -6983,9 +6930,14 @@ async fn create_runtime_worker( api.store.bind_ticket_assignment_operation_worker( &api.config.workspace_id, &assignment.operation_id, - &worker.worker_id, + &worker.worker.worker_id, + )?; + assign_ticket_worker_from_lifecycle( + &api, + assignment, + &runtime_id, + &worker.worker.worker_id, )?; - assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker.worker_id)?; } if worker.working_directory.is_none() { if let Some(workdir_id) = prepared_workdir_id.as_deref() { @@ -7046,9 +6998,10 @@ async fn send_runtime_worker_input( AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, Json(request): Json, ) -> ApiResult> { + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); let result = api .runtime - .send_input(&runtime_id, &worker_id, request) + .send_input(&worker, request) .map_err(|err| err.into_error())?; Ok(Json(result)) } @@ -7058,9 +7011,10 @@ async fn runtime_worker_completions( AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, Json(request): Json, ) -> ApiResult> { + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); let result = api .runtime - .worker_completions(&runtime_id, &worker_id, request) + .worker_completions(&worker, request) .map_err(|err| err.into_error())?; Ok(Json(result)) } @@ -7070,19 +7024,20 @@ async fn stop_runtime_worker( AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, Json(request): Json, ) -> ApiResult> { + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); let result = api .runtime - .stop_worker(&runtime_id, &worker_id, request) + .stop_worker(&worker, request) .map_err(|err| err.into_error())?; - let runtime_worker_id = parse_runtime_worker_id_for_registry(&worker_id)?; - let session_lock = current_worker_session_lock(&api, &runtime_id, runtime_worker_id); + parse_runtime_worker_id_for_registry(&worker.worker_id)?; + let session_lock = current_worker_session_lock(&api, &worker); let _session_guard = session_lock.lock().await; - close_current_worker_session_locked(&api, &runtime_id, runtime_worker_id).await?; - if let Some(record) = - api.store - .get_worker_registry(&api.config.workspace_id, &runtime_id, runtime_worker_id)? + close_current_worker_session_locked(&api, &worker).await?; + if let Some(record) = api + .store + .get_worker_registry(&api.config.workspace_id, &worker)? { - sync_linked_workdir_after_worker_stop(&api, &runtime_id, &record)?; + sync_linked_workdir_after_worker_stop(&api, &worker.runtime_id, &record)?; } Ok(Json(result)) } @@ -7092,9 +7047,10 @@ async fn cancel_runtime_worker( AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, Json(request): Json, ) -> ApiResult> { + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); let result = api .runtime - .cancel_worker(&runtime_id, &worker_id, request) + .cancel_worker(&worker, request) .map_err(|err| err.into_error())?; Ok(Json(result)) } @@ -7104,10 +7060,11 @@ async fn worker_protocol_ws( AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, ws: WebSocketUpgrade, ) -> impl IntoResponse { - let source = match api.observation_proxy.source(&runtime_id, &worker_id) { + let worker = RuntimeWorkerRef::new(&runtime_id, &worker_id); + let source = match api.observation_proxy.source(&worker) { Ok(source) => source, Err(ObservationProxyError::WorkerNotFound(_)) => { - match api.runtime.observation_source(&runtime_id, &worker_id) { + match api.runtime.observation_source(&worker) { Ok(source) => source, Err(error) => return ApiError::from(error.into_error()).into_response(), } @@ -7133,18 +7090,17 @@ pub(crate) struct WorkspaceWorkerProtocolConnection { pub(crate) async fn connect_workspace_worker_protocol( api: &WorkspaceApi, - runtime_id: &str, - worker_id: &str, + worker: &RuntimeWorkerRef, ) -> Result { - let source = match api.observation_proxy.source(runtime_id, worker_id) { + let source = match api.observation_proxy.source(worker) { Ok(source) => source, Err(ObservationProxyError::WorkerNotFound(_)) => api .runtime - .observation_source(runtime_id, worker_id) + .observation_source(worker) .map_err(|error| error.into_error())?, Err(error) => { return Err(Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), + runtime_id: worker.runtime_id.clone(), code: error.code().to_string(), message: error.message().to_string(), }); @@ -7178,7 +7134,7 @@ async fn connect_remote_worker_protocol( connect_async(request) .await .map_err(|error| Error::RuntimeOperationFailed { - runtime_id: config.runtime_id.clone(), + runtime_id: config.worker.runtime_id.clone(), code: "worker_protocol_connect_failed".to_string(), message: error.to_string(), })?; @@ -7217,7 +7173,7 @@ async fn connect_embedded_worker_protocol( RuntimeObservationClient::connect(&RuntimeObservationSource::Embedded(source.clone())) .await .map_err(|error| Error::RuntimeOperationFailed { - runtime_id: source.runtime_id.clone(), + runtime_id: source.worker.runtime_id.clone(), code: error.code().to_string(), message: error.message().to_string(), })?; @@ -7495,14 +7451,7 @@ fn workers_response(api: WorkspaceApi) -> ApiResult ApiResult { let _ = sync_worker_observation(&api, &worker); - observed.insert( - (worker.runtime_id.clone(), record.runtime_worker_id), - worker, - ); + observed.insert(record.worker.clone(), worker); } Err(RuntimeRegistryError::UnknownWorker { .. }) => {} Err(error) => diagnostics.push(RuntimeDiagnostic { @@ -7531,20 +7474,18 @@ fn workers_response(api: WorkspaceApi) -> ApiResult ApiResult { let timestamp = now_registry_timestamp(); - let runtime_worker_id = parse_runtime_worker_id_for_registry(worker.worker_id.as_str())?; - let existing = api.store.get_worker_registry( - &api.config.workspace_id, - worker.runtime_id.as_str(), - runtime_worker_id, - )?; + parse_runtime_worker_id_for_registry(worker.worker.worker_id.as_str())?; + let worker_ref = worker.worker.clone(); + let existing = api + .store + .get_worker_registry(&api.config.workspace_id, &worker_ref)?; let display_name = match (display_name_policy, existing.as_ref()) { (WorkerRegistryDisplayNamePolicy::PreserveExisting, Some(record)) => { record.display_name.clone() @@ -8382,8 +8322,7 @@ fn record_worker_summary( }; let record = WorkerRegistryRecord { workspace_id: api.config.workspace_id.clone(), - runtime_id: worker.runtime_id.as_str().to_string(), - runtime_worker_id, + worker: worker_ref.clone(), display_name, profile, retention_state: existing @@ -8392,8 +8331,8 @@ fn record_worker_summary( .unwrap_or_else(|| "normal".to_string()), transcript_ref: Some(format!( "runtime://{}/workers/{}/transcript", - worker.runtime_id.as_str(), - worker.worker_id.as_str() + worker.worker.runtime_id.as_str(), + worker.worker.worker_id.as_str() )), session_ref: None, summary_ref: None, @@ -8404,18 +8343,13 @@ fn record_worker_summary( api.store.upsert_worker_registry(&record)?; Ok(api .store - .get_worker_registry( - &api.config.workspace_id, - worker.runtime_id.as_str(), - runtime_worker_id, - )? + .get_worker_registry(&api.config.workspace_id, &worker_ref)? .unwrap_or(record)) } fn worker_summary_from_registry(record: &WorkerRegistryRecord) -> WorkerSummary { WorkerSummary { - worker_id: record.runtime_worker_id.to_string(), - runtime_id: record.runtime_id.clone(), + worker: record.worker.clone(), host_id: "backend-registry".to_string(), display_name: record.display_name.clone(), label: record.display_name.clone(), @@ -8468,11 +8402,8 @@ fn merge_worker_registry_projection( .find(|workdir| workdir.workdir_id == link.workdir_id) .map(|workdir| { let mut workdir_summary = workdir_summary_from_record(workdir); - workdir_summary.primary_worker_id = Some(WorkerId::new(record.runtime_worker_id)); workdir_summary.occupied_by = Some(WorkingDirectoryOccupancy { - runtime_id: record.runtime_id.clone(), - runtime_worker_id: record.runtime_worker_id, - worker_id: format!("{}:{}", record.runtime_id, record.runtime_worker_id), + worker: record.worker.clone(), display_name: record.display_name.clone(), linked_at: link.linked_at.clone(), }); @@ -8495,7 +8426,7 @@ fn sync_worker_observation( )?; if let Some(working_directory) = worker.working_directory.as_ref() { let workdir_record = - workdir_record_from_summary(api, worker.runtime_id.as_str(), working_directory); + workdir_record_from_summary(api, worker.worker.runtime_id.as_str(), working_directory); api.store.upsert_workdir_registry(&workdir_record)?; link_worker_to_workdir(api, &record, &working_directory.working_directory_id, None)?; } @@ -8635,11 +8566,9 @@ fn sync_linked_workdir_after_worker_stop( runtime_id: &str, worker_record: &WorkerRegistryRecord, ) -> ApiResult<()> { - let links = api.store.list_worker_workdir_links( - &api.config.workspace_id, - worker_record.runtime_id.as_str(), - worker_record.runtime_worker_id, - )?; + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, &worker_record.worker)?; for link in links { let result = api .runtime @@ -8757,21 +8686,19 @@ fn apply_workdir_occupancy_projection( return Ok(()); }; - let worker = api.store.get_worker_registry( - &api.config.workspace_id, - &link.runtime_id, - link.runtime_worker_id, - )?; - let display_name = worker - .as_ref() - .map(|worker| worker.display_name.clone()) - .unwrap_or_else(|| format!("{}:{}", link.runtime_id, link.runtime_worker_id)); - summary.primary_worker_id = Some(WorkerId::new(link.runtime_worker_id)); + let worker = api + .store + .get_worker_registry(&api.config.workspace_id, &link.worker)? + .ok_or_else(|| { + Error::RegistryInconsistency(format!( + "Workdir {} attachment references missing Worker {}:{}", + link.workdir_id, link.worker.runtime_id, link.worker.worker_id + )) + })?; + summary.primary_worker_id = None; summary.occupied_by = Some(WorkingDirectoryOccupancy { - runtime_id: link.runtime_id.clone(), - runtime_worker_id: link.runtime_worker_id, - worker_id: format!("{}:{}", link.runtime_id, link.runtime_worker_id), - display_name, + worker: link.worker.clone(), + display_name: worker.display_name, linked_at: link.linked_at.clone(), }); Ok(()) @@ -8795,8 +8722,7 @@ fn link_worker_to_workdir( let timestamp = now_registry_timestamp(); let record = WorkerWorkdirLinkRecord { workspace_id: api.config.workspace_id.clone(), - runtime_id: worker_record.runtime_id.clone(), - runtime_worker_id: worker_record.runtime_worker_id, + worker: worker_record.worker.clone(), workdir_id: workdir_id.to_string(), role: "attachment".to_string(), linked_at: timestamp, @@ -9372,8 +9298,7 @@ mod tests { fn backend_worker_projection_preserves_missing_rows_links_and_redacts_paths() { let worker = WorkerRegistryRecord { workspace_id: "workspace-1".to_string(), - runtime_id: "embedded".to_string(), - runtime_worker_id: 1, + worker: RuntimeWorkerRef::new("embedded", "1"), display_name: "Missing Worker".to_string(), profile: Some("builtin:coder".to_string()), retention_state: "pinned".to_string(), @@ -9400,8 +9325,7 @@ mod tests { }; let link = WorkerWorkdirLinkRecord { workspace_id: "workspace-1".to_string(), - runtime_id: worker.runtime_id.clone(), - runtime_worker_id: worker.runtime_worker_id.clone(), + worker: worker.worker.clone(), workdir_id: workdir.workdir_id.clone(), role: "attachment".to_string(), linked_at: "4".to_string(), @@ -9424,8 +9348,12 @@ mod tests { assert_eq!(working_directory.current_selector, None); assert_eq!(working_directory.current_ref.as_deref(), Some("fedcba")); let occupied_by = working_directory.occupied_by.as_ref().unwrap(); - assert_eq!(occupied_by.runtime_id, "embedded"); - assert_eq!(occupied_by.runtime_worker_id, 1); + assert_eq!(occupied_by.worker, RuntimeWorkerRef::new("embedded", "1")); + assert!(working_directory.primary_worker_id.is_none()); + let occupancy = serde_json::to_value(occupied_by).unwrap(); + assert_eq!(occupancy["runtime_id"], "embedded"); + assert_eq!(occupancy["worker_id"], "1"); + assert!(occupancy.get("runtime_worker_id").is_none()); assert_eq!(occupied_by.display_name, "Missing Worker"); assert_eq!(occupied_by.linked_at, "4"); let serialized = serde_json::to_string(&projected).unwrap(); @@ -9482,7 +9410,7 @@ mod tests { dedicated.singleton_key.as_deref(), Some(crate::hosts::WORKSPACE_ORCHESTRATOR_SINGLETON_KEY) ); - assert_ne!(dedicated.worker_id, generic.worker_id); + assert_ne!(dedicated.worker.worker_id, generic.worker_ref.worker_id); let Json(existing) = scoped_start_workspace_orchestrator( State(api.clone()), @@ -9493,7 +9421,10 @@ mod tests { .await .unwrap(); assert_eq!(existing.disposition, "existing"); - assert_eq!(existing.worker.unwrap().worker_id, dedicated.worker_id); + assert_eq!( + existing.worker.unwrap().worker.worker_id, + dedicated.worker.worker_id + ); let Json(status) = scoped_workspace_orchestrator_status( State(api), @@ -9501,7 +9432,10 @@ mod tests { ) .await .unwrap(); - assert_eq!(status.worker.unwrap().worker_id, dedicated.worker_id); + assert_eq!( + status.worker.unwrap().worker.worker_id, + dedicated.worker.worker_id + ); } #[tokio::test] @@ -9543,8 +9477,7 @@ mod tests { api.store .upsert_worker_registry(&WorkerRegistryRecord { workspace_id: TEST_WORKSPACE_ID.to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - runtime_worker_id: 7, + worker: RuntimeWorkerRef::new(EMBEDDED_WORKER_RUNTIME_ID, "7"), display_name: "Worker Seven".to_string(), profile: Some("builtin:coder".to_string()), retention_state: "normal".to_string(), @@ -9559,8 +9492,7 @@ mod tests { api.store .attach_worker_workdir(&WorkerWorkdirLinkRecord { workspace_id: TEST_WORKSPACE_ID.to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - runtime_worker_id: 7, + worker: RuntimeWorkerRef::new(EMBEDDED_WORKER_RUNTIME_ID, "7"), workdir_id: "managed".to_string(), role: "attachment".to_string(), linked_at: "3".to_string(), @@ -9581,8 +9513,10 @@ mod tests { .find(|summary| summary.working_directory_id == "managed") .unwrap(); let occupied_by = managed.occupied_by.as_ref().unwrap(); - assert_eq!(occupied_by.runtime_id, EMBEDDED_WORKER_RUNTIME_ID); - assert_eq!(occupied_by.runtime_worker_id, 7); + assert_eq!( + occupied_by.worker, + RuntimeWorkerRef::new(EMBEDDED_WORKER_RUNTIME_ID, "7") + ); assert_eq!(occupied_by.display_name, "Worker Seven"); assert_eq!(occupied_by.linked_at, "3"); @@ -10235,10 +10169,13 @@ mod tests { ) .unwrap(); assert_eq!(existing.state, WorkerOperationState::Accepted); - let worker_id = existing.worker.unwrap().worker_id; + let worker_id = existing.worker.unwrap().worker.worker_id; let worker = api .runtime - .worker(EMBEDDED_WORKER_RUNTIME_ID, &worker_id) + .worker(&RuntimeWorkerRef::new( + EMBEDDED_WORKER_RUNTIME_ID, + &worker_id, + )) .unwrap(); assert_eq!(worker.state, "idle"); assert_eq!(worker.display_name, "Memory Consolidation"); @@ -10272,7 +10209,7 @@ mod tests { .filter(|worker| is_memory_consolidation_worker(worker)) .collect::>(); assert_eq!(consolidaters_after_second.len(), 1); - assert_eq!(consolidaters_after_second[0].worker_id, worker_id); + assert_eq!(consolidaters_after_second[0].worker.worker_id, worker_id); } fn init_clean_git_workspace(path: &std::path::Path) { @@ -10330,8 +10267,7 @@ mod tests { workspace_id: TEST_WORKSPACE_ID.to_string(), ticket_id: ticket_id.clone(), assignment_id: "assignment-api-1".to_string(), - runtime_id: "embedded".to_string(), - worker_id: "42".to_string(), + worker: RuntimeWorkerRef::new("embedded", "42"), assigned_by: "test-user".to_string(), assigned_at: TEST_CREATED_AT.to_string(), }; @@ -10430,8 +10366,10 @@ mod tests { workspace_id: TEST_WORKSPACE_ID.to_string(), ticket_id: ticket_ref.id.clone(), assignment_id: "notify-assignment".to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - worker_id: recipient_worker.worker_id.clone(), + worker: RuntimeWorkerRef::new( + EMBEDDED_WORKER_RUNTIME_ID, + recipient_worker.worker.worker_id.clone(), + ), assigned_by: "test-user".to_string(), assigned_at: TEST_CREATED_AT.to_string(), }, @@ -10448,7 +10386,7 @@ mod tests { ); headers.insert( "x-yoi-worker-id", - axum::http::HeaderValue::from_str(&source_worker.worker_id).unwrap(), + axum::http::HeaderValue::from_str(&source_worker.worker.worker_id).unwrap(), ); let response = build_router(api.clone()) .oneshot( @@ -10460,7 +10398,7 @@ mod tests { )) .header("content-type", "application/json") .header("x-yoi-runtime-id", EMBEDDED_WORKER_RUNTIME_ID) - .header("x-yoi-worker-id", &source_worker.worker_id) + .header("x-yoi-worker-id", &source_worker.worker.worker_id) .body(Body::from( serde_json::to_vec(&NewTicketEvent::new( TicketEventKind::Comment, @@ -10487,7 +10425,7 @@ mod tests { .attributes .get("source_worker_id") .map(String::as_str), - Some(source_worker.worker_id.as_str()) + Some(source_worker.worker.worker_id.as_str()) ); assert_eq!( committed_event @@ -10531,8 +10469,10 @@ mod tests { workspace_id: TEST_WORKSPACE_ID.to_string(), ticket_id: ticket_ref.id.clone(), assignment_id: "source-assignment".to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - worker_id: source_worker.worker_id.clone(), + worker: RuntimeWorkerRef::new( + EMBEDDED_WORKER_RUNTIME_ID, + source_worker.worker.worker_id.clone(), + ), assigned_by: "test-user".to_string(), assigned_at: TEST_CREATED_AT.to_string(), }, @@ -10635,7 +10575,7 @@ mod tests { ); headers.insert( "x-yoi-worker-id", - axum::http::HeaderValue::from_str(&source.worker_id).unwrap(), + axum::http::HeaderValue::from_str(&source.worker.worker_id).unwrap(), ); let _ = execute_worker_ticket_test_operation( State(api.clone()), @@ -10700,7 +10640,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_id.clone(); + let worker_id = started.worker.as_ref().unwrap().worker.worker_id.clone(); maybe_dispatch_orchestrator_turn_end( &api, &worker_id, @@ -10782,8 +10722,8 @@ mod tests { projected .worker .as_ref() - .map(|worker| worker.worker_id.as_str()), - Some(first_worker.worker_id.as_str()) + .map(|worker| worker.worker.worker_id.as_str()), + Some(first_worker.worker.worker_id.as_str()) ); let Json(retried) = scoped_create_runtime_worker( State(api.clone()), @@ -10795,7 +10735,10 @@ mod tests { ) .await .unwrap(); - assert_eq!(retried.worker.unwrap().worker_id, first_worker.worker_id); + assert_eq!( + retried.worker.unwrap().worker.worker_id, + first_worker.worker.worker_id + ); assert_eq!( api.store .list_ticket_worker_assignment_events(TEST_WORKSPACE_ID, &first_ticket.id, 10,) @@ -10822,8 +10765,7 @@ mod tests { .unwrap(); api.runtime .stop_worker( - EMBEDDED_WORKER_RUNTIME_ID, - &first_worker.worker_id, + &first_worker.worker, WorkerLifecycleRequest { reason: Some("restore assignment test".to_string()), ticket_assignment: None, @@ -10837,8 +10779,10 @@ mod tests { State(api.clone()), AxumPath(ScopedRuntimeWorkerPath { workspace_id: TEST_WORKSPACE_ID.to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - worker_id: first_worker.worker_id.clone(), + worker: RuntimeWorkerRef::new( + EMBEDDED_WORKER_RUNTIME_ID, + first_worker.worker.worker_id.clone(), + ), }), Query(RestoreTicketAssignmentQuery { ticket_id: Some(second_ticket.id.clone()), @@ -10851,8 +10795,10 @@ mod tests { State(api.clone()), AxumPath(ScopedRuntimeWorkerPath { workspace_id: TEST_WORKSPACE_ID.to_string(), - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - worker_id: first_worker.worker_id.clone(), + worker: RuntimeWorkerRef::new( + EMBEDDED_WORKER_RUNTIME_ID, + first_worker.worker.worker_id.clone(), + ), }), Query(RestoreTicketAssignmentQuery { ticket_id: Some(second_ticket.id.clone()), @@ -10861,14 +10807,20 @@ mod tests { ) .await .unwrap(); - assert_eq!(retried_restore.worker_id, first_worker.worker_id); + assert_eq!( + retried_restore.worker_ref.worker_id, + first_worker.worker.worker_id + ); assert_eq!(retried_restore.result.state, WorkerOperationState::Accepted); let restored_assignment = api .store .get_current_ticket_worker_assignment(TEST_WORKSPACE_ID, &second_ticket.id) .unwrap() .unwrap(); - assert_eq!(restored_assignment.worker_id, first_worker.worker_id); + assert_eq!( + restored_assignment.worker.worker_id, + first_worker.worker.worker_id + ); api.store .clear_current_ticket_worker_assignment( @@ -10914,7 +10866,7 @@ mod tests { api.store .get_ticket_assignment_operation(TEST_WORKSPACE_ID, "pending-spawn-operation") .unwrap() - .is_some_and(|operation| operation.worker_id.is_none()) + .is_some_and(|operation| operation.worker.is_none()) ); let worker_count_before_retry = api.runtime.list_workers(100).items.len(); let Json(reconciled) = scoped_create_runtime_worker( @@ -10928,8 +10880,8 @@ mod tests { .await .unwrap(); assert_eq!( - reconciled.worker.unwrap().worker_id, - spawned_before_backend_failure.worker_id + reconciled.worker.unwrap().worker.worker_id, + spawned_before_backend_failure.worker.worker_id ); assert_eq!( api.runtime.list_workers(100).items.len(), @@ -11174,8 +11126,7 @@ mod tests { api.store .upsert_worker_registry(&WorkerRegistryRecord { workspace_id: api.config.workspace_id.clone(), - runtime_id: "runtime-test".to_string(), - runtime_worker_id, + worker: RuntimeWorkerRef::new("runtime-test", runtime_worker_id.to_string()), display_name: runtime_worker_id.to_string(), profile: None, retention_state: retention_state.to_string(), @@ -11215,8 +11166,7 @@ mod tests { api.store .attach_worker_workdir(&WorkerWorkdirLinkRecord { workspace_id: api.config.workspace_id.clone(), - runtime_id: "runtime-test".to_string(), - runtime_worker_id, + worker: RuntimeWorkerRef::new("runtime-test", runtime_worker_id.to_string()), workdir_id: workdir_id.to_string(), role: "attachment".to_string(), linked_at: now_registry_timestamp(), @@ -11240,12 +11190,14 @@ mod tests { #[test] fn workdir_session_owner_is_only_sent_for_same_runtime_worker() { + let embedded_worker = RuntimeWorkerRef::new("embedded-worker-runtime", "5"); assert_eq!( - runtime_local_owner_worker_id("embedded-worker-runtime", "arcadia", 5), + runtime_local_owner_worker_id(&embedded_worker, "arcadia"), None ); + let arcadia_worker = RuntimeWorkerRef::new("arcadia", "30"); assert_eq!( - runtime_local_owner_worker_id("arcadia", "arcadia", 30).as_deref(), + runtime_local_owner_worker_id(&arcadia_worker, "arcadia"), Some("30") ); } @@ -11593,8 +11545,7 @@ mod tests { let runtime_id = "runtime-test"; let handle = broker.issue_profile_source_archive_handle( "workspace-test", - Some(runtime_id), - None, + crate::resource_broker::BackendResourceTarget::Runtime(runtime_id), archive, ); let app = build_router(api); @@ -12852,12 +12803,12 @@ mod tests { ) .expect("spawn worker"); assert_eq!(spawned.state, WorkerOperationState::Accepted); - let worker_id = spawned.worker.expect("created worker").worker_id; + let worker_id = spawned.worker.expect("created worker").worker.worker_id; + let worker_ref = RuntimeWorkerRef::new("embedded-worker-runtime", &worker_id); let sent = api .runtime .send_input( - "embedded-worker-runtime", - &worker_id, + &worker_ref, WorkerInputRequest { kind: WorkerInputKind::User, content: "persist me".to_string(), @@ -12869,10 +12820,7 @@ mod tests { let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2); loop { - let detail = api - .runtime - .worker("embedded-worker-runtime", &worker_id) - .expect("worker detail"); + let detail = api.runtime.worker(&worker_ref).expect("worker detail"); if detail.state == "idle" { break; } @@ -12893,7 +12841,7 @@ mod tests { .expect("restored fs-backed api starts"); let restored_worker = restored .runtime - .worker("embedded-worker-runtime", &worker_id) + .worker(&worker_ref) .expect("restored worker"); assert_eq!(restored_worker.state, "stopped"); assert!(!restored_worker.capabilities.can_stop); @@ -12912,8 +12860,7 @@ mod tests { let rejected_input = restored .runtime .send_input( - "embedded-worker-runtime", - &worker_id, + &worker_ref, WorkerInputRequest { kind: WorkerInputKind::User, content: "should not be routed to corrupted handle".to_string(), @@ -13210,8 +13157,7 @@ mod tests { async fn proxies_worker_protocol_ws_as_raw_events() { let (runtime, worker_ref, endpoint) = spawn_runtime_worker().await; let source = RuntimeObservationSourceConfig { - runtime_id: "runtime-a".into(), - worker_id: "worker-a".into(), + worker: RuntimeWorkerRef::new("runtime-a", "worker-a"), endpoint, bearer_token: None, }; @@ -13257,8 +13203,7 @@ mod tests { let (_runtime, _worker_ref, endpoint) = spawn_runtime_worker().await; let endpoint = endpoint.replace("/protocol/ws", "/missing-worker/protocol/ws"); let source = RuntimeObservationSourceConfig { - runtime_id: "runtime-a".into(), - worker_id: "worker-a".into(), + worker: RuntimeWorkerRef::new("runtime-a", "worker-a"), endpoint, bearer_token: None, }; @@ -13311,8 +13256,8 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let store = SqliteWorkspaceStore::in_memory().unwrap(); let mut config = test_server_config(dir.path()); - let runtime_id = source.runtime_id.clone(); - let worker_id = source.worker_id.clone(); + let runtime_id = source.worker.runtime_id.clone(); + let worker_id = source.worker.worker_id.clone(); config.runtime_event_sources.push(source); let api = WorkspaceApi::new_with_execution_backend( config, @@ -13375,7 +13320,7 @@ mod tests { .spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request) .unwrap(); assert_eq!(spawned.state, WorkerOperationState::Accepted); - let worker_id = spawned.worker.unwrap().worker_id; + let worker_id = spawned.worker.unwrap().worker.worker_id; let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index bfcc517c..5e6a0202 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -6,6 +6,8 @@ use async_trait::async_trait; use rusqlite::{Connection, OptionalExtension, params}; use serde::{Deserialize, Serialize}; +use worker_runtime::identity::RuntimeWorkerRef; + use crate::{Error, Result}; const WORKSPACES_V0_COLUMNS: &[&str] = &[ @@ -268,8 +270,7 @@ pub struct DeviceLoginFlowRecord { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerRegistryRecord { pub workspace_id: String, - pub runtime_id: String, - pub runtime_worker_id: u64, + pub worker: RuntimeWorkerRef, pub display_name: String, pub profile: Option, /// Retention state is explicit so `pinned` can be represented before prune exists. @@ -287,8 +288,7 @@ pub struct TicketWorkerAssignmentRecord { pub workspace_id: String, pub ticket_id: String, pub assignment_id: String, - pub runtime_id: String, - pub worker_id: String, + pub worker: RuntimeWorkerRef, pub assigned_by: String, pub assigned_at: String, } @@ -330,8 +330,7 @@ pub struct WorkdirRegistryRecord { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkerWorkdirLinkRecord { pub workspace_id: String, - pub runtime_id: String, - pub runtime_worker_id: u64, + pub worker: RuntimeWorkerRef, pub workdir_id: String, pub role: String, pub linked_at: String, @@ -541,8 +540,7 @@ pub trait ControlPlaneStore: Send + Sync { fn get_worker_registry( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result>; fn list_worker_registry( &self, @@ -552,17 +550,12 @@ pub trait ControlPlaneStore: Send + Sync { fn update_worker_retention( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, retention_state: &str, updated_at: &str, ) -> Result; - fn delete_worker_registry( - &self, - workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, - ) -> Result; + fn delete_worker_registry(&self, workspace_id: &str, worker: &RuntimeWorkerRef) + -> Result; fn get_ticket_assignment_operation( &self, @@ -653,22 +646,19 @@ pub trait ControlPlaneStore: Send + Sync { fn detach_worker_workdir( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, expected_workdir_id: Option<&str>, unlinked_at: &str, ) -> Result>; fn worker_workdir_link_history_exists( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result; fn list_worker_workdir_links( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result>; fn list_workdir_worker_links( &self, @@ -1637,8 +1627,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { updated_at = excluded.updated_at"#, params![ record.workspace_id, - record.runtime_id, - record.runtime_worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.display_name, record.profile, record.retention_state, @@ -1657,8 +1647,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn get_worker_registry( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result> { self.with_conn(|conn| { conn.query_row( @@ -1666,7 +1655,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { "WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3", ) .as_str(), - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], read_worker_registry_record, ) .optional() @@ -1696,8 +1685,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn update_worker_retention( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, retention_state: &str, updated_at: &str, ) -> Result { @@ -1708,8 +1696,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3"#, params![ workspace_id, - runtime_id, - runtime_worker_id, + worker.runtime_id, + worker.worker_id, retention_state, updated_at ], @@ -1721,8 +1709,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn delete_worker_registry( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result { self.with_conn(|conn| { let tx = conn.unchecked_transaction()?; @@ -1730,11 +1717,11 @@ impl ControlPlaneStore for SqliteWorkspaceStore { r#"UPDATE worker_workdir_links SET unlinked_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now') WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3 AND unlinked_at IS NULL"#, - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], )?; let changed = tx.execute( "DELETE FROM worker_registry WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3", - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], )?; tx.commit()?; Ok(changed > 0) @@ -1787,7 +1774,12 @@ impl ControlPlaneStore for SqliteWorkspaceStore { if existing.action == "assign" && existing.ticket_id == ticket_id && existing.runtime_id.as_deref() == Some(runtime_id) - && (worker_id.is_none() || existing.worker_id.as_deref() == worker_id) + && (worker_id.is_none() + || existing + .worker + .as_ref() + .map(|worker| worker.worker_id.as_str()) + == worker_id) && existing.expected_assignment_id.is_none() && existing.request_fingerprint.as_deref() == Some(request_fingerprint) { @@ -1855,8 +1847,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { { if existing.action != if allow_reassign { "reassign" } else { "assign" } || existing.ticket_id != record.ticket_id - || existing.runtime_id.as_deref() != Some(record.runtime_id.as_str()) - || existing.worker_id.as_deref() != Some(record.worker_id.as_str()) + || existing.worker.as_ref() != Some(&record.worker) || existing.expected_assignment_id.as_deref() != expected_assignment_id { return Err(Error::TicketAssignmentConflict(format!( @@ -1937,8 +1928,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.workspace_id, record.ticket_id, record.assignment_id, - record.runtime_id, - record.worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.assigned_by, record.assigned_at, ], @@ -1952,8 +1943,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.workspace_id, record.ticket_id, record.assignment_id, - record.runtime_id, - record.worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.assigned_at, ], ) @@ -1966,8 +1957,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.workspace_id, record.ticket_id, record.assignment_id, - record.runtime_id, - record.worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.assigned_at, ], ) @@ -1976,7 +1967,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { return Err(map_assignment_constraint( error, &record.ticket_id, - &record.worker_id, + &record.worker.worker_id, )); } tx.execute( @@ -2024,8 +2015,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { operation_id, if allow_reassign { "reassign" } else { "assign" }, record.ticket_id, - record.runtime_id, - record.worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.assignment_id, expected_assignment_id, record.assigned_at, @@ -2116,8 +2107,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { workspace_id, operation_id, ticket_id, - previous.runtime_id, - previous.worker_id, + previous.worker.runtime_id, + previous.worker.worker_id, previous.assignment_id, expected_assignment_id, created_at, @@ -2369,8 +2360,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { unlinked_at = NULL"#, params![ record.workspace_id, - record.runtime_id, - record.runtime_worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.workdir_id, record.role, record.linked_at, @@ -2451,8 +2442,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3 AND unlinked_at IS NULL"#, params![ record.workspace_id, - record.runtime_id, - record.runtime_worker_id, + record.worker.runtime_id, + record.worker.worker_id, ], read_worker_workdir_link_record, ) @@ -2464,7 +2455,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { } return Err(Error::WorkdirAttachmentConflict(format!( "Worker {}:{} is already attached to Workdir {}", - record.runtime_id, record.runtime_worker_id, active.workdir_id + record.worker.runtime_id, record.worker.worker_id, active.workdir_id ))); } let active_for_workdir = tx @@ -2479,7 +2470,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { if let Some(active) = active_for_workdir { return Err(Error::WorkdirAttachmentConflict(format!( "Workdir {} is already attached to Worker {}:{}", - record.workdir_id, active.runtime_id, active.runtime_worker_id + record.workdir_id, active.worker.runtime_id, active.worker.worker_id ))); } let write = tx.execute( @@ -2491,8 +2482,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { unlinked_at = NULL"#, params![ record.workspace_id, - record.runtime_id, - record.runtime_worker_id, + record.worker.runtime_id, + record.worker.worker_id, record.workdir_id, record.role, record.linked_at, @@ -2515,8 +2506,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn detach_worker_workdir( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, expected_workdir_id: Option<&str>, unlinked_at: &str, ) -> Result> { @@ -2530,7 +2520,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { r#"SELECT workspace_id, runtime_id, runtime_worker_id, workdir_id, role, linked_at, unlinked_at FROM worker_workdir_links WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3 AND unlinked_at IS NULL"#, - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], read_worker_workdir_link_record, ) .optional()?; @@ -2541,8 +2531,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { if let Some(expected_workdir_id) = expected_workdir_id { if active.workdir_id != expected_workdir_id { return Err(Error::WorkdirAttachmentConflict(format!( - "Worker {runtime_id}:{runtime_worker_id} is attached to Workdir {}, not {expected_workdir_id}", - active.workdir_id + "Worker {}:{} is attached to Workdir {}, not {expected_workdir_id}", + worker.runtime_id, worker.worker_id, active.workdir_id ))); } } @@ -2550,11 +2540,12 @@ impl ControlPlaneStore for SqliteWorkspaceStore { r#"UPDATE worker_workdir_links SET unlinked_at = ?4 WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3 AND unlinked_at IS NULL"#, - params![workspace_id, runtime_id, runtime_worker_id, unlinked_at], + params![workspace_id, worker.runtime_id, worker.worker_id, unlinked_at], )?; if changed != 1 { return Err(Error::WorkdirAttachmentConflict(format!( - "Worker {runtime_id}:{runtime_worker_id} attachment changed during detach" + "Worker {}:{} attachment changed during detach", + worker.runtime_id, worker.worker_id ))); } tx.commit()?; @@ -2568,8 +2559,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn worker_workdir_link_history_exists( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result { self.with_conn(|conn| { let exists = conn.query_row( @@ -2577,7 +2567,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { SELECT 1 FROM worker_workdir_links WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3 )"#, - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], |row| row.get(0), )?; Ok(exists) @@ -2587,8 +2577,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { fn list_worker_workdir_links( &self, workspace_id: &str, - runtime_id: &str, - runtime_worker_id: u64, + worker: &RuntimeWorkerRef, ) -> Result> { self.with_conn(|conn| { let mut stmt = conn.prepare( @@ -2598,7 +2587,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { ORDER BY linked_at DESC"#, )?; let rows = stmt.query_map( - params![workspace_id, runtime_id, runtime_worker_id], + params![workspace_id, worker.runtime_id, worker.worker_id], read_worker_workdir_link_record, )?; rows.collect::, _>>() @@ -2807,8 +2796,7 @@ fn read_worker_workdir_link_record( ) -> rusqlite::Result { Ok(WorkerWorkdirLinkRecord { workspace_id: row.get(0)?, - runtime_id: row.get(1)?, - runtime_worker_id: row.get(2)?, + worker: RuntimeWorkerRef::new(row.get::<_, String>(1)?, row.get::<_, u64>(2)?.to_string()), workdir_id: row.get(3)?, role: row.get(4)?, linked_at: row.get(5)?, @@ -2900,8 +2888,7 @@ fn worker_registry_select_sql(where_clause: &str) -> String { fn read_worker_registry_record(row: &rusqlite::Row<'_>) -> rusqlite::Result { Ok(WorkerRegistryRecord { workspace_id: row.get(0)?, - runtime_id: row.get(1)?, - runtime_worker_id: row.get(2)?, + worker: RuntimeWorkerRef::new(row.get::<_, String>(1)?, row.get::<_, u64>(2)?.to_string()), display_name: row.get(3)?, profile: row.get(4)?, retention_state: row.get(5)?, @@ -2933,8 +2920,7 @@ fn read_ticket_worker_assignment_record( workspace_id: row.get(0)?, ticket_id: row.get(1)?, assignment_id: row.get(2)?, - runtime_id: row.get(3)?, - worker_id: row.get(4)?, + worker: RuntimeWorkerRef::new(row.get::<_, String>(3)?, row.get::<_, String>(4)?), assigned_by: row.get(5)?, assigned_at: row.get(6)?, }) @@ -2960,7 +2946,7 @@ pub struct TicketAssignmentOperationRecord { pub action: String, pub ticket_id: String, pub runtime_id: Option, - pub worker_id: Option, + pub worker: Option, pub assignment_id: Option, pub expected_assignment_id: Option, pub request_fingerprint: Option, @@ -2978,11 +2964,15 @@ fn read_assignment_operation( WHERE workspace_id = ?1 AND operation_id = ?2"#, params![workspace_id, operation_id], |row| { + let runtime_id: Option = row.get(2)?; + let worker_id: Option = row.get(3)?; Ok(TicketAssignmentOperationRecord { action: row.get(0)?, ticket_id: row.get(1)?, - runtime_id: row.get(2)?, - worker_id: row.get(3)?, + runtime_id: runtime_id.clone(), + worker: runtime_id + .zip(worker_id) + .map(|(runtime_id, worker_id)| RuntimeWorkerRef::new(runtime_id, worker_id)), assignment_id: row.get(4)?, expected_assignment_id: row.get(5)?, request_fingerprint: row.get(6)?, @@ -4204,8 +4194,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); workspace_id: "workspace-a".to_string(), ticket_id: "ticket-1".to_string(), assignment_id: "assignment-1".to_string(), - runtime_id: "runtime-1".to_string(), - worker_id: "worker-1".to_string(), + worker: RuntimeWorkerRef::new("runtime-1", "worker-1"), assigned_by: "user-1".to_string(), assigned_at: "2026-07-31T00:00:01Z".to_string(), }; @@ -4239,7 +4228,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .set_current_ticket_worker_assignment( &TicketWorkerAssignmentRecord { assignment_id: "implicit-reassign".to_string(), - worker_id: "worker-other".to_string(), + worker: RuntimeWorkerRef::new("runtime-1", "worker-other"), ..first.clone() }, None, @@ -4272,8 +4261,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); let second = TicketWorkerAssignmentRecord { assignment_id: "assignment-2".to_string(), - runtime_id: "runtime-2".to_string(), - worker_id: "worker-2".to_string(), + worker: RuntimeWorkerRef::new("runtime-2", "worker-2"), assigned_by: "user-2".to_string(), assigned_at: "2026-07-31T00:00:02Z".to_string(), ..first.clone() @@ -4360,7 +4348,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .get_ticket_assignment_operation("workspace-a", "reserved-operation") .unwrap() .unwrap(); - assert_eq!(pending.worker_id, None); + assert_eq!(pending.worker, None); assert_eq!( pending.request_fingerprint.as_deref(), Some("sha256:reserved") @@ -4376,8 +4364,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); workspace_id: "workspace-a".to_string(), ticket_id: "ticket-3".to_string(), assignment_id: "assignment-3".to_string(), - runtime_id: "runtime-3".to_string(), - worker_id: "worker-3".to_string(), + worker: RuntimeWorkerRef::new("runtime-3", "worker-3"), assigned_by: "runtime".to_string(), assigned_at: "2026-07-31T00:00:06Z".to_string(), }; @@ -4943,8 +4930,7 @@ CREATE TABLE ticket_assignment_operations ( let worker = WorkerRegistryRecord { workspace_id: "local-dev".to_string(), - runtime_id: "embedded".to_string(), - runtime_worker_id: 1, + worker: RuntimeWorkerRef::new("embedded", "1"), display_name: "Browser 1".to_string(), profile: Some("builtin:companion".to_string()), retention_state: "pinned".to_string(), @@ -4996,8 +4982,7 @@ CREATE TABLE ticket_assignment_operations ( let link = WorkerWorkdirLinkRecord { workspace_id: "local-dev".to_string(), - runtime_id: worker.runtime_id.clone(), - runtime_worker_id: worker.runtime_worker_id, + worker: worker.worker.clone(), workdir_id: workdir.workdir_id.clone(), role: "attachment".to_string(), linked_at: "4".to_string(), @@ -5033,7 +5018,7 @@ CREATE TABLE ticket_assignment_operations ( assert_eq!( store - .get_worker_registry("local-dev", "embedded", 1) + .get_worker_registry("local-dev", &worker.worker) .unwrap(), Some(expected_worker.clone()) ); @@ -5049,7 +5034,7 @@ CREATE TABLE ticket_assignment_operations ( ); assert_eq!( store - .list_worker_workdir_links("local-dev", "embedded", 1) + .list_worker_workdir_links("local-dev", &worker.worker) .unwrap(), vec![link.clone()] ); @@ -5065,7 +5050,7 @@ CREATE TABLE ticket_assignment_operations ( )); let second_worker = WorkerRegistryRecord { - runtime_worker_id: 2, + worker: RuntimeWorkerRef::new("embedded", "2"), display_name: "Browser 2".to_string(), created_at: "5".to_string(), updated_at: "5".to_string(), @@ -5073,7 +5058,7 @@ CREATE TABLE ticket_assignment_operations ( }; store.upsert_worker_registry(&second_worker).unwrap(); let workdir_conflict = WorkerWorkdirLinkRecord { - runtime_worker_id: second_worker.runtime_worker_id, + worker: second_worker.worker.clone(), linked_at: "5".to_string(), ..link.clone() }; @@ -5082,17 +5067,17 @@ CREATE TABLE ticket_assignment_operations ( Err(Error::WorkdirAttachmentConflict(_)) )); assert!(matches!( - store.detach_worker_workdir("local-dev", "embedded", 1, Some("wrong-workdir"), "6"), + store.detach_worker_workdir("local-dev", &worker.worker, Some("wrong-workdir"), "6",), Err(Error::WorkdirAttachmentConflict(_)) )); let detached = store - .detach_worker_workdir("local-dev", "embedded", 1, Some(&workdir.workdir_id), "6") + .detach_worker_workdir("local-dev", &worker.worker, Some(&workdir.workdir_id), "6") .unwrap() .unwrap(); assert_eq!(detached.unlinked_at.as_deref(), Some("6")); assert!( store - .worker_workdir_link_history_exists("local-dev", "embedded", 1) + .worker_workdir_link_history_exists("local-dev", &worker.worker) .unwrap() ); assert_eq!( diff --git a/crates/workspace-server/src/workspace_subscription.rs b/crates/workspace-server/src/workspace_subscription.rs index 9584578a..0e5619ca 100644 --- a/crates/workspace-server/src/workspace_subscription.rs +++ b/crates/workspace-server/src/workspace_subscription.rs @@ -8,6 +8,7 @@ use protocol::subscription::{ SubscriptionResponse, SubscriptionSnapshot, SubscriptionTerminationCode, SubscriptionWorker, }; use tokio::sync::mpsc; +use worker_runtime::identity::RuntimeWorkerRef; use crate::runtime_subscription::{BrokerSubscriptionEvent, RuntimeSubscriptionBroker}; use crate::server::{WorkspaceApi, connect_workspace_worker_protocol}; @@ -81,13 +82,8 @@ pub(crate) async fn serve_workspace_subscription(api: WorkspaceApi, socket: WebS worker_id, runtime_id: Some(runtime_id), } => { - match connect_workspace_worker_protocol( - &api, - &runtime_id, - worker_id.as_str(), - ) - .await - { + let worker = RuntimeWorkerRef::new(&runtime_id, worker_id.as_str()); + match connect_workspace_worker_protocol(&api, &worker).await { Ok(connection) => { let methods = connection.methods.clone(); let task = tokio::spawn(run_worker_protocol( @@ -323,16 +319,15 @@ async fn run_workspace_workers( } } - let mut revisions = HashMap::::new(); - let mut initial_workers = workers - .values_mut() - .flat_map(|runtime| runtime.values_mut()) - .map(|worker| { - let key = worker_key(worker.runtime_id.as_deref(), worker.worker_id.as_str()); - worker.subject_revision = next_revision(&mut revisions, &key); - worker.clone() - }) - .collect::>(); + let mut revisions = HashMap::::new(); + let mut initial_workers = Vec::new(); + for (runtime_id, runtime) in &mut workers { + for worker in runtime.values_mut() { + let worker_ref = RuntimeWorkerRef::new(runtime_id, worker.worker_id.as_str()); + worker.subject_revision = next_revision(&mut revisions, &worker_ref); + initial_workers.push(worker.clone()); + } + } sort_workers(&mut initial_workers); if send_frame( &outbound, @@ -359,8 +354,8 @@ async fn run_workspace_workers( BrokerSubscriptionEvent::Snapshot { snapshot, .. } => { let removed = workers.remove(&runtime_id).unwrap_or_default(); for worker in removed.values() { - let key = worker_key(Some(&runtime_id), worker.worker_id.as_str()); - let revision = next_revision(&mut revisions, &key); + let worker_ref = RuntimeWorkerRef::new(&runtime_id, worker.worker_id.as_str()); + let revision = next_revision(&mut revisions, &worker_ref); if send_event( &outbound, &subscription_id, @@ -379,8 +374,9 @@ async fn run_workspace_workers( install_snapshot(&mut workers, &runtime_id, snapshot); if let Some(current) = workers.get_mut(&runtime_id) { for worker in current.values_mut() { - let key = worker_key(Some(&runtime_id), worker.worker_id.as_str()); - let revision = next_revision(&mut revisions, &key); + let worker_ref = + RuntimeWorkerRef::new(&runtime_id, worker.worker_id.as_str()); + let revision = next_revision(&mut revisions, &worker_ref); worker.subject_revision = revision; if send_event( &outbound, @@ -401,8 +397,8 @@ async fn run_workspace_workers( BrokerSubscriptionEvent::Event { payload, .. } => match payload { SubscriptionEventPayload::WorkerUpserted { mut worker } => { worker.runtime_id = Some(runtime_id.clone()); - let key = worker_key(Some(&runtime_id), worker.worker_id.as_str()); - let revision = next_revision(&mut revisions, &key); + let worker_ref = RuntimeWorkerRef::new(&runtime_id, worker.worker_id.as_str()); + let revision = next_revision(&mut revisions, &worker_ref); worker.subject_revision = revision; workers .entry(runtime_id) @@ -425,8 +421,8 @@ async fn run_workspace_workers( .entry(runtime_id.clone()) .or_default() .remove(worker_id.as_str()); - let key = worker_key(Some(&runtime_id), worker_id.as_str()); - let revision = next_revision(&mut revisions, &key); + let worker_ref = RuntimeWorkerRef::new(&runtime_id, worker_id.as_str()); + let revision = next_revision(&mut revisions, &worker_ref); if send_event( &outbound, &subscription_id, @@ -534,14 +530,11 @@ async fn send_frame( .map_err(|_| ()) } -fn next_revision(revisions: &mut HashMap, key: &str) -> u64 { - let revision = revisions.entry(key.to_string()).or_insert(0); +fn next_revision(revisions: &mut HashMap, worker: &RuntimeWorkerRef) -> u64 { + let revision = revisions.entry(worker.clone()).or_insert(0); *revision = revision.saturating_add(1); *revision } -fn worker_key(runtime_id: Option<&str>, worker_id: &str) -> String { - format!("{}:{worker_id}", runtime_id.unwrap_or_default()) -} fn sort_workers(workers: &mut [SubscriptionWorker]) { workers.sort_by(|left, right| { left.runtime_id