fix: harden provider ref review boundaries

This commit is contained in:
2026-09-03 13:49:20 +09:00
parent 7be428d8bf
commit 2884c08466
6 changed files with 364 additions and 81 deletions
+34
View File
@@ -274,6 +274,12 @@ pub struct RegisterReviewerChildSession {
pub reviewer_profile: String, pub reviewer_profile: String,
pub now: DateTime<Utc>, pub now: DateTime<Utc>,
} }
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReviewSubmissionAuthorization {
pub workspace_id: String,
pub subject_ref: String,
}
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct SubmitMergeRequestReview { pub struct SubmitMergeRequestReview {
pub ticket_id: String, pub ticket_id: String,
@@ -535,6 +541,34 @@ impl MergeRequestStore {
t.commit()?; t.commit()?;
Ok(RequestedMergeRequestReview { request_event: e }) Ok(RequestedMergeRequestReview { request_event: e })
} }
pub fn authorize_review_submission(
&self,
ticket_id: &str,
capability_token: &str,
) -> Result<ReviewSubmissionAuthorization, MergeRequestError> {
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( pub fn submit_review(
&self, &self,
i: SubmitMergeRequestReview, i: SubmitMergeRequestReview,
+17
View File
@@ -91,6 +91,23 @@ fn approve(s: &MergeRequestStore, subject: &str, token: &str) -> ReviewEvent {
}) })
.unwrap() .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] #[test]
fn selectors_thread_and_completion_have_no_revision_or_commit_api() { fn selectors_thread_and_completion_have_no_revision_or_commit_api() {
let (d, s) = fixture(); let (d, s) = fixture();
+47 -12
View File
@@ -1970,6 +1970,33 @@ fn status_for_runtime_error(error: &RuntimeError) -> StatusCode {
{ {
StatusCode::NOT_FOUND 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::RuntimeStopped
| RuntimeError::WorkerExecutionUnavailable { .. } | RuntimeError::WorkerExecutionUnavailable { .. }
| RuntimeError::ExecutionBackendUnavailable { .. } | RuntimeError::ExecutionBackendUnavailable { .. }
@@ -1980,8 +2007,8 @@ fn status_for_runtime_error(error: &RuntimeError) -> StatusCode {
| RuntimeError::InvalidInitialInputKind { .. } | RuntimeError::InvalidInitialInputKind { .. }
| RuntimeError::ConfigBundleDigestMismatch { .. } | RuntimeError::ConfigBundleDigestMismatch { .. }
| RuntimeError::InvalidProfileSelector { .. } | RuntimeError::InvalidProfileSelector { .. }
| RuntimeError::UnsupportedConfigDeclaration { .. } | RuntimeError::UnsupportedConfigDeclaration { .. } => StatusCode::BAD_REQUEST,
| RuntimeError::WorkingDirectory(_) => StatusCode::BAD_REQUEST, RuntimeError::WorkingDirectory(_) => StatusCode::BAD_REQUEST,
RuntimeError::StoreIo { .. } RuntimeError::StoreIo { .. }
| RuntimeError::StoreMissing { .. } | RuntimeError::StoreMissing { .. }
| RuntimeError::StoreCorrupt { .. } | RuntimeError::StoreCorrupt { .. }
@@ -2967,17 +2994,25 @@ mod tests {
#[test] #[test]
fn workdir_runtime_errors_preserve_diagnostic_code() { fn workdir_runtime_errors_preserve_diagnostic_code() {
let error = let cases = [
RuntimeError::WorkingDirectory(crate::working_directory::WorkingDirectoryDiagnostic { ("working_directory_not_found", StatusCode::NOT_FOUND),
code: "working_directory_not_found".to_string(), (
message: "working directory missing-workdir was not found".to_string(), "repository_ref_provider_timeout",
}); StatusCode::SERVICE_UNAVAILABLE,
),
assert_eq!(status_for_runtime_error(&error), StatusCode::NOT_FOUND); ("repository_ref_provider_auth_failed", StatusCode::FORBIDDEN),
assert_eq!( ("repository_ref_not_found", StatusCode::NOT_FOUND),
code_for_runtime_error(&error), ];
"working_directory_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);
}
} }
} }
+55 -13
View File
@@ -3002,20 +3002,35 @@ fn worker_status_from_run_state(run_state: WorkerExecutionRunState) -> WorkerSta
} }
fn repository_resource_error(error: BackendResourceError) -> RuntimeError { fn repository_resource_error(error: BackendResourceError) -> RuntimeError {
let category = match error { let (code, message) = match error {
BackendResourceError::Expired => "expired", BackendResourceError::Expired => (
BackendResourceError::Unauthorized { .. } => "unauthorized", "repository_access_credential_expired",
BackendResourceError::UnsupportedKind => "unsupported_kind", "Repository access credential lease expired",
BackendResourceError::MissingResource => "missing_resource", ),
BackendResourceError::Oversized { .. } => "oversized", BackendResourceError::Unauthorized { .. } => (
BackendResourceError::DigestMismatch { .. } => "digest_mismatch", "repository_access_credential_unauthorized",
BackendResourceError::ContentTypeMismatch { .. } => "content_type_mismatch", "Repository access credential lease was rejected",
BackendResourceError::InvalidResponse { .. } => "invalid_response", ),
BackendResourceError::Transport { .. } => "transport", 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!( RuntimeError::WorkingDirectory(
"Backend Repository SSH access resource fetch failed: {category}" crate::working_directory::WorkingDirectoryDiagnostic::rejected(code, message),
)) )
} }
fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest { fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest {
@@ -3217,6 +3232,33 @@ mod tests {
use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex}; 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( fn internal_worker_ref(
session_id: &str, session_id: &str,
parent_session_id: Option<&str>, parent_session_id: Option<&str>,
+44 -8
View File
@@ -954,13 +954,7 @@ impl WorkingDirectoryMaterializer for RuntimeGitCacheMaterializer {
request: &RepositoryRefObservationRequest, request: &RepositoryRefObservationRequest,
) -> Result<RepositoryRefObservation, WorkingDirectoryDiagnostic> { ) -> Result<RepositoryRefObservation, WorkingDirectoryDiagnostic> {
let selector = request.selector.trim(); let selector = request.selector.trim();
validate_selector(selector)?; validate_exact_branch_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 { let working_request = WorkingDirectoryRequest {
repository: request.repository.clone(), repository: request.repository.clone(),
@@ -2100,6 +2094,40 @@ fn repository_git_command(
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<u8> { fn read_bounded_command_output(mut reader: impl Read) -> Vec<u8> {
const MAX_CAPTURE_BYTES: usize = 8192; const MAX_CAPTURE_BYTES: usize = 8192;
let mut captured = Vec::new(); let mut captured = Vec::new();
@@ -2729,12 +2757,20 @@ mod tests {
assert_eq!(missing.code, "repository_ref_not_found"); assert_eq!(missing.code, "repository_ref_not_found");
let non_branch = materializer let non_branch = materializer
.observe_repository_ref(&RepositoryRefObservationRequest { .observe_repository_ref(&RepositoryRefObservationRequest {
repository, repository: repository.clone(),
selector: "HEAD".to_string(), selector: "HEAD".to_string(),
materialization: None, materialization: None,
}) })
.unwrap_err(); .unwrap_err();
assert_eq!(non_branch.code, "repository_ref_selector_invalid"); 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] #[test]
+165 -46
View File
@@ -5593,19 +5593,26 @@ fn require_assigned_workdir_source(
)) ))
})? })?
.summary; .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 { 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(), "merge_request_source_workdir_repository_mismatch: current Coder Workdir belongs to a different Repository".into(),
)) ));
.into());
} }
if workdir.cleanliness.as_deref() != Some("clean") { if workdir.cleanliness.as_deref() != Some("clean") {
return Err( return Err(merge_request::MergeRequestError::Conflict(
Error::MergeRequest(merge_request::MergeRequestError::Conflict(
"merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(), "merge_request_source_workdir_dirty: current Coder Workdir must be clean".into(),
)) ));
.into(),
);
} }
let workdir_selector_matches = workdir let workdir_selector_matches = workdir
.current_selector .current_selector
@@ -5619,10 +5626,9 @@ fn require_assigned_workdir_source(
.is_ok_and(|selector| selector == workdir_selector) .is_ok_and(|selector| selector == workdir_selector)
}); });
if !workdir_selector_matches || workdir.current_ref.as_deref() != Some(revision_ref) { 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(), "merge_request_source_workdir_head_mismatch: current Coder Workdir selector and HEAD must match the provider-published source ref".into(),
)) ));
.into());
} }
Ok(()) Ok(())
} }
@@ -5825,17 +5831,41 @@ async fn scoped_list_merge_requests(
limit: query.limit.unwrap_or(50), limit: query.limit.unwrap_or(50),
}, },
)?; )?;
let reader = api.repository_reader();
let items = page let items = page
.items .items
.into_iter() .into_iter()
.map(|merge_request| -> ApiResult<MergeRequestListItem> { .map(|merge_request| -> ApiResult<MergeRequestListItem> {
let current_subject_ref = merge_request.selector_from.as_deref().and_then(|selector| { let current_subject_ref =
reader if merge_request.state == merge_request::MergeRequestState::Open {
.observe_merge_target(&merge_request.repository_id, Some(selector)) match (
.ok() merge_request.selector_from.as_deref(),
.map(|observation| observation.commit) 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 let repository_key = api
.store .store
.get_repository(&workspace_id, &merge_request.repository_id)? .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::<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,
}
}
async fn scoped_show_merge_request( async fn scoped_show_merge_request(
State(api): State<WorkspaceApi>, State(api): State<WorkspaceApi>,
AxumPath((workspace_id, merge_request_id)): AxumPath<(String, String)>, AxumPath((workspace_id, merge_request_id)): AxumPath<(String, String)>,
@@ -5875,39 +5917,38 @@ async fn scoped_show_merge_request(
query.after, query.after,
query.limit.unwrap_or(100), query.limit.unwrap_or(100),
)?; )?;
let reader = api.repository_reader(); let ticket_id = mr.ticket_ids.first().ok_or_else(|| {
let observed_at = Utc::now().to_rfc3339(); 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() { let source = match mr.selector_from.as_deref() {
Some(selector) => match reader.observe_merge_target(&mr.repository_id, Some(selector)) { Some(selector) => merge_ref_response(observe_published_merge_ref(
Ok(value) => MergeRequestRefResponse { &api,
status: "known".into(), &workspace_id,
revision_ref: Some(value.commit), &assignment.worker.runtime_id,
observed_at: observed_at.clone(), &mr.repository_id,
}, selector,
Err(_) => MergeRequestRefResponse { )?),
status: "unknown".into(),
revision_ref: None,
observed_at: observed_at.clone(),
},
},
None => MergeRequestRefResponse { None => MergeRequestRefResponse {
status: "requires_repair".into(), status: "requires_repair".into(),
revision_ref: None, revision_ref: None,
observed_at: observed_at.clone(), observed_at: Utc::now().to_rfc3339(),
},
};
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 = merge_ref_response(observe_published_merge_ref(
&api,
&workspace_id,
&assignment.worker.runtime_id,
&mr.repository_id,
&mr.selector_to,
)?);
let linked_tickets = mr let linked_tickets = mr
.ticket_ids .ticket_ids
.iter() .iter()
@@ -6209,6 +6250,16 @@ async fn scoped_submit_merge_request_review(
let workspace_id = parse_workspace_id(&workspace_id)?; let workspace_id = parse_workspace_id(&workspace_id)?;
let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?; let ticket_id = resolve_workspace_ticket_reference(&api, &workspace_id, &ticket_id)?;
let store = merge_request_store(&api, &workspace_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 mr = store.get(&workspace_id, &ticket_id)?;
let selector = mr let selector = mr
.selector_from .selector_from
@@ -16543,9 +16594,11 @@ mod tests {
} }
#[test] #[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"); let source = include_str!("server.rs");
for handler in [ for handler in [
"scoped_list_merge_requests",
"scoped_show_merge_request",
"scoped_open_merge_request", "scoped_open_merge_request",
"scoped_repair_merge_request_selector", "scoped_repair_merge_request_selector",
"scoped_register_merge_request_review_capability", "scoped_register_merge_request_review_capability",
@@ -16555,7 +16608,73 @@ mod tests {
let handler = handler_source(source, handler); let handler = handler_source(source, handler);
assert!(handler.contains("observe_published_merge_ref(")); assert!(handler.contains("observe_published_merge_ref("));
assert!(!handler.contains("repository_reader()")); 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 { fn completed_upload_file(sha256: &str) -> protocol::UploadedFileRef {