Merge remote-tracking branch 'refs/remotes/origin/develop' into work/T-566-memory-rest-dto

This commit is contained in:
2026-09-03 16:07:53 +09:00
13 changed files with 1391 additions and 121 deletions
+153 -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,
@@ -2143,6 +2179,28 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime {
}
}
fn observe_repository_ref(
&self,
request: RepositoryRefObservationRequest,
) -> std::result::Result<RepositoryRefObservation, Error> {
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<WorkingDirectoryStatus> {
RuntimeList::new(Vec::new(), Vec::new())
}
@@ -3355,6 +3413,18 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime {
.map_err(|diagnostic| Error::RegistryInconsistency(diagnostic.message))
}
fn observe_repository_ref(
&self,
request: RepositoryRefObservationRequest,
) -> std::result::Result<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()),
@@ -4789,6 +4859,85 @@ mod tests {
}
}
struct ObservingExecutionBackend {
response: RepositoryRefObservation,
observed: Arc<Mutex<Vec<RepositoryRefObservationRequest>>>,
}
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:
+10
View File
@@ -314,12 +314,21 @@ pub struct TicketMergeRequestSummary {
pub review_excerpt: Option<String>,
}
#[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<String>,
pub thread_event_count: usize,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub ref_diagnostics: Vec<MergeRequestRefDiagnostic>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
@@ -470,6 +479,7 @@ pub fn ticket_api_typescript() -> String {
TicketAssignmentPrincipalSummary::decl(&config),
TicketActionEligibility::decl(&config),
TicketMergeRequestSummary::decl(&config),
MergeRequestRefDiagnostic::decl(&config),
MergeRequestListItem::decl(&config),
MergeRequestListResponse::decl(&config),
TicketEvidenceSummary::decl(&config),
+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> {
+577 -80
View File
@@ -121,9 +121,9 @@ use crate::observation::{
RuntimeObservationSource, RuntimeObservationSourceConfig,
};
use crate::records::{
MergeRequestListItem, MergeRequestListResponse, ObjectiveDetail, ObjectiveQueryRequest,
ObjectiveQueryResponse, ObjectiveShowRequest, ProjectRecordList, TicketDetail,
TicketQueryRequest, TicketQueryResponse, TicketShowRequest,
MergeRequestListItem, MergeRequestListResponse, MergeRequestRefDiagnostic, ObjectiveDetail,
ObjectiveQueryRequest, ObjectiveQueryResponse, ObjectiveShowRequest, ProjectRecordList,
TicketDetail, TicketQueryRequest, TicketQueryResponse, TicketShowRequest,
};
use crate::repositories::{
ConfiguredRepository, RepositoryListProjection, RepositoryLogRead, RepositoryLookupError,
@@ -153,10 +153,10 @@ use crate::workdir_removal::{
use crate::workspace_catalog::{WorkspaceCatalogService, WorkspaceCreateRequest};
use crate::{Error, Result};
use worker_runtime::catalog::{
ConfigBundleRef, ProfileSelector, RepositoryMaterializationContext,
RepositorySelector as RuntimeRepositorySelector, RepositorySshMaterializationAccess,
SensitiveString, WorkingDirectoryClaim, WorkingDirectoryRepository, WorkingDirectoryRequest,
WorkspaceApiRef,
ConfigBundleRef, ProfileSelector, RepositoryMaterializationContext, RepositoryRefObservation,
RepositoryRefObservationRequest, RepositorySelector as RuntimeRepositorySelector,
RepositorySshMaterializationAccess, SensitiveString, WorkingDirectoryClaim,
WorkingDirectoryRepository, WorkingDirectoryRequest, WorkspaceApiRef,
};
use worker_runtime::config_bundle::ConfigBundle;
use worker_runtime::http_server::{
@@ -5417,6 +5417,205 @@ fn merge_request_store(
.map_err(Into::into)
}
fn repository_ref_observation_error(error: crate::hosts::RuntimeRegistryError) -> ApiError {
match error {
crate::hosts::RuntimeRegistryError::RuntimeOperationFailed {
runtime_id,
code,
message,
} => Error::RuntimeOperationFailed {
runtime_id,
code,
message,
}
.into(),
other => other.into_error().into(),
}
}
fn observe_published_merge_ref(
api: &WorkspaceApi,
workspace_id: &str,
runtime_id: &str,
repository_id: &str,
selector: &str,
) -> ApiResult<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 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<RepositoryRefObservation> {
observe_published_merge_ref(api, workspace_id, runtime_id, repository_id, selector)
.map_err(remap_source_ref_error)
}
fn require_assigned_workdir_source(
api: &WorkspaceApi,
assignment: &crate::store::TicketCoderAssignmentRecord,
repository_id: &str,
selector: &str,
revision_ref: &str,
) -> ApiResult<()> {
let worker = api
.runtime
.worker(&assignment.worker)
.map_err(|error| error.into_error())?;
let attached_workdir = worker.working_directory.ok_or_else(|| {
Error::MergeRequest(merge_request::MergeRequestError::Conflict(
"merge_request_source_workdir_missing: current Coder has no attached Workdir".into(),
))
})?;
let workdir = api
.runtime
.working_directory(
&assignment.worker.runtime_id,
&attached_workdir.working_directory_id,
)
.map_err(|error| error.into_error())?
.working_directory
.ok_or_else(|| {
Error::MergeRequest(merge_request::MergeRequestError::Conflict(
"merge_request_source_workdir_unavailable: current Coder Workdir could not be observed"
.into(),
))
})?
.summary;
validate_assigned_workdir_source(&workdir, repository_id, selector, revision_ref).map_err(
|diagnostic| Error::RuntimeOperationFailed {
runtime_id: assignment.worker.runtime_id.clone(),
code: diagnostic.code,
message: diagnostic.message,
},
)?;
Ok(())
}
fn validate_assigned_workdir_source(
workdir: &worker_runtime::catalog::WorkingDirectorySummary,
repository_id: &str,
selector: &str,
revision_ref: &str,
) -> std::result::Result<(), worker_runtime::working_directory::WorkingDirectoryDiagnostic> {
if workdir.repository_id != repository_id {
return Err(
worker_runtime::working_directory::WorkingDirectoryDiagnostic {
code: "source_workdir_repository_mismatch".to_string(),
message: "Current Coder Workdir belongs to a different Repository".to_string(),
},
);
}
if workdir.cleanliness.as_deref() != Some("clean") {
return Err(
worker_runtime::working_directory::WorkingDirectoryDiagnostic {
code: "source_workdir_dirty".to_string(),
message: "Current Coder Workdir must be clean".to_string(),
},
);
}
let workdir_selector_matches = workdir
.current_selector
.as_deref()
.and_then(|workdir_selector| {
crate::repositories::normalize_target_branch_selector(repository_id, workdir_selector)
.ok()
})
.is_some_and(|workdir_selector| {
crate::repositories::normalize_target_branch_selector(repository_id, selector)
.is_ok_and(|selector| selector == workdir_selector)
});
if !workdir_selector_matches || workdir.current_ref.as_deref() != Some(revision_ref) {
return Err(
worker_runtime::working_directory::WorkingDirectoryDiagnostic {
code: "source_ref_revision_mismatch".to_string(),
message: "Current Coder Workdir selector and HEAD do not match the provider-published source ref".to_string(),
},
);
}
Ok(())
}
fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError {
Error::InvalidInput(format!(
"repository merge evidence validation failed: {error:?}"
@@ -5536,6 +5735,8 @@ struct MergeRequestRefResponse {
#[serde(rename = "ref")]
revision_ref: Option<String>,
observed_at: String,
#[serde(skip_serializing_if = "Option::is_none")]
diagnostic: Option<MergeRequestRefDiagnostic>,
}
#[derive(Debug, serde::Serialize)]
@@ -5615,17 +5816,44 @@ async fn scoped_list_merge_requests(
limit: query.limit.unwrap_or(50),
},
)?;
let reader = api.repository_reader();
let items = page
.items
.into_iter()
.map(|merge_request| -> ApiResult<MergeRequestListItem> {
let current_subject_ref = merge_request.selector_from.as_deref().and_then(|selector| {
reader
.observe_merge_target(&merge_request.repository_id, Some(selector))
.ok()
.map(|observation| observation.commit)
});
let (current_subject_ref, ref_diagnostics) =
if merge_request.state == merge_request::MergeRequestState::Open {
match (
merge_request.selector_from.as_deref(),
merge_request.ticket_ids.first(),
) {
(Some(selector), Some(ticket_id)) => match api
.store
.get_current_ticket_coder_assignment(&workspace_id, ticket_id)?
{
Some(assignment) => match observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&merge_request.repository_id,
selector,
) {
Ok(observation) => (Some(observation.revision_ref), Vec::new()),
Err(error) => (None, vec![merge_ref_diagnostic(error)]),
},
None => (
None,
vec![MergeRequestRefDiagnostic {
code: "source_ref_runtime_unavailable".to_string(),
message: "No current Coder Runtime is available to observe the source ref"
.to_string(),
}],
),
},
_ => (None, Vec::new()),
}
} else {
(None, Vec::new())
};
let repository_key = api
.store
.get_repository(&workspace_id, &merge_request.repository_id)?
@@ -5641,6 +5869,7 @@ async fn scoped_list_merge_requests(
summary: merge_request_summary(merge_request, repository_key, current_subject_ref),
ticket_ids,
thread_event_count,
ref_diagnostics,
})
})
.collect::<ApiResult<Vec<_>>>()?;
@@ -5650,6 +5879,55 @@ async fn scoped_list_merge_requests(
}))
}
fn merge_ref_diagnostic(error: ApiError) -> MergeRequestRefDiagnostic {
let ApiError { error, diagnostics } = error;
diagnostics
.into_iter()
.next()
.map(|diagnostic| MergeRequestRefDiagnostic {
code: diagnostic.code,
message: diagnostic.message,
})
.unwrap_or_else(|| MergeRequestRefDiagnostic {
code: "merge_ref_observation_unavailable".to_string(),
message: sanitize_backend_error(&error.to_string()),
})
}
fn unknown_merge_ref(code: &str, message: &str) -> MergeRequestRefResponse {
MergeRequestRefResponse {
status: "unknown".to_string(),
revision_ref: None,
observed_at: Utc::now().to_rfc3339(),
diagnostic: Some(MergeRequestRefDiagnostic {
code: code.to_string(),
message: message.to_string(),
}),
}
}
fn unknown_merge_ref_response(error: ApiError) -> MergeRequestRefResponse {
MergeRequestRefResponse {
status: "unknown".to_string(),
revision_ref: None,
observed_at: Utc::now().to_rfc3339(),
diagnostic: Some(merge_ref_diagnostic(error)),
}
}
fn merge_ref_response(observation: RepositoryRefObservation) -> MergeRequestRefResponse {
let observed_at =
chrono::DateTime::<Utc>::from_timestamp(observation.observed_at_epoch_seconds as i64, 0)
.unwrap_or_else(Utc::now)
.to_rfc3339();
MergeRequestRefResponse {
status: "known".into(),
revision_ref: Some(observation.revision_ref),
observed_at,
diagnostic: None,
}
}
async fn scoped_show_merge_request(
State(api): State<WorkspaceApi>,
AxumPath((workspace_id, merge_request_id)): AxumPath<(String, String)>,
@@ -5665,38 +5943,53 @@ async fn scoped_show_merge_request(
query.after,
query.limit.unwrap_or(100),
)?;
let reader = api.repository_reader();
let observed_at = Utc::now().to_rfc3339();
let source = match mr.selector_from.as_deref() {
Some(selector) => match reader.observe_merge_target(&mr.repository_id, Some(selector)) {
Ok(value) => MergeRequestRefResponse {
status: "known".into(),
revision_ref: Some(value.commit),
observed_at: observed_at.clone(),
},
Err(_) => MergeRequestRefResponse {
status: "unknown".into(),
revision_ref: None,
observed_at: observed_at.clone(),
},
},
None => MergeRequestRefResponse {
let ticket_id = mr.ticket_ids.first().ok_or_else(|| {
Error::MergeRequest(merge_request::MergeRequestError::Conflict(
"Merge Request has no linked Ticket".into(),
))
})?;
let assignment = api
.store
.get_current_ticket_coder_assignment(&workspace_id, ticket_id)?;
let source = match (mr.selector_from.as_deref(), assignment.as_ref()) {
(Some(selector), Some(assignment)) => {
match observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
selector,
) {
Ok(observation) => merge_ref_response(observation),
Err(error) => unknown_merge_ref_response(error),
}
}
(Some(_), None) => unknown_merge_ref(
"source_ref_runtime_unavailable",
"No current Coder Runtime is available to observe the source ref",
),
(None, _) => MergeRequestRefResponse {
status: "requires_repair".into(),
revision_ref: None,
observed_at: observed_at.clone(),
observed_at: Utc::now().to_rfc3339(),
diagnostic: None,
},
};
let target = match reader.observe_merge_target(&mr.repository_id, Some(&mr.selector_to)) {
Ok(value) => MergeRequestRefResponse {
status: "known".into(),
revision_ref: Some(value.commit),
observed_at,
},
Err(_) => MergeRequestRefResponse {
status: "unknown".into(),
revision_ref: None,
observed_at,
let target = match assignment.as_ref() {
Some(assignment) => match observe_published_merge_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
&mr.selector_to,
) {
Ok(observation) => merge_ref_response(observation),
Err(error) => unknown_merge_ref_response(error),
},
None => unknown_merge_ref(
"target_ref_runtime_unavailable",
"No current Coder Runtime is available to observe the target ref",
),
};
let linked_tickets = mr
.ticket_ids
@@ -5729,13 +6022,28 @@ async fn scoped_merge_request_readiness(
let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?;
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let current_subject_ref = mr.selector_from.as_deref().and_then(|selector| {
api.repository_reader()
.observe_merge_target(&mr.repository_id, Some(selector))
.ok()
.map(|v| v.commit)
});
Ok(Json(store.readiness(merge_request::ReadinessCheck {
let assignment = api
.store
.get_current_ticket_coder_assignment(&workspace_id, &ticket_id)?;
let (current_subject_ref, source_blocker) = match (mr.selector_from.as_deref(), assignment) {
(Some(selector), Some(assignment)) => match observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
selector,
) {
Ok(observation) => (Some(observation.revision_ref), None),
Err(error) => {
let diagnostic = merge_ref_diagnostic(error);
let blocker = source_ref_readiness_blocker(&diagnostic.code);
(None, Some(blocker.to_string()))
}
},
(Some(_), None) => (None, Some("source_ref_unavailable".to_string())),
(None, _) => (None, None),
};
let mut report = store.readiness(merge_request::ReadinessCheck {
ticket_id,
current_subject_ref,
auth: merge_request::MergeRequestAuth {
@@ -5745,7 +6053,17 @@ async fn scoped_merge_request_readiness(
worker_id: String::new(),
assignment_id: String::new(),
},
})?))
})?;
if let Some(blocker) = source_blocker {
report
.blockers
.retain(|current| current != "selector_unresolved");
if !report.blockers.contains(&blocker) {
report.blockers.push(blocker);
}
report.ready = false;
}
Ok(Json(report))
}
async fn scoped_open_merge_request(
@@ -5785,13 +6103,27 @@ async fn scoped_open_merge_request(
)
.into());
}
let reader = api.repository_reader();
reader
.observe_merge_target(&repository_id, Some(&input.selector_from))
.map_err(repository_merge_evidence_error)?;
reader
.observe_merge_target(&repository_id, Some(&input.selector_to))
.map_err(repository_merge_evidence_error)?;
let source_observation = observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&repository_id,
&input.selector_from,
)?;
require_assigned_workdir_source(
&api,
&assignment,
&repository_id,
&input.selector_from,
&source_observation.revision_ref,
)?;
observe_published_merge_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&repository_id,
&input.selector_to,
)?;
let merge_request = merge_request_store(&api, &workspace_id)?.open_merge_request(
merge_request::OpenMergeRequest {
merge_request_id: Uuid::now_v7().to_string(),
@@ -5850,11 +6182,20 @@ async fn scoped_repair_merge_request_selector(
}
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let resolved_subject_ref = api
.repository_reader()
.observe_merge_target(&mr.repository_id, Some(&input.selector_from))
.map_err(repository_merge_evidence_error)?
.commit;
let assignment = api
.store
.get_current_ticket_coder_assignment(&workspace_id, &ticket_id)?
.ok_or_else(|| {
Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into())
})?;
let resolved_subject_ref = observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
&input.selector_from,
)?
.revision_ref;
let repaired = store.repair_selector_from(merge_request::RepairSelectorFrom {
workspace_id: workspace_id.clone(),
ticket_id,
@@ -5922,11 +6263,21 @@ async fn scoped_register_merge_request_review_capability(
.selector_from
.as_deref()
.ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?;
let subject_ref = api
.repository_reader()
.observe_merge_target(&mr.repository_id, Some(selector))
.map_err(repository_merge_evidence_error)?
.commit;
let source_observation = observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
selector,
)?;
require_assigned_workdir_source(
&api,
&assignment,
&mr.repository_id,
selector,
&source_observation.revision_ref,
)?;
let subject_ref = source_observation.revision_ref;
store.request_review(merge_request::RequestMergeRequestReview {
ticket_id,
subject_ref,
@@ -5952,16 +6303,35 @@ async fn scoped_submit_merge_request_review(
let workspace_id = parse_workspace_id(&workspace_id)?;
let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?;
let store = merge_request_store(&api, &workspace_id)?;
let review_authorization =
store.authorize_review_submission(&ticket_id, &input.capability_token)?;
if review_authorization.workspace_id != workspace_id {
return Err(
Error::MergeRequest(merge_request::MergeRequestError::Unauthorized(
"review grant invalid".into(),
))
.into(),
);
}
let mr = store.get(&workspace_id, &ticket_id)?;
let selector = mr
.selector_from
.as_deref()
.ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?;
let current_subject_ref = api
.repository_reader()
.observe_merge_target(&mr.repository_id, Some(selector))
.map_err(repository_merge_evidence_error)?
.commit;
let assignment = api
.store
.get_current_ticket_coder_assignment(&workspace_id, &ticket_id)?
.ok_or_else(|| {
Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into())
})?;
let current_subject_ref = observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
selector,
)?
.revision_ref;
Ok(Json(store.submit_review(
merge_request::SubmitMergeRequestReview {
ticket_id,
@@ -6034,7 +6404,6 @@ async fn scoped_complete_merge_request(
require_online_workspace_orchestrator_source(&api, &source)?;
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let repositories = api.repository_reader();
if let Some(existing) = recorded_merge_completion(&mr.thread, &input.operation_id) {
let replay = merge_request::CompleteMergeRequest {
ticket_id,
@@ -6066,15 +6435,23 @@ async fn scoped_complete_merge_request(
.selector_from
.as_deref()
.ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?;
let current_source_ref = repositories
.observe_merge_target(&mr.repository_id, Some(selector))
.map_err(repository_merge_evidence_error)?
.commit;
let observed = repositories
.observe_merge_target(&mr.repository_id, Some(&mr.selector_to))
.map_err(repository_merge_evidence_error)?;
let current_source_ref = observe_published_source_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
selector,
)?
.revision_ref;
let observed_target = observe_published_merge_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
&mr.selector_to,
)?;
require_completed_target_observation(
&observed.commit,
&observed_target.revision_ref,
&input.target_ref_before,
&input.target_ref_after,
)?;
@@ -16300,6 +16677,126 @@ 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_http_paths_observe_refs_through_runtime_provider_authority() {
let source = include_str!("server.rs");
for handler in [
"scoped_list_merge_requests",
"scoped_show_merge_request",
"scoped_open_merge_request",
"scoped_repair_merge_request_selector",
"scoped_register_merge_request_review_capability",
"scoped_submit_merge_request_review",
"scoped_complete_merge_request",
] {
let handler = handler_source(source, handler);
assert!(
handler.contains("observe_published_source_ref(")
|| handler.contains("observe_published_merge_ref(")
);
assert!(!handler.contains("repository_reader()"));
}
let submit = handler_source(source, "scoped_submit_merge_request_review");
assert!(
submit
.find("authorize_review_submission(")
.is_some_and(|authorization| {
submit
.find("observe_published_source_ref(")
.is_some_and(|observation| authorization < observation)
}),
"review capability must be validated before provider access",
);
}
#[test]
fn readiness_distinguishes_missing_source_from_unavailable_provider() {
assert_eq!(
source_ref_readiness_blocker("source_ref_not_found"),
"source_ref_not_found"
);
assert_eq!(
source_ref_error_code("repository_ref_not_found"),
"source_ref_not_found"
);
assert_eq!(
source_ref_error_code("repository_ref_provider_timeout"),
"source_ref_provider_timeout"
);
for code in [
"source_ref_provider_unavailable",
"source_ref_provider_timeout",
"source_ref_provider_auth_failed",
"source_ref_credential_expired",
] {
assert_eq!(source_ref_readiness_blocker(code), "source_ref_unavailable");
}
}
#[test]
fn assigned_workdir_must_be_clean_and_match_published_source() {
let mut workdir = worker_runtime::catalog::WorkingDirectorySummary {
working_directory_id: "workdir-1".to_string(),
repository_id: "repository-1".to_string(),
creation_selector: Some("work/T-549".to_string()),
creation_ref: Some("abc123".to_string()),
creation_tree: None,
current_selector: Some("work/T-549".to_string()),
current_ref: Some("abc123".to_string()),
current_tree: None,
observed_at_epoch_seconds: Some(1_767_225_600),
materializer_kind: workspace_api::WorkingDirectoryMaterializerKind::RuntimeGitCache,
cleanup_target: None,
status: worker_runtime::catalog::WorkingDirectoryStatusKind::Active,
cleanliness: Some("clean".to_string()),
primary_worker_id: Some("worker-1".to_string()),
occupied_by: None,
};
assert!(
validate_assigned_workdir_source(
&workdir,
"repository-1",
"refs/heads/work/T-549",
"abc123",
)
.is_ok()
);
workdir.cleanliness = Some("dirty".to_string());
assert!(matches!(
validate_assigned_workdir_source(
&workdir,
"repository-1",
"work/T-549",
"abc123"
),
Err(diagnostic) if diagnostic.code == "source_workdir_dirty"
));
workdir.cleanliness = Some("clean".to_string());
workdir.current_ref = Some("different".to_string());
assert!(matches!(
validate_assigned_workdir_source(
&workdir,
"repository-1",
"work/T-549",
"abc123"
),
Err(diagnostic) if diagnostic.code == "source_ref_revision_mismatch"
));
}
fn completed_upload_file(sha256: &str) -> protocol::UploadedFileRef {
protocol::UploadedFileRef {
artifact_id: "019ca7c8-57b6-7f05-8edf-524147aba7b3".into(),