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 {