From 56798f9fb4bfe18f4fb5986f1a376f253f0cbe8f Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 12:56:06 +0900 Subject: [PATCH 1/6] feat: observe repository refs through runtime providers --- crates/worker-runtime/src/catalog.rs | 24 ++ crates/worker-runtime/src/execution.rs | 18 ++ crates/worker-runtime/src/http_server.rs | 41 ++- crates/worker-runtime/src/runtime.rs | 30 ++- crates/worker-runtime/src/worker_backend.rs | 14 + .../worker-runtime/src/working_directory.rs | 246 +++++++++++++++++- 6 files changed, 366 insertions(+), 7 deletions(-) diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index 2e5d9b9f..14331018 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..81246e5d 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") @@ -2424,6 +2453,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") diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index dad770a2..793f3a37 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -1,6 +1,7 @@ use crate::catalog::{ - ConfigBundleRef, CreateWorkerRequest, ProfileSelector, WorkerDetail, WorkerLifecycleAck, - WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, + ConfigBundleRef, CreateWorkerRequest, ProfileSelector, RepositoryRefObservation, + RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck, WorkerStatus, WorkerSummary, + WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ @@ -392,6 +393,31 @@ 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?; + } + 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, diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index e9141379..192365e6 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..8efe6bbf 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,78 @@ 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_selector(selector)?; + if !selector.starts_with("refs/heads/") { + return Err(WorkingDirectoryDiagnostic::new( + "repository_ref_selector_invalid", + "Repository ref observation requires an exact branch 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 +2100,112 @@ fn repository_git_command( command } +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 +2677,66 @@ 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_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, + selector: "HEAD".to_string(), + materialization: None, + }) + .unwrap_err(); + assert_eq!(non_branch.code, "repository_ref_selector_invalid"); + } + #[test] fn local_git_repo_materializes_detached_worktree_under_runtime_root() { let repo = create_clean_repo(); From 7be428d8bf9ec11ba345d1d19a4d303e63b9634c Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 12:56:17 +0900 Subject: [PATCH 2/6] fix: require provider-published merge request refs --- crates/workspace-server/src/hosts.rs | 56 +++- crates/workspace-server/src/repositories.rs | 2 +- crates/workspace-server/src/server.rs | 294 +++++++++++++++++--- 3 files changed, 306 insertions(+), 46 deletions(-) diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index cc15ba3b..f4419e63 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, @@ -3355,6 +3391,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()), 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 5c04e0af..96c5a500 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -150,10 +150,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::{ @@ -5506,6 +5506,127 @@ 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 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; + if workdir.repository_id != repository_id { + return Err(Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_repository_mismatch: current Coder Workdir belongs to a different Repository".into(), + )) + .into()); + } + if workdir.cleanliness.as_deref() != Some("clean") { + return Err( + Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(), + )) + .into(), + ); + } + 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(Error::MergeRequest(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_head_mismatch: current Coder Workdir selector and HEAD must match the provider-published source ref".into(), + )) + .into()); + } + Ok(()) +} + fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError { Error::InvalidInput(format!( "repository merge evidence validation failed: {error:?}" @@ -5818,12 +5939,26 @@ 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) - }); + 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 = mr + .selector_from + .as_deref() + .map(|selector| { + observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + ) + .map(|observation| observation.revision_ref) + }) + .transpose()?; Ok(Json(store.readiness(merge_request::ReadinessCheck { ticket_id, current_subject_ref, @@ -5874,13 +6009,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_merge_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 +6088,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_merge_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 +6169,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_merge_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, @@ -6046,11 +6214,20 @@ async fn scoped_submit_merge_request_review( .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_merge_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 +6300,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 +6331,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_merge_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, )?; @@ -16346,6 +16530,34 @@ mod tests { SqliteWorkspaceStore, TrustedRuntimeRecord, UserRecord, WorkspaceRecord, }; + fn handler_source<'a>(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_mutations_observe_refs_through_runtime_provider_authority() { + let source = include_str!("server.rs"); + for handler in [ + "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_merge_ref(")); + assert!(!handler.contains("repository_reader()")); + } + } + fn completed_upload_file(sha256: &str) -> protocol::UploadedFileRef { protocol::UploadedFileRef { artifact_id: "019ca7c8-57b6-7f05-8edf-524147aba7b3".into(), From 2884c08466e7a419ac76b655b1974bc071d66562 Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 13:49:20 +0900 Subject: [PATCH 3/6] fix: harden provider ref review boundaries --- crates/merge-request/src/lib.rs | 34 +++ crates/merge-request/tests/store.rs | 17 ++ crates/worker-runtime/src/http_server.rs | 61 +++-- crates/worker-runtime/src/runtime.rs | 68 ++++-- .../worker-runtime/src/working_directory.rs | 52 ++++- crates/workspace-server/src/server.rs | 213 ++++++++++++++---- 6 files changed, 364 insertions(+), 81 deletions(-) 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/http_server.rs b/crates/worker-runtime/src/http_server.rs index 81246e5d..2217e973 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -1970,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 { .. } @@ -1980,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 { .. } @@ -2967,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 793f3a37..d085d2ca 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -3002,20 +3002,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 { @@ -3217,6 +3232,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/working_directory.rs b/crates/worker-runtime/src/working_directory.rs index 8efe6bbf..58890d45 100644 --- a/crates/worker-runtime/src/working_directory.rs +++ b/crates/worker-runtime/src/working_directory.rs @@ -954,13 +954,7 @@ impl WorkingDirectoryMaterializer for RuntimeGitCacheMaterializer { request: &RepositoryRefObservationRequest, ) -> Result { let selector = request.selector.trim(); - validate_selector(selector)?; - if !selector.starts_with("refs/heads/") { - return Err(WorkingDirectoryDiagnostic::new( - "repository_ref_selector_invalid", - "Repository ref observation requires an exact branch selector", - )); - } + validate_exact_branch_selector(selector)?; let working_request = WorkingDirectoryRequest { repository: request.repository.clone(), @@ -2100,6 +2094,40 @@ 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(); @@ -2729,12 +2757,20 @@ mod tests { assert_eq!(missing.code, "repository_ref_not_found"); let non_branch = materializer .observe_repository_ref(&RepositoryRefObservationRequest { - repository, + 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] diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 96c5a500..61b0e9a6 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -5593,19 +5593,26 @@ fn require_assigned_workdir_source( )) })? .summary; + validate_assigned_workdir_source(&workdir, repository_id, selector, revision_ref) + .map_err(Error::MergeRequest)?; + Ok(()) +} + +fn validate_assigned_workdir_source( + workdir: &worker_runtime::catalog::WorkingDirectorySummary, + repository_id: &str, + selector: &str, + revision_ref: &str, +) -> std::result::Result<(), merge_request::MergeRequestError> { if workdir.repository_id != repository_id { - return Err(Error::MergeRequest(merge_request::MergeRequestError::Conflict( + return Err(merge_request::MergeRequestError::Conflict( "merge_request_source_workdir_repository_mismatch: current Coder Workdir belongs to a different Repository".into(), - )) - .into()); + )); } if workdir.cleanliness.as_deref() != Some("clean") { - return Err( - Error::MergeRequest(merge_request::MergeRequestError::Conflict( - "merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(), - )) - .into(), - ); + return Err(merge_request::MergeRequestError::Conflict( + "merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(), + )); } let workdir_selector_matches = workdir .current_selector @@ -5619,10 +5626,9 @@ fn require_assigned_workdir_source( .is_ok_and(|selector| selector == workdir_selector) }); if !workdir_selector_matches || workdir.current_ref.as_deref() != Some(revision_ref) { - return Err(Error::MergeRequest(merge_request::MergeRequestError::Conflict( + return Err(merge_request::MergeRequestError::Conflict( "merge_request_source_workdir_head_mismatch: current Coder Workdir selector and HEAD must match the provider-published source ref".into(), - )) - .into()); + )); } Ok(()) } @@ -5825,17 +5831,41 @@ 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 = + 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)) => { + let assignment = api + .store + .get_current_ticket_coder_assignment(&workspace_id, ticket_id)? + .ok_or_else(|| { + Error::TicketAssignmentConflict( + "Open Merge Request has no current assigned Coder".into(), + ) + })?; + Some( + observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &merge_request.repository_id, + selector, + )? + .revision_ref, + ) + } + _ => None, + } + } else { + None + }; let repository_key = api .store .get_repository(&workspace_id, &merge_request.repository_id)? @@ -5860,6 +5890,18 @@ async fn scoped_list_merge_requests( })) } +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, + } +} + async fn scoped_show_merge_request( State(api): State, AxumPath((workspace_id, merge_request_id)): AxumPath<(String, String)>, @@ -5875,39 +5917,38 @@ 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 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)? + .ok_or_else(|| { + Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) + })?; 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(), - }, - }, + Some(selector) => merge_ref_response(observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + selector, + )?), None => MergeRequestRefResponse { status: "requires_repair".into(), revision_ref: None, - observed_at: observed_at.clone(), - }, - }; - 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, + observed_at: Utc::now().to_rfc3339(), }, }; + let target = merge_ref_response(observe_published_merge_ref( + &api, + &workspace_id, + &assignment.worker.runtime_id, + &mr.repository_id, + &mr.selector_to, + )?); let linked_tickets = mr .ticket_ids .iter() @@ -6209,6 +6250,16 @@ 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 @@ -16543,9 +16594,11 @@ mod tests { } #[test] - fn merge_request_mutations_observe_refs_through_runtime_provider_authority() { + 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", @@ -16555,7 +16608,73 @@ mod tests { let handler = handler_source(source, handler); assert!(handler.contains("observe_published_merge_ref(")); assert!(!handler.contains("repository_reader()")); + assert!(!handler.contains("status: \"unknown\"")); } + let submit = handler_source(source, "scoped_submit_merge_request_review"); + assert!( + submit + .find("authorize_review_submission(") + .is_some_and(|authorization| { + submit + .find("observe_published_merge_ref(") + .is_some_and(|observation| authorization < observation) + }), + "review capability must be validated before provider access", + ); + } + + #[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(merge_request::MergeRequestError::Conflict(message)) + if message.starts_with("merge_request_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(merge_request::MergeRequestError::Conflict(message)) + if message.starts_with("merge_request_source_workdir_head_mismatch:") + )); } fn completed_upload_file(sha256: &str) -> protocol::UploadedFileRef { From da14c82f71ddd78493a6b7e25f613916e845f3e0 Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 14:25:51 +0900 Subject: [PATCH 4/6] fix: expose provider ref blockers without failing reads --- .../worker-runtime/src/working_directory.rs | 83 +++++ crates/workspace-server/src/records.rs | 9 + crates/workspace-server/src/server.rs | 344 +++++++++++++----- 3 files changed, 347 insertions(+), 89 deletions(-) diff --git a/crates/worker-runtime/src/working_directory.rs b/crates/worker-runtime/src/working_directory.rs index 58890d45..797ec679 100644 --- a/crates/worker-runtime/src/working_directory.rs +++ b/crates/worker-runtime/src/working_directory.rs @@ -2740,6 +2740,89 @@ mod tests { ); } + #[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(); diff --git a/crates/workspace-server/src/records.rs b/crates/workspace-server/src/records.rs index 6d0f1690..236d138d 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)] diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 61b0e9a6..6e6127ba 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -118,9 +118,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, @@ -5562,6 +5562,64 @@ fn observe_published_merge_ref( .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, @@ -5593,8 +5651,13 @@ fn require_assigned_workdir_source( )) })? .summary; - validate_assigned_workdir_source(&workdir, repository_id, selector, revision_ref) - .map_err(Error::MergeRequest)?; + 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(()) } @@ -5603,16 +5666,22 @@ fn validate_assigned_workdir_source( repository_id: &str, selector: &str, revision_ref: &str, -) -> std::result::Result<(), merge_request::MergeRequestError> { +) -> std::result::Result<(), worker_runtime::working_directory::WorkingDirectoryDiagnostic> { if workdir.repository_id != repository_id { - return Err(merge_request::MergeRequestError::Conflict( - "merge_request_source_workdir_repository_mismatch: current Coder Workdir belongs to a different Repository".into(), - )); + 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(merge_request::MergeRequestError::Conflict( - "merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(), - )); + 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 @@ -5626,9 +5695,12 @@ fn validate_assigned_workdir_source( .is_ok_and(|selector| selector == workdir_selector) }); if !workdir_selector_matches || workdir.current_ref.as_deref() != Some(revision_ref) { - return Err(merge_request::MergeRequestError::Conflict( - "merge_request_source_workdir_head_mismatch: current Coder Workdir selector and HEAD must match the provider-published source ref".into(), - )); + 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(()) } @@ -5752,6 +5824,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)] @@ -5835,36 +5909,39 @@ async fn scoped_list_merge_requests( .items .into_iter() .map(|merge_request| -> ApiResult { - let current_subject_ref = + 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)) => { - let assignment = api - .store - .get_current_ticket_coder_assignment(&workspace_id, ticket_id)? - .ok_or_else(|| { - Error::TicketAssignmentConflict( - "Open Merge Request has no current assigned Coder".into(), - ) - })?; - Some( - observe_published_merge_ref( - &api, - &workspace_id, - &assignment.worker.runtime_id, - &merge_request.repository_id, - selector, - )? - .revision_ref, - ) - } - _ => None, + (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 + (None, Vec::new()) }; let repository_key = api .store @@ -5881,6 +5958,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::>>()?; @@ -5890,6 +5968,42 @@ 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) @@ -5899,6 +6013,7 @@ fn merge_ref_response(observation: RepositoryRefObservation) -> MergeRequestRefR status: "known".into(), revision_ref: Some(observation.revision_ref), observed_at, + diagnostic: None, } } @@ -5924,31 +6039,47 @@ async fn scoped_show_merge_request( })?; 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 source = match mr.selector_from.as_deref() { - Some(selector) => merge_ref_response(observe_published_merge_ref( + .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: Utc::now().to_rfc3339(), + diagnostic: None, + }, + }; + let target = match assignment.as_ref() { + Some(assignment) => match observe_published_merge_ref( &api, &workspace_id, &assignment.worker.runtime_id, &mr.repository_id, - selector, - )?), - None => MergeRequestRefResponse { - status: "requires_repair".into(), - revision_ref: None, - observed_at: Utc::now().to_rfc3339(), + &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 target = merge_ref_response(observe_published_merge_ref( - &api, - &workspace_id, - &assignment.worker.runtime_id, - &mr.repository_id, - &mr.selector_to, - )?); let linked_tickets = mr .ticket_ids .iter() @@ -5982,25 +6113,26 @@ async fn scoped_merge_request_readiness( let mr = store.get(&workspace_id, &ticket_id)?; 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 = mr - .selector_from - .as_deref() - .map(|selector| { - observe_published_merge_ref( - &api, - &workspace_id, - &assignment.worker.runtime_id, - &mr.repository_id, - selector, - ) - .map(|observation| observation.revision_ref) - }) - .transpose()?; - Ok(Json(store.readiness(merge_request::ReadinessCheck { + .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 { @@ -6010,7 +6142,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( @@ -6050,7 +6192,7 @@ async fn scoped_open_merge_request( ) .into()); } - let source_observation = observe_published_merge_ref( + let source_observation = observe_published_source_ref( &api, &workspace_id, &assignment.worker.runtime_id, @@ -6135,7 +6277,7 @@ async fn scoped_repair_merge_request_selector( .ok_or_else(|| { Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; - let resolved_subject_ref = observe_published_merge_ref( + let resolved_subject_ref = observe_published_source_ref( &api, &workspace_id, &assignment.worker.runtime_id, @@ -6210,7 +6352,7 @@ async fn scoped_register_merge_request_review_capability( .selector_from .as_deref() .ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?; - let source_observation = observe_published_merge_ref( + let source_observation = observe_published_source_ref( &api, &workspace_id, &assignment.worker.runtime_id, @@ -6271,7 +6413,7 @@ async fn scoped_submit_merge_request_review( .ok_or_else(|| { Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; - let current_subject_ref = observe_published_merge_ref( + let current_subject_ref = observe_published_source_ref( &api, &workspace_id, &assignment.worker.runtime_id, @@ -6382,7 +6524,7 @@ async fn scoped_complete_merge_request( .selector_from .as_deref() .ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?; - let current_source_ref = observe_published_merge_ref( + let current_source_ref = observe_published_source_ref( &api, &workspace_id, &assignment.worker.runtime_id, @@ -16606,9 +16748,11 @@ mod tests { "scoped_complete_merge_request", ] { let handler = handler_source(source, handler); - assert!(handler.contains("observe_published_merge_ref(")); + assert!( + handler.contains("observe_published_source_ref(") + || handler.contains("observe_published_merge_ref(") + ); assert!(!handler.contains("repository_reader()")); - assert!(!handler.contains("status: \"unknown\"")); } let submit = handler_source(source, "scoped_submit_merge_request_review"); assert!( @@ -16616,13 +16760,37 @@ mod tests { .find("authorize_review_submission(") .is_some_and(|authorization| { submit - .find("observe_published_merge_ref(") + .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 { @@ -16660,8 +16828,7 @@ mod tests { "work/T-549", "abc123" ), - Err(merge_request::MergeRequestError::Conflict(message)) - if message.starts_with("merge_request_source_workdir_dirty:") + Err(diagnostic) if diagnostic.code == "source_workdir_dirty" )); workdir.cleanliness = Some("clean".to_string()); workdir.current_ref = Some("different".to_string()); @@ -16672,8 +16839,7 @@ mod tests { "work/T-549", "abc123" ), - Err(merge_request::MergeRequestError::Conflict(message)) - if message.starts_with("merge_request_source_workdir_head_mismatch:") + Err(diagnostic) if diagnostic.code == "source_ref_revision_mismatch" )); } From 7c056b1db8db8fff22f99c2b176abdc6db40b60b Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 14:34:09 +0900 Subject: [PATCH 5/6] fix: publish merge request ref diagnostic DTO --- crates/workspace-server/src/records.rs | 1 + web/workspace/src/lib/generated/ticket-api.ts | 4 +++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/crates/workspace-server/src/records.rs b/crates/workspace-server/src/records.rs index 236d138d..2cde25d8 100644 --- a/crates/workspace-server/src/records.rs +++ b/crates/workspace-server/src/records.rs @@ -479,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/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, }; From f1bcd41ad9573ed10111a11b888237da384fa1d9 Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 3 Sep 2026 15:05:30 +0900 Subject: [PATCH 6/6] fix: route embedded repository ref observation --- crates/worker-runtime/src/runtime.rs | 7 ++ crates/workspace-server/src/hosts.rs | 101 +++++++++++++++++++++++++++ 2 files changed, 108 insertions(+) diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index d085d2ca..f9373607 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -404,6 +404,13 @@ impl Runtime { { 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()?; diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index f4419e63..ecc93702 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -2179,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()) } @@ -4838,6 +4860,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: