fix: require provider-published merge request refs

This commit is contained in:
2026-09-03 12:56:17 +09:00
parent 56798f9fb4
commit 7be428d8bf
3 changed files with 306 additions and 46 deletions
+52 -4
View File
@@ -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<RepositoryRefObservation, Error> {
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<WorkingDirectoryStatus> {
RuntimeList::new(Vec::new(), Vec::new())
}
@@ -1449,6 +1460,31 @@ impl RuntimeRegistry {
})
}
pub fn observe_repository_ref(
&self,
runtime_id: &str,
request: RepositoryRefObservationRequest,
) -> Result<RepositoryRefObservation, RuntimeRegistryError> {
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<RepositoryRefObservation, Error> {
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<WorkingDirectoryStatus> {
match self.get_json::<RuntimeHttpWorkingDirectoriesResponse>("/v1/working-directories") {
Ok(response) => RuntimeList::new(response.working_directories, Vec::new()),
+1 -1
View File
@@ -392,7 +392,7 @@ impl RepositoryRegistryReader {
}
}
fn normalize_target_branch_selector(
pub(crate) fn normalize_target_branch_selector(
id: &str,
selector: &str,
) -> Result<String, RepositoryLookupError> {
+253 -41
View File
@@ -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<RepositoryRefObservation> {
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(),