diff --git a/Cargo.lock b/Cargo.lock index dda192e6..bd0014b8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6706,6 +6706,7 @@ dependencies = [ name = "workspace-api" version = "0.1.0" dependencies = [ + "protocol", "serde", "serde_json", "ts-rs", diff --git a/crates/client/src/workspace_product.rs b/crates/client/src/workspace_product.rs index fb96a602..c4a41f20 100644 --- a/crates/client/src/workspace_product.rs +++ b/crates/client/src/workspace_product.rs @@ -1,6 +1,6 @@ use reqwest::Method; +use serde::Serialize; use serde::de::DeserializeOwned; -use serde::{Deserialize, Serialize}; use ticket::{ MarkdownText, NewOrchestrationPlanRecord, NewTicket, NewTicketEvent, NewTicketRelation, OrchestrationPlanKind, OrchestrationPlanRecord, Ticket, TicketBackend, TicketDependencyCheck, @@ -9,39 +9,17 @@ use ticket::{ TicketRelationKind, TicketRelationView, TicketStateChange, TicketStateSelector, TicketSummary, }; use workspace_api::{ - ListResponse, ObjectiveCreateRequest, ObjectiveDetail, ObjectiveEditRequest, - ObjectiveLinkTicketRequest, ObjectiveStateRequest, ObjectiveSummary, + BrowserCreateWorkerResponse, BrowserWorkspaceOrchestratorResponse, + CreateWorkspaceWorkerRequest, ListResponse, ObjectiveCreateRequest, ObjectiveDetail, + ObjectiveEditRequest, ObjectiveLinkTicketRequest, ObjectiveStateRequest, ObjectiveSummary, TICKET_ORCHESTRATION_PLANS_QUERY_PATH, TICKET_RELATIONS_QUERY_PATH, + WorkerLaunchOptionsResponse, }; use crate::{BackendApiClient, BackendWorkspaceClientError}; const DEFAULT_PRODUCT_LIST_LIMIT: usize = 1_000; -#[derive(Debug, Deserialize)] -struct BackendWorkerLaunchOptions { - runtimes: Vec, -} - -#[derive(Debug, Deserialize)] -struct BackendWorkerLaunchRuntime { - runtime_id: String, - worker_creation_available: bool, - working_directory_required: bool, -} - -#[derive(Debug, Deserialize)] -struct BackendCreateWorkerResponse { - runtime_id: String, - worker_id: String, -} - -#[derive(Debug, Deserialize)] -struct BackendWorkspaceOrchestratorResponse { - disposition: String, - worker: Option, -} - /// Workspace-scoped Backend client for Ticket and Objective product state. /// /// Construction requires both the selected Backend URL and Workspace identity. @@ -267,7 +245,7 @@ impl BackendWorkspaceProductClient { &self, ticket_id: &str, ) -> Result { - let options: BackendWorkerLaunchOptions = self.get_json("/workers/launch-options")?; + let options: WorkerLaunchOptionsResponse = self.get_json("/workers/launch-options")?; let runtime = options .runtimes .iter() @@ -278,19 +256,19 @@ impl BackendWorkspaceProductClient { .to_string(), ) })?; - let response: BackendCreateWorkerResponse = self.send_json( - Method::POST, - "/workers", - Some(&serde_json::json!({ - "runtime_id": runtime.runtime_id, - "display_name": format!("intake-{ticket_id}"), - "profile": "builtin:intake", - "initial_submit": [{ - "kind": "text", - "content": format!("Please handle intake for Ticket {ticket_id}.") - }] - })), - )?; + let request = CreateWorkspaceWorkerRequest { + runtime_id: runtime.runtime_id.clone(), + display_name: format!("intake-{ticket_id}"), + profile: Some("builtin:intake".to_string()), + ticket_assignment: None, + initial_submit: vec![protocol::Segment::Text { + content: format!("Please handle intake for Ticket {ticket_id}."), + }], + working_directory: None, + control_operation_id: None, + }; + let response: BrowserCreateWorkerResponse = + self.send_json(Method::POST, "/workers", Some(&request))?; Ok(format!( "Started Intake Worker {}/{} for Ticket {ticket_id}", response.runtime_id, response.worker_id @@ -298,7 +276,7 @@ impl BackendWorkspaceProductClient { } pub fn start_workspace_orchestrator(&self) -> Result { - let response: BackendWorkspaceOrchestratorResponse = + let response: BrowserWorkspaceOrchestratorResponse = self.send_json::<(), _>(Method::POST, "/orchestrator", None)?; let worker = response.worker.ok_or_else(|| { BackendWorkspaceClientError::InvalidTarget( @@ -792,11 +770,11 @@ mod tests { let (base_url, requests, handle) = response_sequence_server(vec![ ( "200 OK", - r#"{"runtimes":[{"runtime_id":"embedded","worker_creation_available":true,"working_directory_required":false}]}"#, + r#"{"workspace_id":"workspace-a","runtimes":[{"runtime_id":"embedded","display_name":"Embedded","built_in":true,"worker_creation_available":true,"working_directory_required":false,"status":"connected","diagnostics":[]}],"default_profile":null,"profiles":[],"repositories":[],"working_directories":[],"diagnostics":[]}"#, ), ( "200 OK", - r#"{"runtime_id":"embedded","worker_id":"worker-1"}"#, + r#"{"workspace_id":"workspace-a","runtime_id":"embedded","worker_id":"worker-1","console_href":"/w/workspace-a/workers/worker-1","worker":{"runtime_id":"embedded","worker_id":"worker-1","host_id":"embedded","display_name":"Intake","label":"worker-1","profile":"builtin:intake","singleton_key":null,"tags":[],"workspace":{"visibility":"workspace","identity":"workspace-a","workspace_id":"workspace-a"},"state":"idle","last_seen_at":null,"pinned":false,"retention_state":"active","implementation":{"kind":"runtime","display_hint":"Runtime Worker"},"capabilities":{"can_stop":true,"can_spawn_followup":false},"diagnostics":[]},"diagnostics":[]}"#, ), ]); let client = BackendWorkspaceProductClient::new_with_access_token( @@ -824,7 +802,7 @@ mod tests { #[test] fn workspace_orchestrator_launch_uses_scoped_backend_route() { - let body = r#"{"disposition":"created","worker":{"runtime_id":"embedded","worker_id":"worker-2"}}"#; + let body = r#"{"workspace_id":"workspace-a","online":true,"disposition":"created","worker":{"runtime_id":"embedded","worker_id":"worker-2","host_id":"embedded","display_name":"Orchestrator","label":"worker-2","profile":"builtin:orchestrator","singleton_key":"workspace-orchestrator","tags":[],"workspace":{"visibility":"workspace","identity":"workspace-a","workspace_id":"workspace-a"},"state":"idle","last_seen_at":null,"pinned":true,"retention_state":"active","implementation":{"kind":"runtime","display_hint":"Runtime Worker"},"capabilities":{"can_stop":true,"can_spawn_followup":false},"diagnostics":[]},"diagnostics":[]}"#; let (base_url, request, handle) = one_response_server("200 OK", body); let client = BackendWorkspaceProductClient::new_with_access_token( base_url, diff --git a/crates/merge-request/src/lib.rs b/crates/merge-request/src/lib.rs index cdb6c99d..f5179323 100644 --- a/crates/merge-request/src/lib.rs +++ b/crates/merge-request/src/lib.rs @@ -274,6 +274,12 @@ pub struct RegisterReviewerChildSession { pub reviewer_profile: String, pub now: DateTime, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ReviewSubmissionAuthorization { + pub workspace_id: String, + pub subject_ref: String, +} + #[derive(Debug, Clone)] pub struct SubmitMergeRequestReview { pub ticket_id: String, @@ -535,6 +541,34 @@ impl MergeRequestStore { t.commit()?; Ok(RequestedMergeRequestReview { request_event: e }) } + pub fn authorize_review_submission( + &self, + ticket_id: &str, + capability_token: &str, + ) -> Result { + let connection = self.lock()?; + connection + .query_row( + "SELECT g.workspace_id,g.subject_ref + FROM merge_request_review_grants g + JOIN merge_request_ticket_relations rel + ON rel.workspace_id=g.workspace_id AND rel.merge_request_id=g.merge_request_id + JOIN merge_requests mr + ON mr.workspace_id=g.workspace_id AND mr.merge_request_id=g.merge_request_id + WHERE g.capability_token=?1 AND rel.ticket_id=?2 + AND g.status='issued' AND mr.state='open'", + params![capability_token, ticket_id], + |row| { + Ok(ReviewSubmissionAuthorization { + workspace_id: row.get(0)?, + subject_ref: row.get(1)?, + }) + }, + ) + .optional()? + .ok_or_else(|| MergeRequestError::Unauthorized("review grant invalid".into())) + } + pub fn submit_review( &self, i: SubmitMergeRequestReview, diff --git a/crates/merge-request/tests/store.rs b/crates/merge-request/tests/store.rs index cc1586bd..fe21a67b 100644 --- a/crates/merge-request/tests/store.rs +++ b/crates/merge-request/tests/store.rs @@ -91,6 +91,23 @@ fn approve(s: &MergeRequestStore, subject: &str, token: &str) -> ReviewEvent { }) .unwrap() } +#[test] +fn review_submission_authorization_rejects_invalid_grants_before_side_effects() { + let (_d, store) = fixture(); + open(&store); + request(&store, "published-source", "valid-token"); + + let invalid = store + .authorize_review_submission("T", "invalid-token") + .unwrap_err(); + assert!(matches!(invalid, MergeRequestError::Unauthorized(_))); + let authorized = store + .authorize_review_submission("T", "valid-token") + .unwrap(); + assert_eq!(authorized.workspace_id, "W"); + assert_eq!(authorized.subject_ref, "published-source"); +} + #[test] fn selectors_thread_and_completion_have_no_revision_or_commit_api() { let (d, s) = fixture(); diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index 1dc76b4b..3cc27d01 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -179,6 +179,30 @@ pub struct WorkingDirectoryRequest { pub materialization: Option, } +/// Backend-authorized request to freshly resolve one Repository provider ref. +/// +/// Runtime executes this against the registered source itself rather than a Workdir +/// or Runtime cache. Secret material is fetched through `materialization` and never +/// appears in the result. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RepositoryRefObservationRequest { + pub repository: WorkingDirectoryRepository, + pub selector: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub materialization: Option, +} + +/// Provider-neutral proof of one freshly observed Repository ref. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RepositoryRefObservation { + pub repository_id: String, + pub source_revision: u64, + pub source_fingerprint: String, + pub selector: String, + pub revision_ref: String, + pub observed_at_epoch_seconds: u64, +} + #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct WorkingDirectoryClaim { pub working_directory_id: String, diff --git a/crates/worker-runtime/src/execution.rs b/crates/worker-runtime/src/execution.rs index ff24ecbb..47660977 100644 --- a/crates/worker-runtime/src/execution.rs +++ b/crates/worker-runtime/src/execution.rs @@ -1,4 +1,5 @@ use crate::catalog::{ + RepositoryRefObservation, RepositoryRefObservationRequest, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus, }; use crate::config_bundle::ConfigBundle; @@ -333,6 +334,16 @@ pub trait WorkerExecutionBackend: Send + Sync + 'static { )) } + fn observe_repository_ref( + &self, + _request: &RepositoryRefObservationRequest, + ) -> Result { + Err(WorkingDirectoryDiagnostic::rejected( + "repository_ref_provider_unavailable", + "Worker execution backend does not support Repository ref observation", + )) + } + fn list_working_directories(&self) -> Vec { Vec::new() } @@ -501,6 +512,13 @@ impl WorkerExecutionBackendRef { .authorize_working_directory_repository_access(request) } + pub(crate) fn observe_repository_ref( + &self, + request: &RepositoryRefObservationRequest, + ) -> Result { + self.backend.observe_repository_ref(request) + } + pub(crate) fn list_working_directories(&self) -> Vec { self.backend.list_working_directories() } diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 1e26482e..2217e973 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -11,9 +11,9 @@ use crate::auth::{ verify_capability_token, }; use crate::catalog::{ - ConfigBundleRef, CreateWorkerRequest, WorkerDetail, WorkerLifecycleAck, WorkerSummary, - WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus, - WorkspaceApiRef, + ConfigBundleRef, CreateWorkerRequest, RepositoryRefObservationRequest, WorkerDetail, + WorkerLifecycleAck, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, + WorkingDirectoryRequest, WorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary}; use crate::error::RuntimeError; @@ -208,6 +208,7 @@ fn runtime_http_router_with_optional_auth( "/v1/working-directories/repository-access", post(authorize_working_directory_repository_access), ) + .route("/v1/repository-refs/observe", post(observe_repository_ref)) .route( "/v1/working-directories/{working_directory_id}/sessions", post(open_workdir_session), @@ -583,6 +584,31 @@ async fn authorize_working_directory_repository_access( })) } +async fn observe_repository_ref( + State(state): State, + Extension(auth): Extension, + body: Result, JsonRejection>, +) -> RestResult { + let Json(request) = body.map_err(RuntimeHttpRestError::json_rejection)?; + if request + .materialization + .as_ref() + .is_some_and(|materialization| materialization.workspace_id != auth.workspace_id) + { + return Err(RuntimeHttpRestError::new( + StatusCode::FORBIDDEN, + "repository_ref_observation_workspace_mismatch", + "Repository ref observation authority does not match the authenticated Workspace", + )); + } + let observation = state + .runtime + .observe_repository_ref_from_resource(request) + .await + .map_err(RuntimeHttpRestError::runtime)?; + Ok(Json(observation)) +} + async fn list_working_directories( State(state): State, ) -> RestResult { @@ -1750,7 +1776,10 @@ fn required_runtime_permission(method: &Method, path: &str) -> Option<&'static s if path == "/v1/workers" && *method == Method::POST { return Some("workers:create"); } - if path == "/v1/working-directories/repository-access" && *method == Method::POST { + if (path == "/v1/working-directories/repository-access" + || path == "/v1/repository-refs/observe") + && *method == Method::POST + { return Some("workdirs:operate"); } if path.starts_with("/v1/workdir-sessions") @@ -1941,6 +1970,33 @@ fn status_for_runtime_error(error: &RuntimeError) -> StatusCode { { StatusCode::NOT_FOUND } + RuntimeError::WorkingDirectory(diagnostic) + if matches!( + diagnostic.code.as_str(), + "repository_ref_provider_unavailable" + | "repository_ref_provider_timeout" + | "repository_access_provider_unavailable" + ) => + { + StatusCode::SERVICE_UNAVAILABLE + } + RuntimeError::WorkingDirectory(diagnostic) + if matches!( + diagnostic.code.as_str(), + "repository_ref_provider_auth_failed" + | "repository_access_credential_expired" + | "repository_access_credential_unavailable" + | "repository_access_credential_unauthorized" + | "repository_access_credential_invalid" + ) => + { + StatusCode::FORBIDDEN + } + RuntimeError::WorkingDirectory(diagnostic) + if diagnostic.code == "repository_ref_not_found" => + { + StatusCode::NOT_FOUND + } RuntimeError::RuntimeStopped | RuntimeError::WorkerExecutionUnavailable { .. } | RuntimeError::ExecutionBackendUnavailable { .. } @@ -1951,8 +2007,8 @@ fn status_for_runtime_error(error: &RuntimeError) -> StatusCode { | RuntimeError::InvalidInitialInputKind { .. } | RuntimeError::ConfigBundleDigestMismatch { .. } | RuntimeError::InvalidProfileSelector { .. } - | RuntimeError::UnsupportedConfigDeclaration { .. } - | RuntimeError::WorkingDirectory(_) => StatusCode::BAD_REQUEST, + | RuntimeError::UnsupportedConfigDeclaration { .. } => StatusCode::BAD_REQUEST, + RuntimeError::WorkingDirectory(_) => StatusCode::BAD_REQUEST, RuntimeError::StoreIo { .. } | RuntimeError::StoreMissing { .. } | RuntimeError::StoreCorrupt { .. } @@ -2424,6 +2480,10 @@ mod tests { required_runtime_permission(&Method::POST, "/v1/working-directories/repository-access",), Some("workdirs:operate") ); + assert_eq!( + required_runtime_permission(&Method::POST, "/v1/repository-refs/observe"), + Some("workdirs:operate") + ); assert_eq!( required_runtime_permission(&Method::POST, "/v1/working-directories/wd-1/sessions"), Some("workdirs:operate") @@ -2934,17 +2994,25 @@ mod tests { #[test] fn workdir_runtime_errors_preserve_diagnostic_code() { - let error = - RuntimeError::WorkingDirectory(crate::working_directory::WorkingDirectoryDiagnostic { - code: "working_directory_not_found".to_string(), - message: "working directory missing-workdir was not found".to_string(), - }); - - assert_eq!(status_for_runtime_error(&error), StatusCode::NOT_FOUND); - assert_eq!( - code_for_runtime_error(&error), - "working_directory_not_found" - ); + let cases = [ + ("working_directory_not_found", StatusCode::NOT_FOUND), + ( + "repository_ref_provider_timeout", + StatusCode::SERVICE_UNAVAILABLE, + ), + ("repository_ref_provider_auth_failed", StatusCode::FORBIDDEN), + ("repository_ref_not_found", StatusCode::NOT_FOUND), + ]; + for (code, expected_status) in cases { + let error = RuntimeError::WorkingDirectory( + crate::working_directory::WorkingDirectoryDiagnostic { + code: code.to_string(), + message: "bounded diagnostic".to_string(), + }, + ); + assert_eq!(status_for_runtime_error(&error), expected_status); + assert_eq!(code_for_runtime_error(&error), code); + } } } diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index ea3f28b4..8997d86f 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -1,8 +1,8 @@ use crate::catalog::{ - ConfigBundleRef, CreateWorkerRequest, ProfileSelector, WorkerDetail, WorkerLifecycleAck, - WorkerRestoreIntent, WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, - WorkingDirectoryRequest, WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, - WorkspaceApiRef, + ConfigBundleRef, CreateWorkerRequest, ProfileSelector, RepositoryRefObservation, + RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck, WorkerRestoreIntent, + WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, + WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary, validate_config_bundle, @@ -380,6 +380,38 @@ impl Runtime { self.create_working_directory(request) } + pub async fn observe_repository_ref_from_resource( + &self, + mut request: RepositoryRefObservationRequest, + ) -> Result { + if let Some(ssh) = request + .materialization + .as_mut() + .and_then(|materialization| materialization.ssh.as_mut()) + { + self.resolve_repository_access_resource(ssh).await?; + } + self.observe_repository_ref(request) + } + + pub fn observe_repository_ref( + &self, + request: RepositoryRefObservationRequest, + ) -> Result { + let backend = { + let state = self.lock()?; + state.ensure_running()?; + state.execution_backend.clone().ok_or_else(|| { + RuntimeError::ExecutionBackendUnavailable { + message: "Repository ref observation requires an execution backend".to_string(), + } + })? + }; + backend + .observe_repository_ref(&request) + .map_err(RuntimeError::from) + } + pub fn authorize_working_directory_repository_access( &self, request: WorkingDirectoryRepositoryAccessRequest, @@ -3036,20 +3068,35 @@ fn worker_status_from_run_state(run_state: WorkerExecutionRunState) -> WorkerSta } fn repository_resource_error(error: BackendResourceError) -> RuntimeError { - let category = match error { - BackendResourceError::Expired => "expired", - BackendResourceError::Unauthorized { .. } => "unauthorized", - BackendResourceError::UnsupportedKind => "unsupported_kind", - BackendResourceError::MissingResource => "missing_resource", - BackendResourceError::Oversized { .. } => "oversized", - BackendResourceError::DigestMismatch { .. } => "digest_mismatch", - BackendResourceError::ContentTypeMismatch { .. } => "content_type_mismatch", - BackendResourceError::InvalidResponse { .. } => "invalid_response", - BackendResourceError::Transport { .. } => "transport", + let (code, message) = match error { + BackendResourceError::Expired => ( + "repository_access_credential_expired", + "Repository access credential lease expired", + ), + BackendResourceError::Unauthorized { .. } => ( + "repository_access_credential_unauthorized", + "Repository access credential lease was rejected", + ), + BackendResourceError::MissingResource => ( + "repository_access_credential_unavailable", + "Repository access credential lease is unavailable or already consumed", + ), + BackendResourceError::Transport { .. } => ( + "repository_access_provider_unavailable", + "Repository access credential provider is unavailable", + ), + BackendResourceError::UnsupportedKind + | BackendResourceError::Oversized { .. } + | BackendResourceError::DigestMismatch { .. } + | BackendResourceError::ContentTypeMismatch { .. } + | BackendResourceError::InvalidResponse { .. } => ( + "repository_access_credential_invalid", + "Repository access credential response is invalid", + ), }; - RuntimeError::InvalidRequest(format!( - "Backend Repository SSH access resource fetch failed: {category}" - )) + RuntimeError::WorkingDirectory( + crate::working_directory::WorkingDirectoryDiagnostic::rejected(code, message), + ) } fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest { @@ -3250,6 +3297,33 @@ mod tests { use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; + #[test] + fn repository_resource_failures_keep_typed_credential_diagnostics() { + let cases = [ + ( + BackendResourceError::Expired, + "repository_access_credential_expired", + ), + ( + BackendResourceError::MissingResource, + "repository_access_credential_unavailable", + ), + ( + BackendResourceError::Unauthorized { + message: "denied".to_string(), + }, + "repository_access_credential_unauthorized", + ), + ]; + for (error, expected_code) in cases { + let RuntimeError::WorkingDirectory(diagnostic) = repository_resource_error(error) + else { + panic!("Repository resource failure lost its typed diagnostic") + }; + assert_eq!(diagnostic.code, expected_code); + } + } + fn internal_worker_ref( session_id: &str, parent_session_id: Option<&str>, diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index ef3da9bb..15d34050 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -20,6 +20,7 @@ use crate::auth::{ }; use crate::catalog::{ CreateWorkerRequest, ProfileSourceArchiveHttpRef, ProfileSourceArchiveSource, + RepositoryRefObservation, RepositoryRefObservationRequest, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus, }; use crate::execution::{ @@ -1645,6 +1646,19 @@ where materializer.authorize_repository_access(request) } + fn observe_repository_ref( + &self, + request: &RepositoryRefObservationRequest, + ) -> Result { + let materializer = self.working_directory_materializer.as_ref().ok_or_else(|| { + WorkingDirectoryDiagnostic::rejected( + "repository_ref_provider_unavailable", + "Repository ref observation requested, but no materializer is configured for this Runtime backend", + ) + })?; + materializer.observe_repository_ref(request) + } + fn list_working_directories(&self) -> Vec { self.working_directory_materializer .as_ref() diff --git a/crates/worker-runtime/src/working_directory.rs b/crates/worker-runtime/src/working_directory.rs index 06702830..797ec679 100644 --- a/crates/worker-runtime/src/working_directory.rs +++ b/crates/worker-runtime/src/working_directory.rs @@ -1,5 +1,6 @@ use crate::catalog::{ - MaterializerKind, RepositorySshMaterializationAccess, WorkingDirectoryCleanupTarget, + MaterializerKind, RepositoryRefObservation, RepositoryRefObservationRequest, + RepositorySshMaterializationAccess, WorkingDirectoryCleanupTarget, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus, WorkingDirectoryStatusKind, WorkingDirectorySummary, }; @@ -196,6 +197,11 @@ pub trait WorkingDirectoryMaterializer: Send + Sync + 'static { request: &WorkingDirectoryRepositoryAccessRequest, ) -> Result<(), WorkingDirectoryDiagnostic>; + fn observe_repository_ref( + &self, + request: &RepositoryRefObservationRequest, + ) -> Result; + fn bind_working_directory( &self, working_directory_id: &str, @@ -943,6 +949,72 @@ impl WorkingDirectoryMaterializer for RuntimeGitCacheMaterializer { self.cache_repository_access(&request.working_directory_id, ssh) } + fn observe_repository_ref( + &self, + request: &RepositoryRefObservationRequest, + ) -> Result { + let selector = request.selector.trim(); + validate_exact_branch_selector(selector)?; + + let working_request = WorkingDirectoryRequest { + repository: request.repository.clone(), + materializer: MaterializerKind::RuntimeGitCache, + backend_workdir_id: None, + materialization: request.materialization.clone(), + }; + Self::validate_request(&working_request)?; + let access = RepositoryCommandAccess::prepare(&self.runtime_root, &working_request)?; + let mut command = repository_git_command(&working_request, access.as_ref()); + command.args([ + "ls-remote", + "--exit-code", + "--refs", + request.repository.source.uri.as_str(), + selector, + ]); + let output = run_repository_git_stdout(command, request.repository.source.kind)?; + let mut lines = output.lines(); + let line = lines.next().ok_or_else(|| { + WorkingDirectoryDiagnostic::new( + "repository_ref_not_found", + "Repository provider did not return the requested ref", + ) + })?; + if lines.next().is_some() { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_response_invalid", + "Repository provider returned an ambiguous ref observation", + )); + } + let (revision_ref, observed_selector) = line.split_once('\t').ok_or_else(|| { + WorkingDirectoryDiagnostic::new( + "repository_ref_response_invalid", + "Repository provider returned an invalid ref observation", + ) + })?; + if observed_selector != selector + || !matches!(revision_ref.len(), 40 | 64) + || !revision_ref.bytes().all(|byte| byte.is_ascii_hexdigit()) + { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_response_invalid", + "Repository provider returned an invalid ref observation", + )); + } + + Ok(RepositoryRefObservation { + repository_id: request.repository.id.clone(), + source_revision: request.repository.source_revision, + source_fingerprint: request.repository.source_fingerprint.clone(), + selector: selector.to_string(), + revision_ref: revision_ref.to_ascii_lowercase(), + observed_at_epoch_seconds: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs(), + }) + } + fn bind_working_directory( &self, working_directory_id: &str, @@ -2022,6 +2094,146 @@ fn repository_git_command( command } +fn validate_exact_branch_selector(selector: &str) -> Result<(), WorkingDirectoryDiagnostic> { + validate_selector(selector).map_err(|_| { + WorkingDirectoryDiagnostic::new( + "repository_ref_selector_invalid", + "Repository ref observation requires a valid exact branch selector", + ) + })?; + if !selector.starts_with("refs/heads/") { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_selector_invalid", + "Repository ref observation requires an exact branch selector", + )); + } + let status = Command::new("git") + .args(["check-ref-format", selector]) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .map_err(|_| { + WorkingDirectoryDiagnostic::new( + "repository_ref_provider_unavailable", + "Git ref validation could not be started", + ) + })?; + if !status.success() { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_selector_invalid", + "Repository ref observation requires a valid exact branch selector", + )); + } + Ok(()) +} + +fn read_bounded_command_output(mut reader: impl Read) -> Vec { + const MAX_CAPTURE_BYTES: usize = 8192; + let mut captured = Vec::new(); + let mut chunk = [0_u8; 4096]; + loop { + match reader.read(&mut chunk) { + Ok(0) | Err(_) => break, + Ok(read) => { + let remaining = MAX_CAPTURE_BYTES.saturating_sub(captured.len()); + captured.extend_from_slice(&chunk[..read.min(remaining)]); + } + } + } + captured +} + +fn run_repository_git_stdout( + mut command: Command, + source_kind: workspace_api::RepositorySourceKind, +) -> Result { + command + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let mut child = command.spawn().map_err(|_| { + WorkingDirectoryDiagnostic::new( + "repository_ref_provider_unavailable", + "Repository provider operation could not be started", + ) + })?; + let stdout = child.stdout.take().ok_or_else(|| { + WorkingDirectoryDiagnostic::new( + "repository_ref_provider_unavailable", + "Repository provider response could not be captured", + ) + })?; + let stderr = child.stderr.take().ok_or_else(|| { + WorkingDirectoryDiagnostic::new( + "repository_ref_provider_unavailable", + "Repository provider diagnostic could not be captured", + ) + })?; + let stdout_reader = std::thread::spawn(move || read_bounded_command_output(stdout)); + let stderr_reader = std::thread::spawn(move || read_bounded_command_output(stderr)); + let started = Instant::now(); + let status = loop { + if let Some(status) = child.try_wait().map_err(|_| { + WorkingDirectoryDiagnostic::new( + "repository_ref_provider_unavailable", + "Repository provider operation status could not be observed", + ) + })? { + break status; + } + if started.elapsed() >= REPOSITORY_COMMAND_TIMEOUT { + let _ = child.kill(); + let _ = child.wait(); + let _ = stdout_reader.join(); + let _ = stderr_reader.join(); + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_provider_timeout", + "Repository provider operation exceeded the Runtime time limit", + )); + } + std::thread::sleep(Duration::from_millis(25)); + }; + let stdout = stdout_reader.join().unwrap_or_default(); + let stderr = stderr_reader.join().unwrap_or_default(); + if status.success() { + return String::from_utf8(stdout).map_err(|_| { + WorkingDirectoryDiagnostic::new( + "repository_ref_response_invalid", + "Repository provider returned a non-UTF-8 ref observation", + ) + }); + } + if status.code() == Some(2) { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_not_found", + "Repository provider did not return the requested ref", + )); + } + let diagnostic = String::from_utf8_lossy(&stderr).to_ascii_lowercase(); + let auth_failed = source_kind.is_remote() + && [ + "authentication failed", + "permission denied", + "could not read username", + "publickey", + ] + .iter() + .any(|marker| diagnostic.contains(marker)); + Err(WorkingDirectoryDiagnostic::new( + if auth_failed { + "repository_ref_provider_auth_failed" + } else { + "repository_ref_provider_unavailable" + }, + if auth_failed { + "Repository provider rejected the operation-scoped authentication" + } else { + "Repository provider operation failed" + }, + )) +} + fn run_repository_git( mut command: Command, code: &'static str, @@ -2493,6 +2705,157 @@ mod tests { WorkerRef::new(WorkerId::from_legacy_u64(sequence)) } + #[test] + fn repository_ref_observation_reads_the_provider_fresh() { + let repo = create_clean_repo(); + git(repo.path(), &["branch", "published"]); + let runtime_root = tempfile::tempdir().unwrap(); + let materializer = RuntimeGitCacheMaterializer::new(runtime_root.path()); + let repository = request(repo.path()).repository; + let observation_request = RepositoryRefObservationRequest { + repository, + selector: "refs/heads/published".to_string(), + materialization: None, + }; + + let first = materializer + .observe_repository_ref(&observation_request) + .unwrap(); + assert_eq!( + first.revision_ref, + git_stdout(repo.path(), ["rev-parse", "published"]).unwrap() + ); + fs::write(repo.path().join("second.txt"), "second\n").unwrap(); + git(repo.path(), &["add", "second.txt"]); + git(repo.path(), &["commit", "-m", "second"]); + git(repo.path(), &["branch", "-f", "published"]); + + let second = materializer + .observe_repository_ref(&observation_request) + .unwrap(); + assert_ne!(first.revision_ref, second.revision_ref); + assert_eq!( + second.revision_ref, + git_stdout(repo.path(), ["rev-parse", "published"]).unwrap() + ); + } + + #[test] + fn repository_ref_observation_ignores_unpublished_and_stale_workdir_or_cache_refs() { + let seed = create_clean_repo(); + let layout = tempfile::tempdir().unwrap(); + let provider = layout.path().join("provider.git"); + git( + layout.path(), + &[ + "clone", + "--bare", + seed.path().to_str().unwrap(), + provider.to_str().unwrap(), + ], + ); + let cache = layout.path().join("cache"); + git( + layout.path(), + &["clone", provider.to_str().unwrap(), cache.to_str().unwrap()], + ); + let workdir = layout.path().join("workdir"); + git( + layout.path(), + &[ + "clone", + provider.to_str().unwrap(), + workdir.to_str().unwrap(), + ], + ); + git(&workdir, &["config", "user.name", "Yoi Test"]); + git(&workdir, &["config", "user.email", "yoi@example.com"]); + git(&workdir, &["switch", "-c", "published-source"]); + fs::write(workdir.join("source.txt"), "first\n").unwrap(); + git(&workdir, &["add", "source.txt"]); + git(&workdir, &["commit", "-m", "source first"]); + + let runtime_root = tempfile::tempdir().unwrap(); + let materializer = RuntimeGitCacheMaterializer::new(runtime_root.path()); + let repository = request(&provider).repository; + let observation_request = RepositoryRefObservationRequest { + repository, + selector: "refs/heads/published-source".to_string(), + materialization: None, + }; + assert_eq!( + materializer + .observe_repository_ref(&observation_request) + .unwrap_err() + .code, + "repository_ref_not_found" + ); + + git( + &workdir, + &["push", "origin", "HEAD:refs/heads/published-source"], + ); + let first = materializer + .observe_repository_ref(&observation_request) + .unwrap(); + fs::write(workdir.join("source.txt"), "second\n").unwrap(); + git(&workdir, &["add", "source.txt"]); + git(&workdir, &["commit", "-m", "source second"]); + let unpublished_second = git_stdout(&workdir, ["rev-parse", "HEAD"]).unwrap(); + let still_first = materializer + .observe_repository_ref(&observation_request) + .unwrap(); + assert_eq!(still_first.revision_ref, first.revision_ref); + assert_ne!(still_first.revision_ref, unpublished_second); + + git( + &workdir, + &["push", "origin", "HEAD:refs/heads/published-source"], + ); + let second = materializer + .observe_repository_ref(&observation_request) + .unwrap(); + assert_eq!(second.revision_ref, unpublished_second); + assert_ne!(second.revision_ref, first.revision_ref); + assert_ne!( + git_stdout(&cache, ["rev-parse", "HEAD"]).unwrap(), + second.revision_ref + ); + } + + #[test] + fn repository_ref_observation_rejects_missing_and_non_branch_selectors() { + let repo = create_clean_repo(); + let runtime_root = tempfile::tempdir().unwrap(); + let materializer = RuntimeGitCacheMaterializer::new(runtime_root.path()); + let repository = request(repo.path()).repository; + + let missing = materializer + .observe_repository_ref(&RepositoryRefObservationRequest { + repository: repository.clone(), + selector: "refs/heads/not-published".to_string(), + materialization: None, + }) + .unwrap_err(); + assert_eq!(missing.code, "repository_ref_not_found"); + let non_branch = materializer + .observe_repository_ref(&RepositoryRefObservationRequest { + repository: repository.clone(), + selector: "HEAD".to_string(), + materialization: None, + }) + .unwrap_err(); + assert_eq!(non_branch.code, "repository_ref_selector_invalid"); + let wildcard = materializer + .observe_repository_ref(&RepositoryRefObservationRequest { + repository, + selector: "refs/heads/release/*".to_string(), + materialization: None, + }) + .unwrap_err(); + assert_eq!(wildcard.code, "repository_ref_selector_invalid"); + } + #[test] fn local_git_repo_materializes_detached_worktree_under_runtime_root() { let repo = create_clean_repo(); diff --git a/crates/workspace-api/Cargo.toml b/crates/workspace-api/Cargo.toml index e00b8a58..979231a2 100644 --- a/crates/workspace-api/Cargo.toml +++ b/crates/workspace-api/Cargo.toml @@ -7,9 +7,10 @@ publish = false [features] default = [] -typescript = ["dep:ts-rs"] +typescript = ["dep:ts-rs", "protocol/typescript"] [dependencies] +protocol.workspace = true serde = { workspace = true, features = ["derive"] } ts-rs = { version = "12.0.1", optional = true } @@ -24,6 +25,10 @@ serde_json.workspace = true name = "generate_workdir_api_types" required-features = ["typescript"] +[[example]] +name = "generate_worker_launch_api_types" +required-features = ["typescript"] + [[example]] name = "generate_companion_api_types" required-features = ["typescript"] diff --git a/crates/workspace-api/examples/generate_worker_launch_api_types.rs b/crates/workspace-api/examples/generate_worker_launch_api_types.rs new file mode 100644 index 00000000..ffcf6248 --- /dev/null +++ b/crates/workspace-api/examples/generate_worker_launch_api_types.rs @@ -0,0 +1,3 @@ +fn main() { + print!("{}", workspace_api::worker_launch_api_typescript()); +} diff --git a/crates/workspace-api/src/lib.rs b/crates/workspace-api/src/lib.rs index 14bfb26a..f6aa4e78 100644 --- a/crates/workspace-api/src/lib.rs +++ b/crates/workspace-api/src/lib.rs @@ -579,6 +579,8 @@ pub struct WorkingDirectoryOccupancy { /// retains the Backend-generated Repository id and is never a Workspace public /// projection. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[cfg_attr(feature = "typescript", ts(optional_fields = nullable))] #[serde(deny_unknown_fields)] pub struct RuntimeWorkingDirectoryCleanupTarget { pub kind: String, @@ -590,6 +592,8 @@ pub struct RuntimeWorkingDirectoryCleanupTarget { /// surfaces must project this through [`WorkingDirectorySummary`] so the UUID is /// replaced with `repository_key`. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[cfg_attr(feature = "typescript", ts(optional_fields = nullable))] #[serde(deny_unknown_fields)] pub struct RuntimeWorkingDirectorySummary { pub working_directory_id: String, @@ -607,6 +611,7 @@ pub struct RuntimeWorkingDirectorySummary { #[serde(default, skip_serializing_if = "Option::is_none")] pub current_tree: Option, #[serde(default, skip_serializing_if = "Option::is_none")] + #[cfg_attr(feature = "typescript", ts(type = "number | null"))] pub observed_at_epoch_seconds: Option, pub materializer_kind: WorkingDirectoryMaterializerKind, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -908,6 +913,7 @@ pub struct RuntimeConnectionTestResponse { } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct WorkerWorkspaceSummary { pub visibility: String, pub identity: String, @@ -916,12 +922,14 @@ pub struct WorkerWorkspaceSummary { } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct WorkerImplementationSummary { pub kind: String, pub display_hint: String, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct WorkerCapabilitySummary { pub can_stop: bool, pub can_spawn_followup: bool, @@ -1095,6 +1103,7 @@ pub struct WorkspaceWorkerDiscoveryPage { /// do not carry one. The Workspace Server must resolve it from Workspace /// authority before constructing this response. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct WorkerSummary { pub runtime_id: String, pub worker_id: String, @@ -1122,6 +1131,142 @@ pub struct WorkerSummary { pub diagnostics: Vec, } +/// Runtime-owned Worker summary embedded in Worker launch responses. +/// +/// This preserves the existing launch wire shape. Workspace-owned Worker list +/// and detail responses use [`WorkerSummary`] instead. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct WorkerLaunchWorkerSummary { + pub runtime_id: String, + pub worker_id: String, + pub host_id: String, + pub display_name: String, + pub label: String, + pub profile: Option, + pub singleton_key: Option, + pub tags: Vec, + pub workspace: WorkerWorkspaceSummary, + pub state: String, + pub last_seen_at: Option, + pub pinned: bool, + pub retention_state: String, + pub implementation: WorkerImplementationSummary, + pub capabilities: WorkerCapabilitySummary, + #[serde(default, skip_serializing_if = "Option::is_none")] + #[cfg_attr(feature = "typescript", ts(optional = nullable))] + pub working_directory: Option, + #[serde(default)] + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct WorkerLaunchOptionsResponse { + pub workspace_id: String, + pub runtimes: Vec, + pub default_profile: Option, + pub profiles: Vec, + pub repositories: Vec, + pub working_directories: Vec, + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct WorkerLaunchRuntimeOption { + pub runtime_id: String, + pub display_name: String, + pub built_in: bool, + pub worker_creation_available: bool, + pub working_directory_required: bool, + pub status: String, + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct WorkerLaunchProfileCandidate { + pub id: String, + pub label: String, + pub description: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct WorkingDirectoryRepositoryOption { + pub repository_key: String, + #[serde(skip_serializing_if = "Option::is_none")] + #[cfg_attr(feature = "typescript", ts(optional = nullable))] + pub default_selector: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct BrowserWorkerWorkingDirectorySelection { + pub working_directory_id: String, + #[serde(default)] + pub relative_cwd: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct CreateWorkspaceWorkerTicketAssignmentRequest { + pub ticket_id: String, + pub operation_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct CreateWorkspaceWorkerRequest { + pub runtime_id: String, + pub display_name: String, + #[serde(default)] + pub profile: Option, + #[serde(default)] + pub ticket_assignment: Option, + #[serde(default)] + pub initial_submit: Vec, + #[serde(default)] + pub working_directory: Option, + /// Backend idempotency key used only for authenticated Worker-owned spawn/control. + #[serde(default)] + pub control_operation_id: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct BrowserCreateWorkerResponse { + pub workspace_id: String, + pub runtime_id: String, + pub worker_id: String, + pub console_href: String, + pub worker: WorkerLaunchWorkerSummary, + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(deny_unknown_fields)] +pub struct BrowserWorkspaceOrchestratorResponse { + pub workspace_id: String, + pub online: bool, + pub disposition: String, + #[serde(skip_serializing_if = "Option::is_none")] + #[cfg_attr(feature = "typescript", ts(optional = nullable))] + pub worker: Option, + pub diagnostics: Vec, +} + #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum WorkerOperationState { @@ -1390,6 +1535,74 @@ pub fn workdir_api_typescript() -> String { ) } +#[cfg(feature = "typescript")] +pub fn worker_launch_api_typescript() -> String { + use ts_rs::TS; + + let config = ts_rs::Config::default(); + let declarations = [ + DiagnosticSeverity::decl(&config), + Diagnostic::decl(&config), + WorkingDirectoryMaterializerKind::decl(&config), + WorkingDirectoryStatusKind::decl(&config), + WorkingDirectoryCleanupTarget::decl(&config), + RuntimeWorkingDirectoryCleanupTarget::decl(&config), + RuntimeWorkingDirectorySummary::decl(&config), + WorkingDirectoryOccupancy::decl(&config), + WorkingDirectorySummary::decl(&config), + WorkerWorkspaceSummary::decl(&config), + WorkerImplementationSummary::decl(&config), + WorkerCapabilitySummary::decl(&config), + WorkerLaunchWorkerSummary::decl(&config), + WorkerLaunchRuntimeOption::decl(&config), + WorkerLaunchProfileCandidate::decl(&config), + WorkingDirectoryRepositoryOption::decl(&config), + WorkerLaunchOptionsResponse::decl(&config), + BrowserWorkerWorkingDirectorySelection::decl(&config), + CreateWorkspaceWorkerTicketAssignmentRequest::decl(&config), + CreateWorkspaceWorkerRequest::decl(&config), + BrowserCreateWorkerResponse::decl(&config), + BrowserWorkspaceOrchestratorResponse::decl(&config), + ]; + format!( + "// Generated from workspace-api. Do not edit by hand.\n// Regenerate: cargo run -q -p workspace-api --features typescript --example generate_worker_launch_api_types > web/workspace/src/lib/generated/worker-launch-api.ts\n\nimport type {{ Segment }} from \"./protocol\";\n\n{}\n", + declarations + .into_iter() + .map(|declaration| format!("export {declaration}")) + .collect::>() + .join("\n\n") + ) +} + +#[cfg(all(test, feature = "typescript"))] +mod worker_launch_typescript_tests { + #[test] + fn generated_worker_launch_api_contract_is_current() { + let expected = super::worker_launch_api_typescript(); + let path = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) + .join("../../web/workspace/src/lib/generated/worker-launch-api.ts"); + let actual = std::fs::read_to_string(&path) + .unwrap_or_else(|error| panic!("failed to read {}: {error}", path.display())); + assert_eq!( + normalize(&actual), + normalize(&expected), + "regenerate Worker launch API TypeScript types with `cargo run -q -p workspace-api --features typescript --example generate_worker_launch_api_types > web/workspace/src/lib/generated/worker-launch-api.ts` and format the generated file", + ); + } + + fn normalize(value: &str) -> String { + value + .chars() + .filter_map(|character| match character { + character if character.is_whitespace() => None, + ',' => Some(';'), + character => Some(character), + }) + .collect::() + .replace("=|", "=") + } +} + #[cfg(all(test, feature = "typescript"))] mod workdir_typescript_tests { #[test] @@ -1423,6 +1636,113 @@ mod workdir_typescript_tests { mod tests { use super::*; + fn worker_launch_summary() -> WorkerLaunchWorkerSummary { + WorkerLaunchWorkerSummary { + runtime_id: "runtime-a".to_string(), + worker_id: "worker-a".to_string(), + host_id: "host-a".to_string(), + display_name: "Worker A".to_string(), + label: "worker-a".to_string(), + profile: None, + singleton_key: None, + tags: Vec::new(), + workspace: WorkerWorkspaceSummary { + visibility: "workspace".to_string(), + identity: "workspace-a".to_string(), + workspace_id: Some("workspace-a".to_string()), + }, + state: "idle".to_string(), + last_seen_at: None, + pinned: false, + retention_state: "active".to_string(), + implementation: WorkerImplementationSummary { + kind: "runtime".to_string(), + display_hint: "Runtime Worker".to_string(), + }, + capabilities: WorkerCapabilitySummary { + can_stop: true, + can_spawn_followup: false, + }, + working_directory: None, + diagnostics: Vec::new(), + } + } + + #[test] + fn worker_launch_optional_omission_and_request_shape_are_stable() { + assert_eq!( + serde_json::to_value(WorkingDirectoryRepositoryOption { + repository_key: "main".to_string(), + default_selector: None, + }) + .unwrap(), + serde_json::json!({ "repository_key": "main" }) + ); + + let orchestrator = serde_json::to_value(BrowserWorkspaceOrchestratorResponse { + workspace_id: "workspace-a".to_string(), + online: false, + disposition: "unavailable".to_string(), + worker: None, + diagnostics: Vec::new(), + }) + .unwrap(); + assert_eq!( + orchestrator, + serde_json::json!({ + "workspace_id": "workspace-a", + "online": false, + "disposition": "unavailable", + "diagnostics": [], + }) + ); + + let worker = serde_json::to_value(worker_launch_summary()).unwrap(); + assert!( + !worker + .as_object() + .unwrap() + .contains_key("working_directory") + ); + assert_eq!(worker["profile"], serde_json::Value::Null); + assert_eq!(worker["singleton_key"], serde_json::Value::Null); + assert_eq!(worker["last_seen_at"], serde_json::Value::Null); + + let request = serde_json::to_value(CreateWorkspaceWorkerRequest { + runtime_id: "runtime-a".to_string(), + display_name: "Worker A".to_string(), + profile: None, + ticket_assignment: None, + initial_submit: Vec::new(), + working_directory: None, + control_operation_id: None, + }) + .unwrap(); + assert_eq!( + request, + serde_json::json!({ + "runtime_id": "runtime-a", + "display_name": "Worker A", + "profile": null, + "ticket_assignment": null, + "initial_submit": [], + "working_directory": null, + "control_operation_id": null, + }) + ); + } + + #[test] + fn worker_launch_request_rejects_unknown_fields() { + let error = serde_json::from_value::(serde_json::json!({ + "runtime_id": "runtime-a", + "display_name": "Worker A", + "unexpected": true, + })) + .unwrap_err(); + assert!(error.to_string().contains("unknown field")); + } + #[test] fn repository_key_validation_is_canonical_and_bounded() { let max = "a".repeat(64); diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index c378ce5d..7f530a8d 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -23,10 +23,10 @@ use worker_runtime::RuntimeWorkspaceScope; use worker_runtime::auth::{CapabilityTokenSigner, capability_claims}; use worker_runtime::catalog::{ ConfigBundleRef, CreateWorkerRequest, ProfileSelector, ProfileSourceArchiveHttpRef, - ProfileSourceArchiveSource, WorkerDetail as EmbeddedWorkerDetail, - WorkerStatus as EmbeddedWorkerStatus, WorkingDirectoryClaim, - WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus, - WorkingDirectorySummary, WorkspaceApiRef, + ProfileSourceArchiveSource, RepositoryRefObservation, RepositoryRefObservationRequest, + WorkerDetail as EmbeddedWorkerDetail, WorkerStatus as EmbeddedWorkerStatus, + WorkingDirectoryClaim, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, + WorkingDirectoryStatus, WorkingDirectorySummary, WorkspaceApiRef, }; use worker_runtime::config_bundle::{ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary}; #[cfg(test)] @@ -830,6 +830,17 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { )) } + fn observe_repository_ref( + &self, + _request: RepositoryRefObservationRequest, + ) -> std::result::Result { + Err(Error::RuntimeOperationFailed { + runtime_id: self.runtime_id().to_string(), + code: "repository_ref_provider_unavailable".to_string(), + message: "Runtime does not support Repository ref observation".to_string(), + }) + } + fn list_working_directories(&self) -> RuntimeList { RuntimeList::new(Vec::new(), Vec::new()) } @@ -1449,6 +1460,31 @@ impl RuntimeRegistry { }) } + pub fn observe_repository_ref( + &self, + runtime_id: &str, + request: RepositoryRefObservationRequest, + ) -> Result { + validate_backend_identifier("runtime_id", runtime_id)?; + let runtime = self.runtime(runtime_id)?; + runtime + .observe_repository_ref(request) + .map_err(|error| match error { + Error::RuntimeOperationFailed { code, message, .. } => { + RuntimeRegistryError::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code, + message, + } + } + other => RuntimeRegistryError::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "repository_ref_provider_unavailable".to_string(), + message: other.to_string(), + }, + }) + } + pub fn list_working_directories( &self, runtime_id: &str, @@ -2143,6 +2179,28 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { } } + fn observe_repository_ref( + &self, + request: RepositoryRefObservationRequest, + ) -> std::result::Result { + self.runtime + .observe_repository_ref(request) + .map_err(|error| match error { + worker_runtime::error::RuntimeError::WorkingDirectory(diagnostic) => { + Error::RuntimeOperationFailed { + runtime_id: self.runtime_id.clone(), + code: diagnostic.code, + message: diagnostic.message, + } + } + error => Error::RuntimeOperationFailed { + runtime_id: self.runtime_id.clone(), + code: "repository_ref_provider_unavailable".to_string(), + message: error.to_string(), + }, + }) + } + fn list_working_directories(&self) -> RuntimeList { RuntimeList::new(Vec::new(), Vec::new()) } @@ -3355,6 +3413,18 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { .map_err(|diagnostic| Error::RegistryInconsistency(diagnostic.message)) } + fn observe_repository_ref( + &self, + request: RepositoryRefObservationRequest, + ) -> std::result::Result { + self.post_json::<_, RepositoryRefObservation>("/v1/repository-refs/observe", &request) + .map_err(|diagnostic| Error::RuntimeOperationFailed { + runtime_id: self.runtime_id.clone(), + code: diagnostic.code, + message: diagnostic.message, + }) + } + fn list_working_directories(&self) -> RuntimeList { match self.get_json::("/v1/working-directories") { Ok(response) => RuntimeList::new(response.working_directories, Vec::new()), @@ -4789,6 +4859,85 @@ mod tests { } } + struct ObservingExecutionBackend { + response: RepositoryRefObservation, + observed: Arc>>, + } + + impl worker_runtime::execution::WorkerExecutionBackend for ObservingExecutionBackend { + fn backend_id(&self) -> &str { + "repository-observation-test-backend" + } + + fn spawn_worker( + &self, + _request: worker_runtime::execution::WorkerExecutionSpawnRequest, + ) -> worker_runtime::execution::WorkerExecutionSpawnResult { + unreachable!("Repository observation test does not spawn Workers") + } + + fn dispatch_input( + &self, + _handle: &worker_runtime::execution::WorkerExecutionHandle, + _input: EmbeddedWorkerInput, + ) -> worker_runtime::execution::WorkerExecutionResult { + unreachable!("Repository observation test does not dispatch Worker input") + } + + fn observe_repository_ref( + &self, + request: &RepositoryRefObservationRequest, + ) -> Result< + RepositoryRefObservation, + worker_runtime::working_directory::WorkingDirectoryDiagnostic, + > { + self.observed.lock().unwrap().push(request.clone()); + Ok(self.response.clone()) + } + } + + #[test] + fn embedded_runtime_forwards_repository_ref_observation_to_execution_backend() { + let observed = Arc::new(Mutex::new(Vec::new())); + let expected = RepositoryRefObservation { + repository_id: "repository-1".to_string(), + source_revision: 7, + source_fingerprint: "sha256:source".to_string(), + selector: "refs/heads/published".to_string(), + revision_ref: "0123456789012345678901234567890123456789".to_string(), + observed_at_epoch_seconds: 42, + }; + let runtime = EmbeddedWorkerRuntime::new_memory_with_execution_backend( + "workspace-test", + Arc::new(ObservingExecutionBackend { + response: expected.clone(), + observed: observed.clone(), + }), + ) + .unwrap(); + let request = RepositoryRefObservationRequest { + repository: worker_runtime::catalog::WorkingDirectoryRepository { + id: "repository-1".to_string(), + provider: "git".to_string(), + source: workspace_api::RepositorySource { + kind: workspace_api::RepositorySourceKind::LocalPath, + uri: "/provider/repository.git".to_string(), + }, + source_revision: 7, + source_fingerprint: "sha256:source".to_string(), + selector: None, + }, + selector: "refs/heads/published".to_string(), + materialization: None, + }; + + assert_eq!( + runtime.observe_repository_ref(request.clone()).unwrap(), + expected + ); + assert_eq!(observed.lock().unwrap().as_slice(), &[request]); + } + #[derive(Default)] struct AcceptingExecutionBackend { contexts: diff --git a/crates/workspace-server/src/records.rs b/crates/workspace-server/src/records.rs index 6d0f1690..2cde25d8 100644 --- a/crates/workspace-server/src/records.rs +++ b/crates/workspace-server/src/records.rs @@ -314,12 +314,21 @@ pub struct TicketMergeRequestSummary { pub review_excerpt: Option, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +pub struct MergeRequestRefDiagnostic { + pub code: String, + pub message: String, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct MergeRequestListItem { pub summary: TicketMergeRequestSummary, pub ticket_ids: Vec, pub thread_event_count: usize, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub ref_diagnostics: Vec, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -470,6 +479,7 @@ pub fn ticket_api_typescript() -> String { TicketAssignmentPrincipalSummary::decl(&config), TicketActionEligibility::decl(&config), TicketMergeRequestSummary::decl(&config), + MergeRequestRefDiagnostic::decl(&config), MergeRequestListItem::decl(&config), MergeRequestListResponse::decl(&config), TicketEvidenceSummary::decl(&config), diff --git a/crates/workspace-server/src/repositories.rs b/crates/workspace-server/src/repositories.rs index 0044f568..3fc43c66 100644 --- a/crates/workspace-server/src/repositories.rs +++ b/crates/workspace-server/src/repositories.rs @@ -392,7 +392,7 @@ impl RepositoryRegistryReader { } } -fn normalize_target_branch_selector( +pub(crate) fn normalize_target_branch_selector( id: &str, selector: &str, ) -> Result { diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 30f79ef5..83e8904d 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -58,25 +58,28 @@ use worker::feature::builtin::{WorkerObservationSubject, WorkerObservationSubjec use worker_runtime::resource::{BackendResourceError, BackendResourceFetchRequest}; use worker_runtime::worker_backend::{ProfileRuntimeWorkerFactory, WorkerRuntimeExecutionBackend}; use workspace_api::{ - CreateRemoteRuntimeRequest, CreateRepositorySshCredentialRequest, - CreateWorkspaceRepositoryRequest, CreateWorkspaceRepositoryResponse, - DeleteRepositorySshCredentialRequest, DeleteRepositorySshHostTrustRequest, - ObjectiveCreateRequest, ObjectiveEditRequest, ObjectiveLinkTicketRequest, - ObjectiveStateRequest, ProfileSettingsResponse, PutRepositorySshHostTrustRequest, - RepositoryAccessProjection, RepositoryDetailResponse, RepositoryListResponse, - RepositoryLogResponse, RepositorySshCredential, RepositorySshHostTrust, + BrowserCreateWorkerResponse, BrowserWorkspaceOrchestratorResponse, CreateRemoteRuntimeRequest, + CreateRepositorySshCredentialRequest, CreateWorkspaceRepositoryRequest, + CreateWorkspaceRepositoryResponse, CreateWorkspaceWorkerRequest, + CreateWorkspaceWorkerTicketAssignmentRequest, DeleteRepositorySshCredentialRequest, + DeleteRepositorySshHostTrustRequest, ObjectiveCreateRequest, ObjectiveEditRequest, + ObjectiveLinkTicketRequest, ObjectiveStateRequest, ProfileSettingsResponse, + PutRepositorySshHostTrustRequest, RepositoryAccessProjection, RepositoryDetailResponse, + RepositoryListResponse, RepositoryLogResponse, RepositorySshCredential, RepositorySshHostTrust, RotateRepositorySshCredentialRequest, RuntimeConnectionTestResponse, RuntimeManagementSummary, TICKET_ORCHESTRATION_PLANS_QUERY_PATH, TICKET_RELATIONS_QUERY_PATH, - UpdateWorkspaceMetadataRequest, + UpdateWorkspaceMetadataRequest, WorkerLaunchOptionsResponse, WorkerLaunchProfileCandidate, + WorkerLaunchRuntimeOption, WorkerLaunchWorkerSummary, WorkingDirectoryCreateRequest as BrowserWorkingDirectoryCreateRequest, WorkingDirectoryCreateResponse as BrowserWorkingDirectoryCreateResponse, WorkingDirectoryDetailResponse as BrowserWorkingDirectoryDetailResponse, WorkingDirectoryListResponse as BrowserWorkingDirectoryListResponse, WorkingDirectoryRemovalDisposition, WorkingDirectoryRemovalRequest, - WorkingDirectoryRemovalResponse, WorkspaceCatalogListResponse, WorkspaceCreateResponse, - WorkspaceExtensionPointState, WorkspaceExtensionPoints, WorkspaceMetadataMutationResponse, - WorkspaceMetadataSettingsResponse, WorkspacePermissionSummary, WorkspaceRepositoryRecord, - WorkspaceResponse, WorkspaceRuntimeResource, WorkspaceSummary, WorkspaceWorkerDiscoveryItem, + WorkingDirectoryRemovalResponse, WorkingDirectoryRepositoryOption, + WorkspaceCatalogListResponse, WorkspaceCreateResponse, WorkspaceExtensionPointState, + WorkspaceExtensionPoints, WorkspaceMetadataMutationResponse, WorkspaceMetadataSettingsResponse, + WorkspacePermissionSummary, WorkspaceRepositoryRecord, WorkspaceResponse, + WorkspaceRuntimeResource, WorkspaceSummary, WorkspaceWorkerDiscoveryItem, WorkspaceWorkerDiscoveryPage, WorkspaceWorkerSubject, }; @@ -118,9 +121,9 @@ use crate::observation::{ RuntimeObservationSource, RuntimeObservationSourceConfig, }; use crate::records::{ - MergeRequestListItem, MergeRequestListResponse, ObjectiveDetail, ObjectiveQueryRequest, - ObjectiveQueryResponse, ObjectiveShowRequest, ProjectRecordList, TicketDetail, - TicketQueryRequest, TicketQueryResponse, TicketShowRequest, + MergeRequestListItem, MergeRequestListResponse, MergeRequestRefDiagnostic, ObjectiveDetail, + ObjectiveQueryRequest, ObjectiveQueryResponse, ObjectiveShowRequest, ProjectRecordList, + TicketDetail, TicketQueryRequest, TicketQueryResponse, TicketShowRequest, }; use crate::repositories::{ ConfiguredRepository, RepositoryListProjection, RepositoryLogRead, RepositoryLookupError, @@ -150,10 +153,10 @@ use crate::workdir_removal::{ use crate::workspace_catalog::{WorkspaceCatalogService, WorkspaceCreateRequest}; use crate::{Error, Result}; use worker_runtime::catalog::{ - ConfigBundleRef, ProfileSelector, RepositoryMaterializationContext, - RepositorySelector as RuntimeRepositorySelector, RepositorySshMaterializationAccess, - SensitiveString, WorkingDirectoryClaim, WorkingDirectoryRepository, WorkingDirectoryRequest, - WorkspaceApiRef, + ConfigBundleRef, ProfileSelector, RepositoryMaterializationContext, RepositoryRefObservation, + RepositoryRefObservationRequest, RepositorySelector as RuntimeRepositorySelector, + RepositorySshMaterializationAccess, SensitiveString, WorkingDirectoryClaim, + WorkingDirectoryRepository, WorkingDirectoryRequest, WorkspaceApiRef, }; use worker_runtime::config_bundle::ConfigBundle; use worker_runtime::http_server::{ @@ -3097,98 +3100,6 @@ pub struct WorkerRetentionResponse { pub retention_state: String, } -#[derive(Debug, Serialize, Deserialize)] -pub struct WorkerLaunchOptionsResponse { - pub workspace_id: String, - pub runtimes: Vec, - pub default_profile: Option, - pub profiles: Vec, - pub repositories: Vec, - pub working_directories: Vec, - pub diagnostics: Vec, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct WorkerLaunchRuntimeOption { - pub runtime_id: String, - pub display_name: String, - pub built_in: bool, - pub worker_creation_available: bool, - pub working_directory_required: bool, - pub status: String, - pub diagnostics: Vec, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct WorkerLaunchProfileCandidate { - pub id: String, - pub label: String, - pub description: String, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct WorkingDirectoryRepositoryOption { - pub repository_key: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub default_selector: Option, -} - -#[derive(Debug, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct BrowserWorkerWorkingDirectorySelection { - pub working_directory_id: String, - #[serde(default)] - pub relative_cwd: Option, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct BrowserWorkspaceOrchestratorResponse { - pub workspace_id: String, - pub online: bool, - pub disposition: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub worker: Option, - pub diagnostics: Vec, -} - -#[derive(Debug, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct CreateWorkspaceWorkerTicketAssignmentRequest { - pub ticket_id: String, - pub operation_id: String, -} - -#[derive(Debug, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct CreateWorkspaceWorkerRequest { - pub runtime_id: String, - pub display_name: String, - #[serde(default)] - pub profile: Option, - #[serde(default)] - pub ticket_assignment: Option, - #[serde(default)] - pub initial_submit: Vec, - #[serde(default)] - pub working_directory: Option, - /// Backend idempotency key used only for authenticated Worker-owned spawn/control. - #[serde(default)] - pub control_operation_id: Option, - /// Trusted resolution populated only by the authenticated worker-control handler. - #[serde(skip, default)] - pub resolved_control_operation: Option, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct BrowserCreateWorkerResponse { - pub workspace_id: String, - #[serde(flatten)] - pub worker_ref: RuntimeWorkerRef, - pub console_href: String, - pub worker: WorkerSummary, - pub diagnostics: Vec, -} - #[derive(Debug, Deserialize)] struct LogQuery { limit: Option, @@ -5506,6 +5417,205 @@ fn merge_request_store( .map_err(Into::into) } +fn repository_ref_observation_error(error: crate::hosts::RuntimeRegistryError) -> ApiError { + match error { + crate::hosts::RuntimeRegistryError::RuntimeOperationFailed { + runtime_id, + code, + message, + } => Error::RuntimeOperationFailed { + runtime_id, + code, + message, + } + .into(), + other => other.into_error().into(), + } +} + +fn observe_published_merge_ref( + api: &WorkspaceApi, + workspace_id: &str, + runtime_id: &str, + repository_id: &str, + selector: &str, +) -> ApiResult { + let canonical_selector = + crate::repositories::normalize_target_branch_selector(repository_id, selector) + .map_err(repository_merge_evidence_error)?; + let operation_id = format!("repository-ref-observation-{}", Uuid::now_v7()); + let repository = api + .config + .repositories + .iter() + .find(|repository| repository.id == repository_id) + .ok_or_else(|| Error::UnknownRepository(repository_id.to_string()))?; + let mut materialization_request = + working_directory_request_from_repository(repository, Some(&canonical_selector)); + materialization_request.backend_workdir_id = Some(operation_id.clone()); + let projection = active_repository_access_projection(api, workspace_id)?; + authorize_repository_materialization( + api, + runtime_id, + &operation_id, + &projection, + &mut materialization_request, + )?; + api.runtime + .observe_repository_ref( + runtime_id, + RepositoryRefObservationRequest { + repository: materialization_request.repository, + selector: canonical_selector, + materialization: materialization_request.materialization, + }, + ) + .map_err(repository_ref_observation_error) +} + +fn source_ref_error_code(code: &str) -> String { + match code { + "repository_ref_not_found" => "source_ref_not_found".to_string(), + "repository_ref_provider_timeout" => "source_ref_provider_timeout".to_string(), + "repository_ref_provider_auth_failed" => "source_ref_provider_auth_failed".to_string(), + "repository_ref_response_invalid" => "source_ref_response_invalid".to_string(), + "repository_ref_selector_invalid" => "source_ref_selector_invalid".to_string(), + "repository_access_credential_expired" => "source_ref_credential_expired".to_string(), + "repository_access_credential_unavailable" => { + "source_ref_credential_unavailable".to_string() + } + "repository_access_credential_unauthorized" => { + "source_ref_credential_unauthorized".to_string() + } + "repository_access_credential_invalid" => "source_ref_credential_invalid".to_string(), + "repository_access_provider_unavailable" | "repository_ref_provider_unavailable" => { + "source_ref_provider_unavailable".to_string() + } + other => other.to_string(), + } +} + +fn source_ref_readiness_blocker(code: &str) -> &'static str { + if code == "source_ref_not_found" { + "source_ref_not_found" + } else { + "source_ref_unavailable" + } +} + +fn remap_source_ref_error(error: ApiError) -> ApiError { + let ApiError { error, diagnostics } = error; + match error { + Error::RuntimeOperationFailed { + runtime_id, + code, + message, + } => Error::RuntimeOperationFailed { + runtime_id, + code: source_ref_error_code(&code), + message, + } + .into(), + error => ApiError { error, diagnostics }, + } +} + +fn observe_published_source_ref( + api: &WorkspaceApi, + workspace_id: &str, + runtime_id: &str, + repository_id: &str, + selector: &str, +) -> ApiResult { + observe_published_merge_ref(api, workspace_id, runtime_id, repository_id, selector) + .map_err(remap_source_ref_error) +} + +fn require_assigned_workdir_source( + api: &WorkspaceApi, + assignment: &crate::store::TicketCoderAssignmentRecord, + repository_id: &str, + selector: &str, + revision_ref: &str, +) -> ApiResult<()> { + let worker = api + .runtime + .worker(&assignment.worker) + .map_err(|error| error.into_error())?; + let attached_workdir = worker.working_directory.ok_or_else(|| { + Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_missing: current Coder has no attached Workdir".into(), + )) + })?; + let workdir = api + .runtime + .working_directory( + &assignment.worker.runtime_id, + &attached_workdir.working_directory_id, + ) + .map_err(|error| error.into_error())? + .working_directory + .ok_or_else(|| { + Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_unavailable: current Coder Workdir could not be observed" + .into(), + )) + })? + .summary; + validate_assigned_workdir_source(&workdir, repository_id, selector, revision_ref).map_err( + |diagnostic| Error::RuntimeOperationFailed { + runtime_id: assignment.worker.runtime_id.clone(), + code: diagnostic.code, + message: diagnostic.message, + }, + )?; + Ok(()) +} + +fn validate_assigned_workdir_source( + workdir: &worker_runtime::catalog::WorkingDirectorySummary, + repository_id: &str, + selector: &str, + revision_ref: &str, +) -> std::result::Result<(), worker_runtime::working_directory::WorkingDirectoryDiagnostic> { + if workdir.repository_id != repository_id { + return Err( + worker_runtime::working_directory::WorkingDirectoryDiagnostic { + code: "source_workdir_repository_mismatch".to_string(), + message: "Current Coder Workdir belongs to a different Repository".to_string(), + }, + ); + } + if workdir.cleanliness.as_deref() != Some("clean") { + return Err( + worker_runtime::working_directory::WorkingDirectoryDiagnostic { + code: "source_workdir_dirty".to_string(), + message: "Current Coder Workdir must be clean".to_string(), + }, + ); + } + let workdir_selector_matches = workdir + .current_selector + .as_deref() + .and_then(|workdir_selector| { + crate::repositories::normalize_target_branch_selector(repository_id, workdir_selector) + .ok() + }) + .is_some_and(|workdir_selector| { + crate::repositories::normalize_target_branch_selector(repository_id, selector) + .is_ok_and(|selector| selector == workdir_selector) + }); + if !workdir_selector_matches || workdir.current_ref.as_deref() != Some(revision_ref) { + return Err( + worker_runtime::working_directory::WorkingDirectoryDiagnostic { + code: "source_ref_revision_mismatch".to_string(), + message: "Current Coder Workdir selector and HEAD do not match the provider-published source ref".to_string(), + }, + ); + } + Ok(()) +} + fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError { Error::InvalidInput(format!( "repository merge evidence validation failed: {error:?}" @@ -5625,6 +5735,8 @@ struct MergeRequestRefResponse { #[serde(rename = "ref")] revision_ref: Option, observed_at: String, + #[serde(skip_serializing_if = "Option::is_none")] + diagnostic: Option, } #[derive(Debug, serde::Serialize)] @@ -5704,17 +5816,44 @@ async fn scoped_list_merge_requests( limit: query.limit.unwrap_or(50), }, )?; - let reader = api.repository_reader(); let items = page .items .into_iter() .map(|merge_request| -> ApiResult { - let current_subject_ref = merge_request.selector_from.as_deref().and_then(|selector| { - reader - .observe_merge_target(&merge_request.repository_id, Some(selector)) - .ok() - .map(|observation| observation.commit) - }); + let (current_subject_ref, ref_diagnostics) = + if merge_request.state == merge_request::MergeRequestState::Open { + match ( + merge_request.selector_from.as_deref(), + merge_request.ticket_ids.first(), + ) { + (Some(selector), Some(ticket_id)) => match api + .store + .get_current_ticket_coder_assignment(&workspace_id, ticket_id)? + { + Some(assignment) => match observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &merge_request.repository_id, + selector, + ) { + Ok(observation) => (Some(observation.revision_ref), Vec::new()), + Err(error) => (None, vec![merge_ref_diagnostic(error)]), + }, + None => ( + None, + vec![MergeRequestRefDiagnostic { + code: "source_ref_runtime_unavailable".to_string(), + message: "No current Coder Runtime is available to observe the source ref" + .to_string(), + }], + ), + }, + _ => (None, Vec::new()), + } + } else { + (None, Vec::new()) + }; let repository_key = api .store .get_repository(&workspace_id, &merge_request.repository_id)? @@ -5730,6 +5869,7 @@ async fn scoped_list_merge_requests( summary: merge_request_summary(merge_request, repository_key, current_subject_ref), ticket_ids, thread_event_count, + ref_diagnostics, }) }) .collect::>>()?; @@ -5739,6 +5879,55 @@ async fn scoped_list_merge_requests( })) } +fn merge_ref_diagnostic(error: ApiError) -> MergeRequestRefDiagnostic { + let ApiError { error, diagnostics } = error; + diagnostics + .into_iter() + .next() + .map(|diagnostic| MergeRequestRefDiagnostic { + code: diagnostic.code, + message: diagnostic.message, + }) + .unwrap_or_else(|| MergeRequestRefDiagnostic { + code: "merge_ref_observation_unavailable".to_string(), + message: sanitize_backend_error(&error.to_string()), + }) +} + +fn unknown_merge_ref(code: &str, message: &str) -> MergeRequestRefResponse { + MergeRequestRefResponse { + status: "unknown".to_string(), + revision_ref: None, + observed_at: Utc::now().to_rfc3339(), + diagnostic: Some(MergeRequestRefDiagnostic { + code: code.to_string(), + message: message.to_string(), + }), + } +} + +fn unknown_merge_ref_response(error: ApiError) -> MergeRequestRefResponse { + MergeRequestRefResponse { + status: "unknown".to_string(), + revision_ref: None, + observed_at: Utc::now().to_rfc3339(), + diagnostic: Some(merge_ref_diagnostic(error)), + } +} + +fn merge_ref_response(observation: RepositoryRefObservation) -> MergeRequestRefResponse { + let observed_at = + chrono::DateTime::::from_timestamp(observation.observed_at_epoch_seconds as i64, 0) + .unwrap_or_else(Utc::now) + .to_rfc3339(); + MergeRequestRefResponse { + status: "known".into(), + revision_ref: Some(observation.revision_ref), + observed_at, + diagnostic: None, + } +} + async fn scoped_show_merge_request( State(api): State, AxumPath((workspace_id, merge_request_id)): AxumPath<(String, String)>, @@ -5754,38 +5943,53 @@ async fn scoped_show_merge_request( query.after, query.limit.unwrap_or(100), )?; - let reader = api.repository_reader(); - let observed_at = Utc::now().to_rfc3339(); - let source = match mr.selector_from.as_deref() { - Some(selector) => match reader.observe_merge_target(&mr.repository_id, Some(selector)) { - Ok(value) => MergeRequestRefResponse { - status: "known".into(), - revision_ref: Some(value.commit), - observed_at: observed_at.clone(), - }, - Err(_) => MergeRequestRefResponse { - status: "unknown".into(), - revision_ref: None, - observed_at: observed_at.clone(), - }, - }, - None => MergeRequestRefResponse { + let ticket_id = mr.ticket_ids.first().ok_or_else(|| { + Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "Merge Request has no linked Ticket".into(), + )) + })?; + let assignment = api + .store + .get_current_ticket_coder_assignment(&workspace_id, ticket_id)?; + let source = match (mr.selector_from.as_deref(), assignment.as_ref()) { + (Some(selector), Some(assignment)) => { + match observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + ) { + Ok(observation) => merge_ref_response(observation), + Err(error) => unknown_merge_ref_response(error), + } + } + (Some(_), None) => unknown_merge_ref( + "source_ref_runtime_unavailable", + "No current Coder Runtime is available to observe the source ref", + ), + (None, _) => MergeRequestRefResponse { status: "requires_repair".into(), revision_ref: None, - observed_at: observed_at.clone(), + observed_at: Utc::now().to_rfc3339(), + diagnostic: None, }, }; - let target = match reader.observe_merge_target(&mr.repository_id, Some(&mr.selector_to)) { - Ok(value) => MergeRequestRefResponse { - status: "known".into(), - revision_ref: Some(value.commit), - observed_at, - }, - Err(_) => MergeRequestRefResponse { - status: "unknown".into(), - revision_ref: None, - observed_at, + let target = match assignment.as_ref() { + Some(assignment) => match observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + &mr.selector_to, + ) { + Ok(observation) => merge_ref_response(observation), + Err(error) => unknown_merge_ref_response(error), }, + None => unknown_merge_ref( + "target_ref_runtime_unavailable", + "No current Coder Runtime is available to observe the target ref", + ), }; let linked_tickets = mr .ticket_ids @@ -5818,13 +6022,28 @@ async fn scoped_merge_request_readiness( let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?; let store = merge_request_store(&api, &workspace_id)?; let mr = store.get(&workspace_id, &ticket_id)?; - let current_subject_ref = mr.selector_from.as_deref().and_then(|selector| { - api.repository_reader() - .observe_merge_target(&mr.repository_id, Some(selector)) - .ok() - .map(|v| v.commit) - }); - Ok(Json(store.readiness(merge_request::ReadinessCheck { + let assignment = api + .store + .get_current_ticket_coder_assignment(&workspace_id, &ticket_id)?; + let (current_subject_ref, source_blocker) = match (mr.selector_from.as_deref(), assignment) { + (Some(selector), Some(assignment)) => match observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + ) { + Ok(observation) => (Some(observation.revision_ref), None), + Err(error) => { + let diagnostic = merge_ref_diagnostic(error); + let blocker = source_ref_readiness_blocker(&diagnostic.code); + (None, Some(blocker.to_string())) + } + }, + (Some(_), None) => (None, Some("source_ref_unavailable".to_string())), + (None, _) => (None, None), + }; + let mut report = store.readiness(merge_request::ReadinessCheck { ticket_id, current_subject_ref, auth: merge_request::MergeRequestAuth { @@ -5834,7 +6053,17 @@ async fn scoped_merge_request_readiness( worker_id: String::new(), assignment_id: String::new(), }, - })?)) + })?; + if let Some(blocker) = source_blocker { + report + .blockers + .retain(|current| current != "selector_unresolved"); + if !report.blockers.contains(&blocker) { + report.blockers.push(blocker); + } + report.ready = false; + } + Ok(Json(report)) } async fn scoped_open_merge_request( @@ -5874,13 +6103,27 @@ async fn scoped_open_merge_request( ) .into()); } - let reader = api.repository_reader(); - reader - .observe_merge_target(&repository_id, Some(&input.selector_from)) - .map_err(repository_merge_evidence_error)?; - reader - .observe_merge_target(&repository_id, Some(&input.selector_to)) - .map_err(repository_merge_evidence_error)?; + let source_observation = observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &repository_id, + &input.selector_from, + )?; + require_assigned_workdir_source( + &api, + &assignment, + &repository_id, + &input.selector_from, + &source_observation.revision_ref, + )?; + observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &repository_id, + &input.selector_to, + )?; let merge_request = merge_request_store(&api, &workspace_id)?.open_merge_request( merge_request::OpenMergeRequest { merge_request_id: Uuid::now_v7().to_string(), @@ -5939,11 +6182,20 @@ async fn scoped_repair_merge_request_selector( } let store = merge_request_store(&api, &workspace_id)?; let mr = store.get(&workspace_id, &ticket_id)?; - let resolved_subject_ref = api - .repository_reader() - .observe_merge_target(&mr.repository_id, Some(&input.selector_from)) - .map_err(repository_merge_evidence_error)? - .commit; + let assignment = api + .store + .get_current_ticket_coder_assignment(&workspace_id, &ticket_id)? + .ok_or_else(|| { + Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) + })?; + let resolved_subject_ref = observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + &input.selector_from, + )? + .revision_ref; let repaired = store.repair_selector_from(merge_request::RepairSelectorFrom { workspace_id: workspace_id.clone(), ticket_id, @@ -6011,11 +6263,21 @@ async fn scoped_register_merge_request_review_capability( .selector_from .as_deref() .ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?; - let subject_ref = api - .repository_reader() - .observe_merge_target(&mr.repository_id, Some(selector)) - .map_err(repository_merge_evidence_error)? - .commit; + let source_observation = observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + )?; + require_assigned_workdir_source( + &api, + &assignment, + &mr.repository_id, + selector, + &source_observation.revision_ref, + )?; + let subject_ref = source_observation.revision_ref; store.request_review(merge_request::RequestMergeRequestReview { ticket_id, subject_ref, @@ -6041,16 +6303,35 @@ async fn scoped_submit_merge_request_review( let workspace_id = parse_workspace_id(&workspace_id)?; let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?; let store = merge_request_store(&api, &workspace_id)?; + let review_authorization = + store.authorize_review_submission(&ticket_id, &input.capability_token)?; + if review_authorization.workspace_id != workspace_id { + return Err( + Error::MergeRequest(merge_request::MergeRequestError::Unauthorized( + "review grant invalid".into(), + )) + .into(), + ); + } let mr = store.get(&workspace_id, &ticket_id)?; let selector = mr .selector_from .as_deref() .ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?; - let current_subject_ref = api - .repository_reader() - .observe_merge_target(&mr.repository_id, Some(selector)) - .map_err(repository_merge_evidence_error)? - .commit; + let assignment = api + .store + .get_current_ticket_coder_assignment(&workspace_id, &ticket_id)? + .ok_or_else(|| { + Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) + })?; + let current_subject_ref = observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + )? + .revision_ref; Ok(Json(store.submit_review( merge_request::SubmitMergeRequestReview { ticket_id, @@ -6123,7 +6404,6 @@ async fn scoped_complete_merge_request( require_online_workspace_orchestrator_source(&api, &source)?; let store = merge_request_store(&api, &workspace_id)?; let mr = store.get(&workspace_id, &ticket_id)?; - let repositories = api.repository_reader(); if let Some(existing) = recorded_merge_completion(&mr.thread, &input.operation_id) { let replay = merge_request::CompleteMergeRequest { ticket_id, @@ -6155,15 +6435,23 @@ async fn scoped_complete_merge_request( .selector_from .as_deref() .ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?; - let current_source_ref = repositories - .observe_merge_target(&mr.repository_id, Some(selector)) - .map_err(repository_merge_evidence_error)? - .commit; - let observed = repositories - .observe_merge_target(&mr.repository_id, Some(&mr.selector_to)) - .map_err(repository_merge_evidence_error)?; + let current_source_ref = observe_published_source_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + )? + .revision_ref; + let observed_target = observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + &mr.selector_to, + )?; require_completed_target_observation( - &observed.commit, + &observed_target.revision_ref, &input.target_ref_before, &input.target_ref_after, )?; @@ -7758,9 +8046,9 @@ fn start_memory_staging_consolidation( resolved_config_bundle, resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }, )?; if result.state != WorkerOperationState::Accepted { @@ -8711,18 +8999,20 @@ async fn spawn_known_worker( .map(|byte| format!("{byte:02x}")) .collect::() ); - request.resolved_control_operation = Some(WorkerControlOperation { + let resolved_control_operation = Some(WorkerControlOperation { operation_id: scoped_worker_control_operation_id(&controller, &operation_id), input_fingerprint, }); - let response = create_workspace_worker(State(api.clone()), headers, Json(request)).await?; + let response = + create_workspace_worker_inner(api.clone(), headers, request, resolved_control_operation) + .await?; if let Err(error) = api .store .create_worker_control_grant(&WorkerControlGrantRecord { workspace_id: path.workspace_id.clone(), grant_id: new_id("wcg"), controller, - subject: response.0.worker_ref.clone(), + subject: RuntimeWorkerRef::new(&response.0.runtime_id, &response.0.worker_id), relation: relation.to_string(), origin: "worker_spawn".to_string(), permissions: vec![ @@ -9055,9 +9345,9 @@ async fn scoped_start_workspace_orchestrator( resolved_config_bundle: None, resolved_worker_observation_enabled: true, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }, )?; if result.state != WorkerOperationState::Accepted || result.worker.is_none() { @@ -9085,6 +9375,42 @@ async fn scoped_start_workspace_orchestrator( Ok(Json(workspace_orchestrator_response(&api, "created"))) } +fn worker_launch_worker_summary(worker: WorkerSummary) -> WorkerLaunchWorkerSummary { + WorkerLaunchWorkerSummary { + runtime_id: worker.worker.runtime_id, + worker_id: worker.worker.worker_id, + host_id: worker.host_id, + display_name: worker.display_name, + label: worker.label, + profile: worker.profile, + singleton_key: worker.singleton_key, + tags: worker.tags, + workspace: workspace_api::WorkerWorkspaceSummary { + visibility: worker.workspace.visibility, + identity: worker.workspace.identity, + workspace_id: worker.workspace.workspace_id, + }, + state: worker.state, + last_seen_at: worker.last_seen_at, + pinned: worker.pinned, + retention_state: worker.retention_state, + implementation: workspace_api::WorkerImplementationSummary { + kind: worker.implementation.kind, + display_hint: worker.implementation.display_hint, + }, + capabilities: workspace_api::WorkerCapabilitySummary { + can_stop: worker.capabilities.can_stop, + can_spawn_followup: worker.capabilities.can_spawn_followup, + }, + working_directory: worker.working_directory, + diagnostics: worker + .diagnostics + .into_iter() + .map(workspace_api::Diagnostic::from) + .collect(), + } +} + fn workspace_orchestrator_response( api: &WorkspaceApi, disposition: &str, @@ -9095,13 +9421,20 @@ fn workspace_orchestrator_response( .is_some_and(workspace_orchestrator_is_online); let diagnostics = worker .as_ref() - .map(|worker| worker.diagnostics.clone()) + .map(|worker| { + worker + .diagnostics + .iter() + .cloned() + .map(workspace_api::Diagnostic::from) + .collect() + }) .unwrap_or_default(); BrowserWorkspaceOrchestratorResponse { workspace_id: api.config.workspace_id.clone(), online, disposition: disposition.to_string(), - worker, + worker: worker.map(worker_launch_worker_summary), diagnostics, } } @@ -12458,6 +12791,15 @@ async fn create_workspace_worker( State(api): State, headers: HeaderMap, Json(request): Json, +) -> ApiResult> { + create_workspace_worker_inner(api, headers, request, None).await +} + +async fn create_workspace_worker_inner( + api: WorkspaceApi, + headers: HeaderMap, + request: CreateWorkspaceWorkerRequest, + resolved_control_operation: Option, ) -> ApiResult> { let CreateWorkspaceWorkerRequest { runtime_id, @@ -12467,7 +12809,6 @@ async fn create_workspace_worker( initial_submit, working_directory, control_operation_id: _, - resolved_control_operation, } = request; let config_state = api .config_store @@ -12765,10 +13106,14 @@ fn browser_worker_response_from_summary( ); Ok(BrowserCreateWorkerResponse { workspace_id, - worker_ref: RuntimeWorkerRef::new(&runtime_id, &worker_id), + runtime_id, + worker_id, console_href, - worker, - diagnostics, + worker: worker_launch_worker_summary(worker), + diagnostics: diagnostics + .into_iter() + .map(workspace_api::Diagnostic::from) + .collect(), }) } @@ -14801,7 +15146,11 @@ fn worker_launch_options_response(api: &WorkspaceApi) -> ApiResult(source: &'a str, name: &str) -> &'a str { + let start = source + .find(&format!("async fn {name}")) + .unwrap_or_else(|| panic!("missing handler {name}")); + let tail = &source[start..]; + let end = tail[1..] + .find("\nasync fn ") + .map(|offset| offset + 1) + .unwrap_or(tail.len()); + &tail[..end] + } + + #[test] + fn merge_request_http_paths_observe_refs_through_runtime_provider_authority() { + let source = include_str!("server.rs"); + for handler in [ + "scoped_list_merge_requests", + "scoped_show_merge_request", + "scoped_open_merge_request", + "scoped_repair_merge_request_selector", + "scoped_register_merge_request_review_capability", + "scoped_submit_merge_request_review", + "scoped_complete_merge_request", + ] { + let handler = handler_source(source, handler); + assert!( + handler.contains("observe_published_source_ref(") + || handler.contains("observe_published_merge_ref(") + ); + assert!(!handler.contains("repository_reader()")); + } + let submit = handler_source(source, "scoped_submit_merge_request_review"); + assert!( + submit + .find("authorize_review_submission(") + .is_some_and(|authorization| { + submit + .find("observe_published_source_ref(") + .is_some_and(|observation| authorization < observation) + }), + "review capability must be validated before provider access", + ); + } + + #[test] + fn readiness_distinguishes_missing_source_from_unavailable_provider() { + assert_eq!( + source_ref_readiness_blocker("source_ref_not_found"), + "source_ref_not_found" + ); + assert_eq!( + source_ref_error_code("repository_ref_not_found"), + "source_ref_not_found" + ); + assert_eq!( + source_ref_error_code("repository_ref_provider_timeout"), + "source_ref_provider_timeout" + ); + for code in [ + "source_ref_provider_unavailable", + "source_ref_provider_timeout", + "source_ref_provider_auth_failed", + "source_ref_credential_expired", + ] { + assert_eq!(source_ref_readiness_blocker(code), "source_ref_unavailable"); + } + } + + #[test] + fn assigned_workdir_must_be_clean_and_match_published_source() { + let mut workdir = worker_runtime::catalog::WorkingDirectorySummary { + working_directory_id: "workdir-1".to_string(), + repository_id: "repository-1".to_string(), + creation_selector: Some("work/T-549".to_string()), + creation_ref: Some("abc123".to_string()), + creation_tree: None, + current_selector: Some("work/T-549".to_string()), + current_ref: Some("abc123".to_string()), + current_tree: None, + observed_at_epoch_seconds: Some(1_767_225_600), + materializer_kind: workspace_api::WorkingDirectoryMaterializerKind::RuntimeGitCache, + cleanup_target: None, + status: worker_runtime::catalog::WorkingDirectoryStatusKind::Active, + cleanliness: Some("clean".to_string()), + primary_worker_id: Some("worker-1".to_string()), + occupied_by: None, + }; + assert!( + validate_assigned_workdir_source( + &workdir, + "repository-1", + "refs/heads/work/T-549", + "abc123", + ) + .is_ok() + ); + + workdir.cleanliness = Some("dirty".to_string()); + assert!(matches!( + validate_assigned_workdir_source( + &workdir, + "repository-1", + "work/T-549", + "abc123" + ), + Err(diagnostic) if diagnostic.code == "source_workdir_dirty" + )); + workdir.cleanliness = Some("clean".to_string()); + workdir.current_ref = Some("different".to_string()); + assert!(matches!( + validate_assigned_workdir_source( + &workdir, + "repository-1", + "work/T-549", + "abc123" + ), + Err(diagnostic) if diagnostic.code == "source_ref_revision_mismatch" + )); + } + fn completed_upload_file(sha256: &str) -> protocol::UploadedFileRef { protocol::UploadedFileRef { artifact_id: "019ca7c8-57b6-7f05-8edf-524147aba7b3".into(), @@ -17201,9 +17670,9 @@ mod tests { resolved_config_bundle: None, resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }; assert!( api.validate_worker_spawn_repository_scope(&workdir_flow_launch) @@ -17447,9 +17916,9 @@ mod tests { resolved_config_bundle: None, resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }; assert!( @@ -17496,7 +17965,6 @@ mod tests { }], working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await @@ -17512,7 +17980,8 @@ mod tests { .get_current_ticket_coder_assignment(&api.config.workspace_id, &ticket.id) .unwrap() .unwrap(); - assert_eq!(current.worker, response.worker_ref); + let response_ref = RuntimeWorkerRef::new(&response.runtime_id, &response.worker_id); + assert_eq!(current.worker, response_ref); let operation = api .store .get_ticket_assignment_operation( @@ -17522,7 +17991,7 @@ mod tests { .unwrap() .unwrap(); assert_eq!(operation.assignment_id, Some(current.assignment_id)); - assert_eq!(operation.worker, Some(response.worker_ref)); + assert_eq!(operation.worker, Some(response_ref)); } #[tokio::test] @@ -17552,7 +18021,6 @@ mod tests { }], working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await; @@ -17586,7 +18054,6 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await @@ -17594,11 +18061,11 @@ mod tests { let mut headers = HeaderMap::new(); headers.insert( "x-yoi-runtime-id", - axum::http::HeaderValue::from_str(&created.worker_ref.runtime_id).unwrap(), + axum::http::HeaderValue::from_str(&created.runtime_id).unwrap(), ); headers.insert( "x-yoi-worker-id", - axum::http::HeaderValue::from_str(&created.worker_ref.worker_id).unwrap(), + axum::http::HeaderValue::from_str(&created.worker_id).unwrap(), ); let error = authenticate_worker_mutation_source(&api, "other-workspace", &headers).unwrap_err(); @@ -17622,14 +18089,14 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await .unwrap(); + let generic_ref = RuntimeWorkerRef::new(&generic.runtime_id, &generic.worker_id); assert!(matches!( - require_online_workspace_orchestrator_source(&api, &generic.worker_ref), + require_online_workspace_orchestrator_source(&api, &generic_ref), Err(Error::TicketAssignmentConflict(_)) )); @@ -17641,10 +18108,11 @@ mod tests { ) .await .unwrap(); - let orchestrator = started.worker.unwrap().worker; + let orchestrator = started.worker.unwrap(); + let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id); require_online_workspace_orchestrator_source(&api, &orchestrator).unwrap(); assert!(matches!( - require_online_workspace_orchestrator_source(&api, &generic.worker_ref), + require_online_workspace_orchestrator_source(&api, &generic_ref), Err(Error::TicketAssignmentConflict(_)) )); @@ -17713,8 +18181,8 @@ mod tests { assert!(started.online); let worker = started .worker - .expect("production Workspace Orchestrator Worker") - .worker; + .expect("production Workspace Orchestrator Worker"); + let worker = RuntimeWorkerRef::new(&worker.runtime_id, &worker.worker_id); let stopped = api .runtime @@ -17814,12 +18282,12 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await .unwrap(); - let controller = controller_worker.worker_ref; + let controller = + RuntimeWorkerRef::new(&controller_worker.runtime_id, &controller_worker.worker_id); assert_ne!( scoped_worker_control_operation_id(&controller, "same-operation"), scoped_worker_control_operation_id( @@ -17845,7 +18313,6 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: Some("control-spawn-retry".to_string()), - resolved_control_operation: None, }; let Json(first) = spawn_known_worker( @@ -17869,7 +18336,8 @@ mod tests { .await .unwrap(); - assert_eq!(retried.worker_ref, first.worker_ref); + assert_eq!(retried.runtime_id, first.runtime_id); + assert_eq!(retried.worker_id, first.worker_id); let mut conflicting_request = request(); conflicting_request.display_name = "Different controlled child".to_string(); let conflict = spawn_known_worker( @@ -17899,7 +18367,10 @@ mod tests { .list_active_worker_control_grants(&workspace_id, &controller, 10) .unwrap(); assert_eq!(grants.len(), 1); - assert_eq!(grants[0].subject, first.worker_ref); + assert_eq!( + grants[0].subject, + RuntimeWorkerRef::new(&first.runtime_id, &first.worker_id) + ); assert_eq!(grants[0].operation_id, "control-spawn-retry"); } @@ -17921,7 +18392,6 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await @@ -17939,7 +18409,6 @@ mod tests { initial_submit: Vec::new(), working_directory: None, control_operation_id: None, - resolved_control_operation: None, }), ) .await @@ -17960,14 +18429,16 @@ mod tests { dedicated.singleton_key.as_deref(), Some(crate::hosts::WORKSPACE_ORCHESTRATOR_SINGLETON_KEY) ); - assert_ne!(dedicated.worker.worker_id, generic.worker_ref.worker_id); + let dedicated_ref = RuntimeWorkerRef::new(&dedicated.runtime_id, &dedicated.worker_id); + let generic_ref = RuntimeWorkerRef::new(&generic.runtime_id, &generic.worker_id); + assert_ne!(dedicated.worker_id, generic.worker_id); api.store .create_worker_control_grant(&WorkerControlGrantRecord { workspace_id: workspace_id.clone(), grant_id: "orchestrator-controls-generic".to_string(), - controller: dedicated.worker.clone(), - subject: generic.worker_ref.clone(), + controller: dedicated_ref.clone(), + subject: generic_ref.clone(), relation: "spawned".to_string(), origin: "test".to_string(), permissions: vec!["observe".to_string()], @@ -17980,11 +18451,11 @@ mod tests { let mut observation_headers = HeaderMap::new(); observation_headers.insert( "x-yoi-runtime-id", - axum::http::HeaderValue::from_str(&dedicated.worker.runtime_id).unwrap(), + axum::http::HeaderValue::from_str(&dedicated.runtime_id).unwrap(), ); observation_headers.insert( "x-yoi-worker-id", - axum::http::HeaderValue::from_str(&dedicated.worker.worker_id).unwrap(), + axum::http::HeaderValue::from_str(&dedicated.worker_id).unwrap(), ); let Json(known) = list_known_workers( State(api.clone()), @@ -17996,7 +18467,7 @@ mod tests { .await .unwrap(); assert_eq!(known.items.len(), 1); - assert_eq!(known.items[0].subject, generic.worker_ref); + assert_eq!(known.items[0].subject, generic_ref); assert_eq!(known.items[0].permissions, ["observe"]); let Json(sessions) = scoped_list_worker_observation_sessions( @@ -18015,8 +18486,8 @@ mod tests { .iter() .any(|session| { session["subject"]["kind"] == "runtime_worker" - && session["subject"]["runtime_id"] == generic.worker_ref.runtime_id - && session["subject"]["worker_id"] == generic.worker_ref.worker_id + && session["subject"]["runtime_id"] == generic.runtime_id + && session["subject"]["worker_id"] == generic.worker_id }) ); let Json(capture) = scoped_capture_worker_observation_session( @@ -18026,8 +18497,8 @@ mod tests { }), observation_headers.clone(), Json(WorkerObservationSubjectRef::RuntimeWorker { - runtime_id: generic.worker_ref.runtime_id.clone(), - worker_id: generic.worker_ref.worker_id.clone(), + runtime_id: generic.runtime_id.clone(), + worker_id: generic.worker_id.clone(), }), ) .await @@ -18050,8 +18521,8 @@ mod tests { }), observation_headers.clone(), Json(WorkerObservationSubjectRef::RuntimeWorker { - runtime_id: generic.worker_ref.runtime_id.clone(), - worker_id: generic.worker_ref.worker_id.clone(), + runtime_id: generic.runtime_id.clone(), + worker_id: generic.worker_id.clone(), }), ) .await @@ -18064,11 +18535,11 @@ mod tests { let mut unauthorized_headers = HeaderMap::new(); unauthorized_headers.insert( "x-yoi-runtime-id", - axum::http::HeaderValue::from_str(&generic.worker_ref.runtime_id).unwrap(), + axum::http::HeaderValue::from_str(&generic.runtime_id).unwrap(), ); unauthorized_headers.insert( "x-yoi-worker-id", - axum::http::HeaderValue::from_str(&generic.worker_ref.worker_id).unwrap(), + axum::http::HeaderValue::from_str(&generic.worker_id).unwrap(), ); let Json(unauthorized) = scoped_list_worker_observation_sessions( State(api.clone()), @@ -18090,10 +18561,7 @@ mod tests { .await .unwrap(); assert_eq!(existing.disposition, "existing"); - assert_eq!( - existing.worker.unwrap().worker.worker_id, - dedicated.worker.worker_id - ); + assert_eq!(existing.worker.unwrap().worker_id, dedicated.worker_id); let Json(status) = scoped_workspace_orchestrator_status( State(api), @@ -18101,10 +18569,7 @@ mod tests { ) .await .unwrap(); - assert_eq!( - status.worker.unwrap().worker.worker_id, - dedicated.worker.worker_id - ); + assert_eq!(status.worker.unwrap().worker_id, dedicated.worker_id); } #[tokio::test] @@ -19908,7 +20373,8 @@ mod tests { ) .await .unwrap(); - let orchestrator = started.worker.unwrap().worker; + let orchestrator = started.worker.unwrap(); + let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id); execution.take_inputs(); let mut input = ticket::NewTicket::new("Bounded notification"); @@ -21019,8 +21485,8 @@ mod tests { .unwrap() .0 .worker - .expect("Workspace Orchestrator should be available") - .worker; + .expect("Workspace Orchestrator should be available"); + let orchestrator = RuntimeWorkerRef::new(&orchestrator.runtime_id, &orchestrator.worker_id); let _ = execution.take_inputs(); let backend = browser_ticket_backend(&api).unwrap(); let mut input = ticket::NewTicket::new("Queued notification"); @@ -21216,7 +21682,7 @@ mod tests { Some(ticket_ref.id.as_str()) ); *api.orchestrator_attention_fingerprint.lock().unwrap() = None; - let worker_id = started.worker.as_ref().unwrap().worker.worker_id.clone(); + let worker_id = started.worker.as_ref().unwrap().worker_id.clone(); maybe_dispatch_orchestrator_turn_end( &api, &worker_id, @@ -21347,9 +21813,9 @@ mod tests { resolved_config_bundle: None, resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }; let Json(first) = scoped_create_runtime_worker( State(api.clone()), @@ -21512,7 +21978,6 @@ mod tests { ticket_id: second_ticket.id.clone(), operation_id: "pending-spawn-operation".to_string(), }), - resolved_control_operation: None, ..request }; pending_request.resolved_workspace_api = @@ -21630,9 +22095,9 @@ mod tests { resolved_config_bundle: None, resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }; let Json(created) = scoped_create_runtime_worker( State(api.clone()), @@ -22340,7 +22805,8 @@ mod tests { ) .await .unwrap(); - let source = orchestrator.worker.unwrap().worker; + let source = orchestrator.worker.unwrap(); + let source = RuntimeWorkerRef::new(&source.runtime_id, &source.worker_id); let verified_source = || crate::worker_source::VerifiedWorkerMutationSource { runtime_id: source.runtime_id.clone(), worker_id: source.worker_id.clone(), @@ -22416,7 +22882,8 @@ mod tests { ) .await .unwrap(); - let source = orchestrator.worker.unwrap().worker; + let source = orchestrator.worker.unwrap(); + let source = RuntimeWorkerRef::new(&source.runtime_id, &source.worker_id); let spawned = api .spawn_workspace_worker( @@ -26939,9 +27406,9 @@ mod tests { resolved_config_bundle: Some(runtime_test_bundle()), resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), - resolved_control_operation: None, resolved_workspace_api: None, resolved_memory_settings: None, + resolved_control_operation: None, }; let spawned = api .spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request) diff --git a/web/workspace/deno.json b/web/workspace/deno.json index 28c50342..76398e59 100644 --- a/web/workspace/deno.json +++ b/web/workspace/deno.json @@ -6,7 +6,7 @@ "dev": "deno run -A npm:vite@7.2.7 dev", "dev:backend": "cd ../.. && cargo run -p yoi-workspace-server --bin yoi-server -- serve --listen 127.0.0.1:8787", "check": "deno run -A npm:@sveltejs/kit@2.49.4 sync && deno run -A npm:svelte-check@4.3.4 --tsconfig ./tsconfig.json", - "test": "deno test --allow-read=src,test,tests --allow-env=LOG,VSCODE_TEXTMATE_DEBUG,NODE_ENV tests/workspace-model.test.ts tests/workspace-catalog.test.ts tests/profile-api.test.ts src/lib/workspace/auth/model.test.ts src/lib/workspace/api/http.test.ts src/lib/workspace/header/breadcrumb-model.test.ts src/lib/workspace/console/chat-submit.test.ts test/composer-history.test.ts tests/composer-paste.test.ts src/lib/workspace/console/composer-command.test.ts src/lib/workspace/console/composer-draft.test.ts src/lib/workspace/console/composer-completion.test.ts src/lib/workspace/console/markdown.test.ts test/console/ansi.test.ts src/lib/workspace/console/model.test.ts src/lib/workspace/companion/api.test.ts tests/workdir-api.test.ts src/lib/workspace/console/tasks.test.ts test/ticket-detail-route-reuse.test.ts test/repositories/ui.test.ts src/lib/workspace/console/worker-console.ui.test.ts src/lib/workspace/settings/model.test.ts src/lib/workspace/sidebar/override-stack.test.ts src/lib/workspace/sidebar/workers.test.ts src/lib/workspace/sidebar/workspace-switcher.test.ts src/lib/workspace/sidebar/worker-subscription.test.ts src/lib/workspace/sidebar/worker-launch.test.ts test/sidebar/worker-actions.test.ts src/lib/workspace/tickets/merge-request-resources.test.ts src/lib/workspace/tickets/ticket-panel.test.ts test/merge-request-status.test.ts test/config-source/decodal-grammar.test.ts test/config-source/editor-state.test.ts test/config-source/fixed-schema-wrapper.test.ts test/config-source/toolchain.test.ts test/config-source/wasm-parity.test.ts test/repository-access/api.test.ts test/repository-access/loader.test.ts test/repository-access/ui.test.ts", + "test": "deno test --allow-read=src,test,tests --allow-env=LOG,VSCODE_TEXTMATE_DEBUG,NODE_ENV tests/workspace-model.test.ts tests/workspace-catalog.test.ts tests/profile-api.test.ts src/lib/workspace/auth/model.test.ts src/lib/workspace/api/http.test.ts src/lib/workspace/api/workers.test.ts src/lib/workspace/header/breadcrumb-model.test.ts src/lib/workspace/console/chat-submit.test.ts test/composer-history.test.ts tests/composer-paste.test.ts src/lib/workspace/console/composer-command.test.ts src/lib/workspace/console/composer-draft.test.ts src/lib/workspace/console/composer-completion.test.ts src/lib/workspace/console/markdown.test.ts test/console/ansi.test.ts src/lib/workspace/console/model.test.ts src/lib/workspace/companion/api.test.ts tests/workdir-api.test.ts src/lib/workspace/console/tasks.test.ts test/ticket-detail-route-reuse.test.ts test/repositories/ui.test.ts src/lib/workspace/console/worker-console.ui.test.ts src/lib/workspace/settings/model.test.ts src/lib/workspace/sidebar/override-stack.test.ts src/lib/workspace/sidebar/workers.test.ts src/lib/workspace/sidebar/workspace-switcher.test.ts src/lib/workspace/sidebar/worker-subscription.test.ts src/lib/workspace/sidebar/worker-launch.test.ts test/sidebar/worker-actions.test.ts src/lib/workspace/tickets/merge-request-resources.test.ts src/lib/workspace/tickets/ticket-panel.test.ts test/merge-request-status.test.ts test/config-source/decodal-grammar.test.ts test/config-source/editor-state.test.ts test/config-source/fixed-schema-wrapper.test.ts test/config-source/toolchain.test.ts test/config-source/wasm-parity.test.ts test/repository-access/api.test.ts test/repository-access/loader.test.ts test/repository-access/ui.test.ts", "build": "deno run -A npm:vite@7.2.7 build", "preview": "deno run -A npm:vite@7.2.7 preview" }, diff --git a/web/workspace/src/lib/generated/ticket-api.ts b/web/workspace/src/lib/generated/ticket-api.ts index dc366555..1901ee15 100644 --- a/web/workspace/src/lib/generated/ticket-api.ts +++ b/web/workspace/src/lib/generated/ticket-api.ts @@ -25,7 +25,9 @@ export type TicketActionEligibility = { can_assign_orchestrator: boolean, can_un export type TicketMergeRequestSummary = { merge_request_id: string, repository_key: string, state: string, review_status: string, selector_from: string | null, selector_to: string, updated_at: string, current_subject_ref: string | null, review_subject_ref: string | null, review_requested_at: string | null, review_submitted_at: string | null, review_excerpt: string | null, }; -export type MergeRequestListItem = { summary: TicketMergeRequestSummary, ticket_ids: Array, thread_event_count: number, }; +export type MergeRequestRefDiagnostic = { code: string, message: string, }; + +export type MergeRequestListItem = { summary: TicketMergeRequestSummary, ticket_ids: Array, thread_event_count: number, ref_diagnostics?: Array, }; export type MergeRequestListResponse = { items: Array, next_cursor: string | null, }; diff --git a/web/workspace/src/lib/generated/worker-launch-api.ts b/web/workspace/src/lib/generated/worker-launch-api.ts new file mode 100644 index 00000000..7fe002b6 --- /dev/null +++ b/web/workspace/src/lib/generated/worker-launch-api.ts @@ -0,0 +1,185 @@ +// Generated from workspace-api. Do not edit by hand. +// Regenerate: cargo run -q -p workspace-api --features typescript --example generate_worker_launch_api_types > web/workspace/src/lib/generated/worker-launch-api.ts + +import type { Segment } from "./protocol"; + +export type DiagnosticSeverity = "info" | "warning" | "error"; + +export type Diagnostic = { + code: string; + severity: DiagnosticSeverity; + message: string; +}; + +export type WorkingDirectoryMaterializerKind = + | "runtime_git_cache" + | "local_git_worktree"; + +export type WorkingDirectoryStatusKind = + | "active" + | "cleanup_pending" + | "corrupted" + | "not_found" + | "unknown"; + +export type WorkingDirectoryCleanupTarget = { + kind: string; + working_directory_id: string; + repository_key: string; +}; + +export type RuntimeWorkingDirectoryCleanupTarget = { + kind: string; + working_directory_id: string; + repository_id: string; +}; + +export type RuntimeWorkingDirectorySummary = { + working_directory_id: string; + repository_id: string; + creation_selector?: string | null; + creation_ref?: string | null; + creation_tree?: string | null; + current_selector?: string | null; + current_ref?: string | null; + current_tree?: string | null; + observed_at_epoch_seconds?: number | null; + materializer_kind: WorkingDirectoryMaterializerKind; + cleanup_target?: RuntimeWorkingDirectoryCleanupTarget | null; + status: WorkingDirectoryStatusKind; + cleanliness?: string | null; + primary_worker_id?: string | null; + occupied_by?: WorkingDirectoryOccupancy | null; +}; + +export type WorkingDirectoryOccupancy = { + runtime_id: string; + worker_id: string; + display_name: string; + linked_at: string; +}; + +export type WorkingDirectorySummary = { + working_directory_id: string; + repository_key: string; + creation_selector?: string | null; + creation_ref?: string | null; + creation_tree?: string | null; + current_selector?: string | null; + current_ref?: string | null; + current_tree?: string | null; + observed_at_epoch_seconds?: number | null; + materializer_kind: WorkingDirectoryMaterializerKind; + cleanup_target?: WorkingDirectoryCleanupTarget | null; + status: WorkingDirectoryStatusKind; + cleanliness?: string | null; + primary_worker_id?: string | null; + occupied_by?: WorkingDirectoryOccupancy | null; +}; + +export type WorkerWorkspaceSummary = { + visibility: string; + identity: string; + workspace_id?: string | null; +}; + +export type WorkerImplementationSummary = { + kind: string; + display_hint: string; +}; + +export type WorkerCapabilitySummary = { + can_stop: boolean; + can_spawn_followup: boolean; +}; + +export type WorkerLaunchWorkerSummary = { + runtime_id: string; + worker_id: string; + host_id: string; + display_name: string; + label: string; + profile: string | null; + singleton_key: string | null; + tags: Array; + workspace: WorkerWorkspaceSummary; + state: string; + last_seen_at: string | null; + pinned: boolean; + retention_state: string; + implementation: WorkerImplementationSummary; + capabilities: WorkerCapabilitySummary; + working_directory?: RuntimeWorkingDirectorySummary | null; + diagnostics: Array; +}; + +export type WorkerLaunchRuntimeOption = { + runtime_id: string; + display_name: string; + built_in: boolean; + worker_creation_available: boolean; + working_directory_required: boolean; + status: string; + diagnostics: Array; +}; + +export type WorkerLaunchProfileCandidate = { + id: string; + label: string; + description: string; +}; + +export type WorkingDirectoryRepositoryOption = { + repository_key: string; + default_selector?: string | null; +}; + +export type WorkerLaunchOptionsResponse = { + workspace_id: string; + runtimes: Array; + default_profile: string | null; + profiles: Array; + repositories: Array; + working_directories: Array; + diagnostics: Array; +}; + +export type BrowserWorkerWorkingDirectorySelection = { + working_directory_id: string; + relative_cwd: string | null; +}; + +export type CreateWorkspaceWorkerTicketAssignmentRequest = { + ticket_id: string; + operation_id: string; +}; + +export type CreateWorkspaceWorkerRequest = { + runtime_id: string; + display_name: string; + profile: string | null; + ticket_assignment: CreateWorkspaceWorkerTicketAssignmentRequest | null; + initial_submit: Array; + working_directory: BrowserWorkerWorkingDirectorySelection | null; + /** + * Backend idempotency key used only for authenticated Worker-owned spawn/control. + */ + control_operation_id: string | null; +}; + +export type BrowserCreateWorkerResponse = { + workspace_id: string; + runtime_id: string; + worker_id: string; + console_href: string; + worker: WorkerLaunchWorkerSummary; + diagnostics: Array; +}; + +export type BrowserWorkspaceOrchestratorResponse = { + workspace_id: string; + online: boolean; + disposition: string; + worker?: WorkerLaunchWorkerSummary | null; + diagnostics: Array; +}; diff --git a/web/workspace/src/lib/workspace/api/workdirs.ts b/web/workspace/src/lib/workspace/api/workdirs.ts index 4381ab27..d8769eb6 100644 --- a/web/workspace/src/lib/workspace/api/workdirs.ts +++ b/web/workspace/src/lib/workspace/api/workdirs.ts @@ -55,7 +55,7 @@ export function parseWorkingDirectoryListResponse( ); return { workspace_id: stringField(record, "workspace_id"), - items: arrayField(record, "items").map(parseSummary), + items: arrayField(record, "items").map(parseWorkingDirectorySummary), diagnostics: arrayField(record, "diagnostics").map(parseDiagnostic), }; } @@ -101,12 +101,14 @@ function parseDetailLike( return { workspace_id: stringField(record, "workspace_id"), runtime_id: stringField(record, "runtime_id"), - item: parseSummary(record.item), + item: parseWorkingDirectorySummary(record.item), diagnostics: arrayField(record, "diagnostics").map(parseDiagnostic), }; } -function parseSummary(value: unknown): WorkingDirectorySummary { +export function parseWorkingDirectorySummary( + value: unknown, +): WorkingDirectorySummary { const record = exactRecord(value, SUMMARY_KEYS, "Workdir summary"); const summary: WorkingDirectorySummary = { working_directory_id: stringField(record, "working_directory_id"), diff --git a/web/workspace/src/lib/workspace/api/workers.test.ts b/web/workspace/src/lib/workspace/api/workers.test.ts new file mode 100644 index 00000000..f80f632c --- /dev/null +++ b/web/workspace/src/lib/workspace/api/workers.test.ts @@ -0,0 +1,183 @@ +declare const Deno: { + test(name: string, fn: () => Promise | void): void; +}; + +function assertEquals(actual: unknown, expected: unknown): void { + const actualJson = JSON.stringify(actual); + const expectedJson = JSON.stringify(expected); + if (actualJson !== expectedJson) { + throw new Error(`Expected ${expectedJson}, received ${actualJson}`); + } +} + +function assertThrows( + operation: () => unknown, + errorClass: typeof Error, + message: string, +): void { + try { + operation(); + } catch (error) { + if (!(error instanceof errorClass) || !error.message.includes(message)) { + throw error; + } + return; + } + throw new Error(`Expected operation to throw ${errorClass.name}: ${message}`); +} + +import { + parseBrowserCreateWorkerResponse, + parseBrowserWorkspaceOrchestratorResponse, + parseCreateWorkspaceWorkerRequest, + parseWorkerLaunchOptionsResponse, +} from "./workers.ts"; + +const worker = { + runtime_id: "runtime-a", + worker_id: "worker-a", + host_id: "host-a", + display_name: "Worker A", + label: "worker-a", + profile: "builtin:coder", + singleton_key: null, + tags: [], + workspace: { + visibility: "workspace", + identity: "workspace-a", + workspace_id: "workspace-a", + }, + state: "idle", + last_seen_at: null, + pinned: false, + retention_state: "active", + implementation: { + kind: "runtime", + display_hint: "Runtime Worker", + }, + capabilities: { + can_stop: true, + can_spawn_followup: false, + }, + diagnostics: [], +}; + +Deno.test("Worker launch options parser accepts the generated wire shape", () => { + const parsed = parseWorkerLaunchOptionsResponse({ + workspace_id: "workspace-a", + runtimes: [{ + runtime_id: "runtime-a", + display_name: "Runtime A", + built_in: false, + worker_creation_available: true, + working_directory_required: true, + status: "connected", + diagnostics: [], + }], + default_profile: null, + profiles: [{ id: "builtin:coder", label: "Coder", description: "Code" }], + repositories: [{ repository_key: "main" }], + working_directories: [], + diagnostics: [], + }); + + assertEquals(parsed.runtimes[0].runtime_id, "runtime-a"); + assertEquals(parsed.repositories[0].default_selector, undefined); +}); + +Deno.test("Worker launch response parsers reject missing and unknown fields", () => { + assertThrows( + () => + parseWorkerLaunchOptionsResponse({ + workspace_id: "workspace-a", + runtimes: [], + profiles: [], + repositories: [], + working_directories: [], + diagnostics: [], + }), + Error, + "default_profile", + ); + + assertThrows( + () => + parseBrowserCreateWorkerResponse({ + workspace_id: "workspace-a", + runtime_id: "runtime-a", + worker_id: "worker-a", + console_href: "/workers/worker-a", + worker, + diagnostics: [], + unexpected: true, + }), + Error, + "unknown field unexpected", + ); + + assertThrows( + () => + parseBrowserWorkspaceOrchestratorResponse({ + workspace_id: "workspace-a", + online: false, + disposition: "missing", + diagnostics: [], + extra: false, + }), + Error, + "unknown field extra", + ); +}); + +Deno.test("Worker create request parser requires the complete shared request", () => { + const request = { + runtime_id: "runtime-a", + display_name: "Worker A", + profile: "builtin:coder", + ticket_assignment: null, + initial_submit: [ + { kind: "text", content: "Implement T-565." }, + { kind: "flow", selector: "builtin:coder-review" }, + ], + working_directory: { + working_directory_id: "workdir-a", + relative_cwd: null, + }, + control_operation_id: null, + }; + + assertEquals(parseCreateWorkspaceWorkerRequest(request), request); + assertThrows( + () => + parseCreateWorkspaceWorkerRequest({ + ...request, + operation_id: "legacy-literal", + }), + Error, + "unknown field operation_id", + ); + const { initial_submit: _initialSubmit, ...missingInitialSubmit } = request; + assertThrows( + () => parseCreateWorkspaceWorkerRequest(missingInitialSubmit), + Error, + "initial_submit", + ); + assertThrows( + () => + parseCreateWorkspaceWorkerRequest({ + ...request, + initial_submit: [{ kind: "flow" }], + }), + Error, + "selector", + ); + assertThrows( + () => + parseCreateWorkspaceWorkerRequest({ + ...request, + initial_submit: [{ kind: "newer_client_segment" }], + }), + Error, + "kind is invalid", + ); +}); diff --git a/web/workspace/src/lib/workspace/api/workers.ts b/web/workspace/src/lib/workspace/api/workers.ts new file mode 100644 index 00000000..af5c4066 --- /dev/null +++ b/web/workspace/src/lib/workspace/api/workers.ts @@ -0,0 +1,665 @@ +import type { + BrowserCreateWorkerResponse, + BrowserWorkerWorkingDirectorySelection, + BrowserWorkspaceOrchestratorResponse, + CreateWorkspaceWorkerRequest, + CreateWorkspaceWorkerTicketAssignmentRequest, + Diagnostic, + DiagnosticSeverity, + RuntimeWorkingDirectoryCleanupTarget, + RuntimeWorkingDirectorySummary, + WorkerCapabilitySummary, + WorkerImplementationSummary, + WorkerLaunchOptionsResponse, + WorkerLaunchProfileCandidate, + WorkerLaunchRuntimeOption, + WorkerLaunchWorkerSummary, + WorkerWorkspaceSummary, + WorkingDirectoryRepositoryOption, +} from "$lib/generated/worker-launch-api"; +import type { Segment } from "$lib/generated/protocol"; +import { parseWorkingDirectorySummary } from "$lib/workspace/api/workdirs"; + +const DIAGNOSTIC_SEVERITIES = new Set([ + "info", + "warning", + "error", +]); + +function record(value: unknown, label: string): Record { + if (typeof value !== "object" || value === null || Array.isArray(value)) { + throw new Error(`${label} must be an object`); + } + return value as Record; +} + +function exact( + value: Record, + allowed: readonly string[], + label: string, +): void { + const unexpected = Object.keys(value).filter((key) => !allowed.includes(key)); + if (unexpected.length > 0) { + throw new Error(`${label} contains unknown field ${unexpected[0]}`); + } +} + +function string(value: unknown, label: string): string { + if (typeof value !== "string") throw new Error(`${label} must be a string`); + return value; +} + +function boolean(value: unknown, label: string): boolean { + if (typeof value !== "boolean") throw new Error(`${label} must be a boolean`); + return value; +} + +function number(value: unknown, label: string): number { + if (typeof value !== "number" || !Number.isFinite(value)) { + throw new Error(`${label} must be a finite number`); + } + return value; +} + +function nullableString(value: unknown, label: string): string | null { + return value === null ? null : string(value, label); +} + +function array( + value: unknown, + label: string, + parse: (item: unknown, label: string) => T, +): T[] { + if (!Array.isArray(value)) throw new Error(`${label} must be an array`); + return value.map((item, index) => parse(item, `${label}[${index}]`)); +} + +function optional( + value: unknown, + label: string, + parse: (item: unknown, label: string) => T, +): T | null | undefined { + return value === undefined + ? undefined + : value === null + ? null + : parse(value, label); +} + +function diagnostic(value: unknown, label: string): Diagnostic { + const item = record(value, label); + exact(item, ["code", "severity", "message"], label); + const severity = string(item.severity, `${label}.severity`); + if (!DIAGNOSTIC_SEVERITIES.has(severity as DiagnosticSeverity)) { + throw new Error(`${label}.severity is invalid`); + } + return { + code: string(item.code, `${label}.code`), + severity: severity as DiagnosticSeverity, + message: string(item.message, `${label}.message`), + }; +} + +function runtimeOption( + value: unknown, + label: string, +): WorkerLaunchRuntimeOption { + const item = record(value, label); + exact( + item, + [ + "runtime_id", + "display_name", + "built_in", + "worker_creation_available", + "working_directory_required", + "status", + "diagnostics", + ], + label, + ); + return { + runtime_id: string(item.runtime_id, `${label}.runtime_id`), + display_name: string(item.display_name, `${label}.display_name`), + built_in: boolean(item.built_in, `${label}.built_in`), + worker_creation_available: boolean( + item.worker_creation_available, + `${label}.worker_creation_available`, + ), + working_directory_required: boolean( + item.working_directory_required, + `${label}.working_directory_required`, + ), + status: string(item.status, `${label}.status`), + diagnostics: array(item.diagnostics, `${label}.diagnostics`, diagnostic), + }; +} + +function profileCandidate( + value: unknown, + label: string, +): WorkerLaunchProfileCandidate { + const item = record(value, label); + exact(item, ["id", "label", "description"], label); + return { + id: string(item.id, `${label}.id`), + label: string(item.label, `${label}.label`), + description: string(item.description, `${label}.description`), + }; +} + +function repositoryOption( + value: unknown, + label: string, +): WorkingDirectoryRepositoryOption { + const item = record(value, label); + exact(item, ["repository_key", "default_selector"], label); + return { + repository_key: string(item.repository_key, `${label}.repository_key`), + default_selector: optional( + item.default_selector, + `${label}.default_selector`, + string, + ), + }; +} + +export function parseWorkerLaunchOptionsResponse( + value: unknown, +): WorkerLaunchOptionsResponse { + const item = record(value, "Worker launch options response"); + exact( + item, + [ + "workspace_id", + "runtimes", + "default_profile", + "profiles", + "repositories", + "working_directories", + "diagnostics", + ], + "Worker launch options response", + ); + return { + workspace_id: string(item.workspace_id, "workspace_id"), + runtimes: array(item.runtimes, "runtimes", runtimeOption), + default_profile: nullableString(item.default_profile, "default_profile"), + profiles: array(item.profiles, "profiles", profileCandidate), + repositories: array(item.repositories, "repositories", repositoryOption), + working_directories: array( + item.working_directories, + "working_directories", + parseWorkingDirectorySummary, + ), + diagnostics: array(item.diagnostics, "diagnostics", diagnostic), + }; +} + +function workspaceSummary( + value: unknown, + label: string, +): WorkerWorkspaceSummary { + const item = record(value, label); + exact(item, ["visibility", "identity", "workspace_id"], label); + return { + visibility: string(item.visibility, `${label}.visibility`), + identity: string(item.identity, `${label}.identity`), + workspace_id: optional(item.workspace_id, `${label}.workspace_id`, string), + }; +} + +function implementationSummary( + value: unknown, + label: string, +): WorkerImplementationSummary { + const item = record(value, label); + exact(item, ["kind", "display_hint"], label); + return { + kind: string(item.kind, `${label}.kind`), + display_hint: string(item.display_hint, `${label}.display_hint`), + }; +} + +function capabilitySummary( + value: unknown, + label: string, +): WorkerCapabilitySummary { + const item = record(value, label); + exact(item, ["can_stop", "can_spawn_followup"], label); + return { + can_stop: boolean(item.can_stop, `${label}.can_stop`), + can_spawn_followup: boolean( + item.can_spawn_followup, + `${label}.can_spawn_followup`, + ), + }; +} + +function runtimeCleanupTarget( + value: unknown, + label: string, +): RuntimeWorkingDirectoryCleanupTarget { + const item = record(value, label); + exact(item, ["kind", "working_directory_id", "repository_id"], label); + return { + kind: string(item.kind, `${label}.kind`), + working_directory_id: string( + item.working_directory_id, + `${label}.working_directory_id`, + ), + repository_id: string(item.repository_id, `${label}.repository_id`), + }; +} + +function runtimeWorkingDirectory( + value: unknown, + label: string, +): RuntimeWorkingDirectorySummary { + const item = record(value, label); + exact( + item, + [ + "working_directory_id", + "repository_id", + "creation_selector", + "creation_ref", + "creation_tree", + "current_selector", + "current_ref", + "current_tree", + "observed_at_epoch_seconds", + "materializer_kind", + "cleanup_target", + "status", + "cleanliness", + "primary_worker_id", + "occupied_by", + ], + label, + ); + const materializerKind = string( + item.materializer_kind, + `${label}.materializer_kind`, + ); + if ( + materializerKind !== "runtime_git_cache" && + materializerKind !== "local_git_worktree" + ) { + throw new Error(`${label}.materializer_kind is invalid`); + } + const status = string(item.status, `${label}.status`); + if ( + !["active", "cleanup_pending", "corrupted", "not_found", "unknown"] + .includes(status) + ) { + throw new Error(`${label}.status is invalid`); + } + const occupied = optional( + item.occupied_by, + `${label}.occupied_by`, + (value, occupiedLabel) => { + const occupancy = record(value, occupiedLabel); + exact( + occupancy, + ["runtime_id", "worker_id", "display_name", "linked_at"], + occupiedLabel, + ); + return { + runtime_id: string(occupancy.runtime_id, `${occupiedLabel}.runtime_id`), + worker_id: string(occupancy.worker_id, `${occupiedLabel}.worker_id`), + display_name: string( + occupancy.display_name, + `${occupiedLabel}.display_name`, + ), + linked_at: string(occupancy.linked_at, `${occupiedLabel}.linked_at`), + }; + }, + ); + return { + working_directory_id: string( + item.working_directory_id, + `${label}.working_directory_id`, + ), + repository_id: string(item.repository_id, `${label}.repository_id`), + creation_selector: optional( + item.creation_selector, + `${label}.creation_selector`, + string, + ), + creation_ref: optional(item.creation_ref, `${label}.creation_ref`, string), + creation_tree: optional( + item.creation_tree, + `${label}.creation_tree`, + string, + ), + current_selector: optional( + item.current_selector, + `${label}.current_selector`, + string, + ), + current_ref: optional(item.current_ref, `${label}.current_ref`, string), + current_tree: optional(item.current_tree, `${label}.current_tree`, string), + observed_at_epoch_seconds: optional( + item.observed_at_epoch_seconds, + `${label}.observed_at_epoch_seconds`, + number, + ), + materializer_kind: materializerKind, + cleanup_target: optional( + item.cleanup_target, + `${label}.cleanup_target`, + runtimeCleanupTarget, + ), + status: status as RuntimeWorkingDirectorySummary["status"], + cleanliness: optional(item.cleanliness, `${label}.cleanliness`, string), + primary_worker_id: optional( + item.primary_worker_id, + `${label}.primary_worker_id`, + string, + ), + occupied_by: occupied, + }; +} + +function workerSummary( + value: unknown, + label: string, +): WorkerLaunchWorkerSummary { + const item = record(value, label); + exact( + item, + [ + "runtime_id", + "worker_id", + "host_id", + "display_name", + "label", + "profile", + "singleton_key", + "tags", + "workspace", + "state", + "last_seen_at", + "pinned", + "retention_state", + "implementation", + "capabilities", + "working_directory", + "diagnostics", + ], + label, + ); + return { + runtime_id: string(item.runtime_id, `${label}.runtime_id`), + worker_id: string(item.worker_id, `${label}.worker_id`), + host_id: string(item.host_id, `${label}.host_id`), + display_name: string(item.display_name, `${label}.display_name`), + label: string(item.label, `${label}.label`), + profile: nullableString(item.profile, `${label}.profile`), + singleton_key: nullableString(item.singleton_key, `${label}.singleton_key`), + tags: array(item.tags, `${label}.tags`, string), + workspace: workspaceSummary(item.workspace, `${label}.workspace`), + state: string(item.state, `${label}.state`), + last_seen_at: nullableString(item.last_seen_at, `${label}.last_seen_at`), + pinned: boolean(item.pinned, `${label}.pinned`), + retention_state: string(item.retention_state, `${label}.retention_state`), + implementation: implementationSummary( + item.implementation, + `${label}.implementation`, + ), + capabilities: capabilitySummary(item.capabilities, `${label}.capabilities`), + working_directory: optional( + item.working_directory, + `${label}.working_directory`, + runtimeWorkingDirectory, + ), + diagnostics: array(item.diagnostics, `${label}.diagnostics`, diagnostic), + }; +} + +export function parseBrowserCreateWorkerResponse( + value: unknown, +): BrowserCreateWorkerResponse { + const item = record(value, "Worker create response"); + exact( + item, + [ + "workspace_id", + "runtime_id", + "worker_id", + "console_href", + "worker", + "diagnostics", + ], + "Worker create response", + ); + return { + workspace_id: string(item.workspace_id, "workspace_id"), + runtime_id: string(item.runtime_id, "runtime_id"), + worker_id: string(item.worker_id, "worker_id"), + console_href: string(item.console_href, "console_href"), + worker: workerSummary(item.worker, "worker"), + diagnostics: array(item.diagnostics, "diagnostics", diagnostic), + }; +} + +export function parseBrowserWorkspaceOrchestratorResponse( + value: unknown, +): BrowserWorkspaceOrchestratorResponse { + const item = record(value, "Workspace Orchestrator response"); + exact( + item, + ["workspace_id", "online", "disposition", "worker", "diagnostics"], + "Workspace Orchestrator response", + ); + return { + workspace_id: string(item.workspace_id, "workspace_id"), + online: boolean(item.online, "online"), + disposition: string(item.disposition, "disposition"), + worker: optional(item.worker, "worker", workerSummary), + diagnostics: array(item.diagnostics, "diagnostics", diagnostic), + }; +} + +function workingDirectorySelection( + value: unknown, + label: string, +): BrowserWorkerWorkingDirectorySelection { + const item = record(value, label); + exact(item, ["working_directory_id", "relative_cwd"], label); + return { + working_directory_id: string( + item.working_directory_id, + `${label}.working_directory_id`, + ), + relative_cwd: nullableString(item.relative_cwd, `${label}.relative_cwd`), + }; +} + +function ticketAssignment( + value: unknown, + label: string, +): CreateWorkspaceWorkerTicketAssignmentRequest { + const item = record(value, label); + exact(item, ["ticket_id", "operation_id"], label); + return { + ticket_id: string(item.ticket_id, `${label}.ticket_id`), + operation_id: string(item.operation_id, `${label}.operation_id`), + }; +} + +function unsignedInteger(value: unknown, label: string): number { + const parsed = number(value, label); + if (!Number.isSafeInteger(parsed) || parsed < 0) { + throw new Error(`${label} must be a non-negative safe integer`); + } + return parsed; +} + +function pasteArtifact( + value: unknown, + label: string, +): Extract["artifact"] { + const item = record(value, label); + exact( + item, + [ + "artifact_id", + "created_at_ms", + "media_type", + "availability", + "byte_len", + "char_count", + "line_count", + "sha256", + "source_entry_id", + ], + label, + ); + const mediaType = string(item.media_type, `${label}.media_type`); + if (mediaType !== "text_plain_utf8") { + throw new Error(`${label}.media_type is invalid`); + } + const availability = string(item.availability, `${label}.availability`); + if ( + !["available", "unavailable", "integrity_failed"].includes(availability) + ) { + throw new Error(`${label}.availability is invalid`); + } + return { + artifact_id: string(item.artifact_id, `${label}.artifact_id`), + created_at_ms: unsignedInteger( + item.created_at_ms, + `${label}.created_at_ms`, + ), + media_type: mediaType, + availability: availability as Extract[ + "artifact" + ]["availability"], + byte_len: unsignedInteger(item.byte_len, `${label}.byte_len`), + char_count: unsignedInteger(item.char_count, `${label}.char_count`), + line_count: unsignedInteger(item.line_count, `${label}.line_count`), + sha256: string(item.sha256, `${label}.sha256`), + source_entry_id: string(item.source_entry_id, `${label}.source_entry_id`), + }; +} + +function uploadedFile( + value: unknown, + label: string, +): Extract["file"] { + const item = record(value, label); + exact( + item, + [ + "artifact_id", + "file_name", + "media_type", + "created_at_ms", + "availability", + "byte_len", + "sha256", + "source_entry_id", + ], + label, + ); + const availability = string(item.availability, `${label}.availability`); + if ( + !["available", "unavailable", "integrity_failed"].includes(availability) + ) { + throw new Error(`${label}.availability is invalid`); + } + return { + artifact_id: string(item.artifact_id, `${label}.artifact_id`), + file_name: string(item.file_name, `${label}.file_name`), + media_type: string(item.media_type, `${label}.media_type`), + created_at_ms: unsignedInteger( + item.created_at_ms, + `${label}.created_at_ms`, + ), + availability: availability as Extract[ + "file" + ]["availability"], + byte_len: unsignedInteger(item.byte_len, `${label}.byte_len`), + sha256: string(item.sha256, `${label}.sha256`), + source_entry_id: optional( + item.source_entry_id, + `${label}.source_entry_id`, + string, + ), + }; +} + +function segment(value: unknown, label: string): Segment { + const item = record(value, label); + const kind = string(item.kind, `${label}.kind`) as Segment["kind"]; + switch (kind) { + case "text": + exact(item, ["kind", "content"], label); + return { kind, content: string(item.content, `${label}.content`) }; + case "paste": + exact(item, ["kind", "id", "chars", "lines", "content"], label); + return { + kind, + id: unsignedInteger(item.id, `${label}.id`), + chars: unsignedInteger(item.chars, `${label}.chars`), + lines: unsignedInteger(item.lines, `${label}.lines`), + content: string(item.content, `${label}.content`), + }; + case "paste_artifact": + exact(item, ["kind", "artifact"], label); + return { + kind, + artifact: pasteArtifact(item.artifact, `${label}.artifact`), + }; + case "uploaded_file": + exact(item, ["kind", "file"], label); + return { kind, file: uploadedFile(item.file, `${label}.file`) }; + case "file_ref": + exact(item, ["kind", "path"], label); + return { kind, path: string(item.path, `${label}.path`) }; + case "flow": + exact(item, ["kind", "selector"], label); + return { kind, selector: string(item.selector, `${label}.selector`) }; + case "unknown": + throw new Error(`${label}.kind is not supported by Worker creation`); + } + const exhaustive: never = kind; + throw new Error(`${label}.kind is invalid: ${exhaustive}`); +} + +export function parseCreateWorkspaceWorkerRequest( + value: unknown, +): CreateWorkspaceWorkerRequest { + const item = record(value, "Worker create request"); + exact( + item, + [ + "runtime_id", + "display_name", + "profile", + "ticket_assignment", + "initial_submit", + "working_directory", + "control_operation_id", + ], + "Worker create request", + ); + return { + runtime_id: string(item.runtime_id, "runtime_id"), + display_name: string(item.display_name, "display_name"), + profile: nullableString(item.profile, "profile"), + ticket_assignment: item.ticket_assignment === null + ? null + : ticketAssignment(item.ticket_assignment, "ticket_assignment"), + initial_submit: array(item.initial_submit, "initial_submit", segment), + working_directory: item.working_directory === null + ? null + : workingDirectorySelection(item.working_directory, "working_directory"), + control_operation_id: nullableString( + item.control_operation_id, + "control_operation_id", + ), + }; +} diff --git a/web/workspace/src/lib/workspace/sidebar/types.ts b/web/workspace/src/lib/workspace/sidebar/types.ts index 2076d719..ec758cd7 100644 --- a/web/workspace/src/lib/workspace/sidebar/types.ts +++ b/web/workspace/src/lib/workspace/sidebar/types.ts @@ -1,3 +1,12 @@ +import type { + BrowserCreateWorkerResponse as SharedBrowserCreateWorkerResponse, + BrowserWorkerWorkingDirectorySelection + as SharedBrowserWorkerWorkingDirectorySelection, + WorkerLaunchOptionsResponse as SharedWorkerLaunchOptionsResponse, + WorkerLaunchProfileCandidate as SharedWorkerLaunchProfileCandidate, + WorkerLaunchRuntimeOption as SharedWorkerLaunchRuntimeOption, + WorkingDirectoryRepositoryOption as SharedWorkingDirectoryRepositoryOption, +} from "$lib/generated/worker-launch-api"; import type { WorkingDirectoryCreateRequest, WorkingDirectoryCreateResponse, @@ -101,26 +110,10 @@ export type Worker = { export type WorkerOperationState = "accepted" | "unsupported" | "rejected"; -export type WorkerLaunchRuntimeOption = { - runtime_id: string; - display_name: string; - built_in: boolean; - worker_creation_available: boolean; - working_directory_required: boolean; - status: string; - diagnostics: Diagnostic[]; -}; - -export type WorkerLaunchProfileCandidate = { - id: string; - label: string; - description: string; -}; - -export type WorkingDirectoryRepositoryOption = { - repository_key: string; - default_selector?: string | null; -}; +export type WorkerLaunchRuntimeOption = SharedWorkerLaunchRuntimeOption; +export type WorkerLaunchProfileCandidate = SharedWorkerLaunchProfileCandidate; +export type WorkingDirectoryRepositoryOption = + SharedWorkingDirectoryRepositoryOption; export type CleanupTargetKind = | "worker_delete" @@ -185,29 +178,10 @@ export type RuntimeCleanupExecutionResponse = { diagnostics: Diagnostic[]; }; -export type BrowserWorkerWorkingDirectorySelection = { - working_directory_id: string; - relative_cwd?: string | null; -}; - -export type WorkerLaunchOptionsResponse = { - workspace_id: string; - runtimes: WorkerLaunchRuntimeOption[]; - default_profile?: string | null; - profiles: WorkerLaunchProfileCandidate[]; - repositories: WorkingDirectoryRepositoryOption[]; - working_directories: WorkingDirectorySummary[]; - diagnostics: Diagnostic[]; -}; - -export type BrowserCreateWorkerResponse = { - workspace_id: string; - runtime_id: string; - worker_id: string; - console_href: string; - worker: Worker; - diagnostics: Diagnostic[]; -}; +export type BrowserWorkerWorkingDirectorySelection = + SharedBrowserWorkerWorkingDirectorySelection; +export type WorkerLaunchOptionsResponse = SharedWorkerLaunchOptionsResponse; +export type BrowserCreateWorkerResponse = SharedBrowserCreateWorkerResponse; export type WorkerInputResult = { state: WorkerOperationState; diff --git a/web/workspace/src/lib/workspace/sidebar/worker-launch.test.ts b/web/workspace/src/lib/workspace/sidebar/worker-launch.test.ts index a52ce503..c5dfe6f1 100644 --- a/web/workspace/src/lib/workspace/sidebar/worker-launch.test.ts +++ b/web/workspace/src/lib/workspace/sidebar/worker-launch.test.ts @@ -196,11 +196,13 @@ Deno.test("buildCreateWorkspaceWorkerRequest sends working_directory id and rela runtime_id: "embedded", display_name: "Worker", profile: "builtin:coder", + ticket_assignment: null, initial_submit: [{ kind: "text", content: "go" }], working_directory: { working_directory_id: "wd-1-repo", relative_cwd: "crates/yoi", }, + control_operation_id: null, }); }); @@ -219,7 +221,7 @@ Deno.test("buildCreateWorkspaceWorkerRequest sends no initial segments for an em assertEquals(request.initial_submit, []); }); -Deno.test("buildCreateWorkspaceWorkerRequest omits working_directory for embedded no-workdir launches", () => { +Deno.test("buildCreateWorkspaceWorkerRequest emits null for embedded no-workdir launches", () => { const request = buildCreateWorkspaceWorkerRequest({ runtime_id: "embedded", display_name: "Worker", @@ -235,6 +237,9 @@ Deno.test("buildCreateWorkspaceWorkerRequest omits working_directory for embedde runtime_id: "embedded", display_name: "Worker", profile: "builtin:companion", + ticket_assignment: null, initial_submit: [{ kind: "text", content: "chat" }], + working_directory: null, + control_operation_id: null, }); }); diff --git a/web/workspace/src/lib/workspace/sidebar/worker-launch.ts b/web/workspace/src/lib/workspace/sidebar/worker-launch.ts index 2926bc1e..ed9adbb3 100644 --- a/web/workspace/src/lib/workspace/sidebar/worker-launch.ts +++ b/web/workspace/src/lib/workspace/sidebar/worker-launch.ts @@ -1,9 +1,7 @@ -import type { Segment } from "$lib/generated/protocol"; +import type { CreateWorkspaceWorkerRequest } from "$lib/generated/worker-launch-api"; +import { parseCreateWorkspaceWorkerRequest } from "$lib/workspace/api/workers"; -import type { - BrowserWorkerWorkingDirectorySelection, - WorkerLaunchOptionsResponse, -} from "./types"; +import type { WorkerLaunchOptionsResponse } from "./types"; export type WorkerLaunchFormState = { runtime_id: string; @@ -16,14 +14,6 @@ export type WorkerLaunchFormState = { relative_cwd: string; }; -export type CreateWorkspaceWorkerRequest = { - runtime_id: string; - display_name: string; - profile: string; - initial_submit: Segment[]; - working_directory?: BrowserWorkerWorkingDirectorySelection; -}; - export function defaultWorkerLaunchForm( options: WorkerLaunchOptionsResponse | null, current: WorkerLaunchFormState, @@ -86,7 +76,8 @@ export function defaultWorkerLaunchForm( ) ? current.working_directory_id : preferredWorkingDirectory?.working_directory_id || "", - working_directory_repository_key: current.working_directory_repository_key || + working_directory_repository_key: + current.working_directory_repository_key || preferredRepository?.repository_key || "", working_directory_selector: current.working_directory_selector || preferredRepository?.default_selector || "HEAD", @@ -97,22 +88,21 @@ export function defaultWorkerLaunchForm( export function buildCreateWorkspaceWorkerRequest( form: WorkerLaunchFormState, ): CreateWorkspaceWorkerRequest { - const request: CreateWorkspaceWorkerRequest = { - runtime_id: form.runtime_id, - display_name: form.display_name, - profile: form.profile, - initial_submit: form.initial_text.trim() + const initialMessage = form.initial_text.trim(); + return parseCreateWorkspaceWorkerRequest({ + runtime_id: form.runtime_id.trim(), + display_name: form.display_name.trim(), + profile: form.profile.trim() || null, + ticket_assignment: null, + initial_submit: initialMessage ? [{ kind: "text", content: form.initial_text }] : [], - }; - if (form.working_directory_id) { - request.working_directory = { - working_directory_id: form.working_directory_id, - }; - const relativeCwd = form.relative_cwd.trim(); - if (relativeCwd) { - request.working_directory.relative_cwd = relativeCwd; - } - } - return request; + working_directory: form.working_directory_id + ? { + working_directory_id: form.working_directory_id, + relative_cwd: form.relative_cwd.trim() || null, + } + : null, + control_operation_id: null, + }); } diff --git a/web/workspace/src/lib/workspace/tickets/ticket-panel.ts b/web/workspace/src/lib/workspace/tickets/ticket-panel.ts index 0df64ebc..428a9f1d 100644 --- a/web/workspace/src/lib/workspace/tickets/ticket-panel.ts +++ b/web/workspace/src/lib/workspace/tickets/ticket-panel.ts @@ -1,3 +1,4 @@ +import type { BrowserWorkspaceOrchestratorResponse } from "$lib/generated/worker-launch-api"; import type { TicketDetail, TicketSummary } from "$lib/generated/ticket-api"; export const TICKET_STATES = [ @@ -12,22 +13,7 @@ export const TICKET_STATES = [ export type TicketState = (typeof TICKET_STATES)[number]; export type TicketWorkerRole = "coder" | "reviewer"; -export type WorkspaceOrchestratorStatus = { - workspace_id: string; - online: boolean; - disposition: string; - worker?: { - runtime_id: string; - worker_id: string; - state: string; - display_name: string; - } | null; - diagnostics: Array<{ - code: string; - severity: string; - message: string; - }>; -}; +export type WorkspaceOrchestratorStatus = BrowserWorkspaceOrchestratorResponse; const LANE_DEFINITIONS = [ { diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte index f20a36d1..a0d34adf 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte @@ -2,6 +2,7 @@ import { untrack } from "svelte"; import type { ApiResult } from "$lib/workspace/api/http"; import { loadJson, workspaceApiPath } from "$lib/workspace/api/http"; + import { parseBrowserWorkspaceOrchestratorResponse } from "$lib/workspace/api/workers"; import type { QueryPage, TicketListResponse, @@ -99,6 +100,7 @@ fetch, workspaceApiPath(data.workspaceId, "/orchestrator"), { method: "POST" }, + parseBrowserWorkspaceOrchestratorResponse, ); orchestratorStarting = false; } diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts index 04ab28e8..2717de97 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts @@ -1,3 +1,4 @@ +import { parseBrowserWorkspaceOrchestratorResponse } from "$lib/workspace/api/workers"; import { loadJson, workspaceApiPath } from "$lib/workspace/api/http"; import type { TicketListResponse } from "$lib/generated/ticket-api"; import type { WorkspaceOrchestratorStatus } from "$lib/workspace/tickets/ticket-panel"; @@ -42,6 +43,8 @@ export const load: PageLoad = async ({ fetch, params }) => { loadJson( fetch, workspaceApiPath(workspaceId, "/orchestrator"), + undefined, + parseBrowserWorkspaceOrchestratorResponse, ), ]); diff --git a/web/workspace/src/routes/w/[workspaceId]/workers/new/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/workers/new/+page.svelte index 5ced53cd..bf500233 100644 --- a/web/workspace/src/routes/w/[workspaceId]/workers/new/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/workers/new/+page.svelte @@ -6,10 +6,13 @@ parseWorkingDirectoryCreateResponse, validateWorkingDirectoryCreateRequest, } from '$lib/workspace/api/workdirs'; + import { + parseBrowserCreateWorkerResponse, + parseWorkerLaunchOptionsResponse, + } from '$lib/workspace/api/workers'; import { formatCurrentWorkdirRevision } from '$lib/workspace/settings/workdir-revision'; import { buildCreateWorkspaceWorkerRequest, defaultWorkerLaunchForm } from '$lib/workspace/sidebar/worker-launch'; import type { - BrowserCreateWorkerResponse, Diagnostic, WorkerLaunchOptionsResponse, WorkingDirectorySummary, @@ -115,7 +118,7 @@ if (!response.ok) { throw new Error(`worker launch options failed (${response.status})`); } - const payload = (await response.json()) as WorkerLaunchOptionsResponse; + const payload = parseWorkerLaunchOptionsResponse(await response.json()); options = payload; const form = defaultWorkerLaunchForm(payload, { runtime_id: runtimeId, @@ -229,7 +232,7 @@ submitError = await responseDisplayError(response, 'worker create failed'); return; } - const payload = (await response.json()) as BrowserCreateWorkerResponse; + const payload = parseBrowserCreateWorkerResponse(await response.json()); await goto(payload.console_href); } catch (err) { submitError = exceptionDisplayError(err, 'worker create failed');