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(),