fix: align merge requests with selector thread contract

This commit is contained in:
2026-08-17 07:41:03 +09:00
parent 9e48cae759
commit 1fe1b7463f
13 changed files with 1354 additions and 2658 deletions
+240 -542
View File
@@ -94,8 +94,8 @@ use crate::observation::{
use crate::profile_settings::UpdateWorkspaceMetadataRequest;
use crate::records::{ObjectiveDetail, ProjectRecordList, TicketDetail};
use crate::repositories::{
ConfiguredRepository, MergeTargetObservation, RepositoryListProjection, RepositoryLogRead,
RepositoryLookupError, RepositoryRegistryReader, RepositorySummary,
ConfiguredRepository, RepositoryListProjection, RepositoryLogRead, RepositoryLookupError,
RepositoryRegistryReader, RepositorySummary,
};
use crate::resource_broker::BackendResourceBroker;
use crate::runtime_subscription::RuntimeSubscriptionBroker;
@@ -1324,8 +1324,12 @@ pub fn build_router(api: WorkspaceApi) -> Router {
get(scoped_merge_request_readiness),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/review-requests",
post(scoped_request_merge_request_review),
"/api/w/{workspace_id}/tickets/{id}/merge-request/thread",
get(scoped_merge_request_thread),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/repair-source",
post(scoped_repair_merge_request_selector),
)
.route(
"/api/w/{workspace_id}/internal/reviewer-child-sessions",
@@ -1339,18 +1343,15 @@ pub fn build_router(api: WorkspaceApi) -> Router {
"/api/w/{workspace_id}/tickets/{id}/merge-request/reviews",
post(scoped_submit_merge_request_review),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/reviews/revoke",
post(scoped_revoke_merge_request_review),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/complete",
post(scoped_complete_merge_request),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/close",
post(scoped_close_merge_request),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/merge-request/reopen",
post(scoped_reopen_merge_request),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/workflow/close",
post(scoped_close_ticket_record),
@@ -3583,23 +3584,28 @@ struct OpenMergeRequestRequest {
repository_id: String,
selector_from: String,
selector_to: String,
base_commit: String,
head_commit: String,
#[serde(default)]
changed_paths: Vec<String>,
#[serde(default)]
summary: String,
}
#[derive(Debug, serde::Deserialize)]
struct RequestMergeRequestReviewRequest {
expected_head_commit: String,
base_commit: String,
head_commit: String,
#[serde(default)]
changed_paths: Vec<String>,
#[serde(default)]
summary: String,
struct RepairMergeRequestSelectorRequest {
selector_from: String,
reason: String,
explicit_confirmation: bool,
}
#[derive(Debug, serde::Deserialize)]
struct RevokeMergeRequestReviewRequest {
review_event_id: String,
reason: String,
explicit_confirmation: bool,
}
#[derive(Debug, serde::Deserialize)]
struct MergeRequestThreadQuery {
after: Option<u64>,
limit: Option<usize>,
}
#[derive(Debug, serde::Deserialize)]
@@ -3609,14 +3615,12 @@ struct RegisterReviewerChildSessionRequest {
#[derive(Debug, serde::Deserialize)]
struct RegisterMergeRequestReviewCapabilityRequest {
expected_head_commit: String,
child_session_id: String,
capability_token: String,
}
#[derive(Debug, serde::Deserialize)]
struct SubmitMergeRequestReviewRequest {
expected_head_commit: String,
capability_token: String,
decision: merge_request::ReviewDecision,
#[serde(default)]
@@ -3628,21 +3632,13 @@ struct SubmitMergeRequestReviewRequest {
#[derive(Debug, serde::Deserialize)]
struct CompleteMergeRequestRequest {
operation_id: String,
expected_head_commit: String,
target_commit: String,
source_commit: String,
result_commit: String,
approval_event_id: String,
target_ref_before: String,
target_ref_after: String,
strategy: merge_request::MergeStrategy,
resolution: merge_request::ConflictResolution,
}
#[derive(Debug, serde::Deserialize)]
struct MergeRequestStateRequest {
#[serde(default)]
body: String,
explicit_confirmation: bool,
}
fn parse_workspace_id(value: &str) -> ApiResult<String> {
if value.trim().is_empty() {
return Err(Error::InvalidInput("workspace_id must not be empty".to_string()).into());
@@ -3699,26 +3695,6 @@ impl merge_request::RepositorySource for MergeRequestRepositorySource {
}
Ok(self.reader.summary(repository_id).is_ok())
}
fn is_ancestor(
&self,
workspace_id: &str,
repository_id: &str,
ancestor: &str,
descendant: &str,
) -> std::result::Result<bool, String> {
if workspace_id != self.workspace_id {
return Ok(false);
}
match self
.reader
.ensure_ancestor(repository_id, ancestor, descendant)
{
Ok(()) => Ok(true),
Err(RepositoryLookupError::InvalidCommitRelation { .. }) => Ok(false),
Err(error) => Err(format!("{error:?}")),
}
}
}
fn merge_request_store(
@@ -3747,65 +3723,36 @@ fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError {
.into()
}
fn validate_open_merge_request_evidence(
api: &WorkspaceApi,
ticket_id: &str,
repository_id: &str,
base_commit: &str,
head_commit: &str,
) -> ApiResult<(String, String, MergeTargetObservation)> {
let ticket = browser_ticket_backend(api)?
.show(TicketIdOrSlug::Id(ticket_id.into()))
.map_err(Error::from)?;
if ticket.meta.repository_id.as_deref() != Some(repository_id) {
return Err(Error::InvalidInput(
"Merge Request repository must match the authoritative Ticket target".into(),
)
.into());
}
let reader = api.repository_reader();
let target = reader
.observe_merge_target(repository_id, ticket.meta.ref_selector.as_deref())
.map_err(repository_merge_evidence_error)?;
let base = reader
.observe_commit(repository_id, base_commit)
.map_err(repository_merge_evidence_error)?;
let source = reader
.observe_commit(repository_id, head_commit)
.map_err(repository_merge_evidence_error)?;
reader
.ensure_ancestor(repository_id, &base.commit, &source.commit)
.map_err(repository_merge_evidence_error)?;
Ok((base.commit, source.commit, target))
}
fn validate_revision_evidence(
api: &WorkspaceApi,
repository_id: &str,
base_commit: &str,
head_commit: &str,
) -> ApiResult<(String, String)> {
let reader = api.repository_reader();
let base = reader
.observe_commit(repository_id, base_commit)
.map_err(repository_merge_evidence_error)?;
let source = reader
.observe_commit(repository_id, head_commit)
.map_err(repository_merge_evidence_error)?;
reader
.ensure_ancestor(repository_id, &base.commit, &source.commit)
.map_err(repository_merge_evidence_error)?;
Ok((base.commit, source.commit))
}
async fn scoped_show_merge_request(
State(api): State<WorkspaceApi>,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
) -> ApiResult<Json<merge_request::MergeRequest>> {
) -> ApiResult<Json<serde_json::Value>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
Ok(Json(
merge_request_store(&api, &workspace_id)?.get(&workspace_id, &ticket_id)?,
))
let mr = merge_request_store(&api, &workspace_id)?.get(&workspace_id, &ticket_id)?;
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) => {
serde_json::json!({"status":"known","ref":value.commit,"observed_at":observed_at})
}
Err(_) => serde_json::json!({"status":"unknown","observed_at":observed_at}),
},
None => serde_json::json!({"status":"requires_repair","observed_at":observed_at}),
};
let target = match reader.observe_merge_target(&mr.repository_id, Some(&mr.selector_to)) {
Ok(value) => {
serde_json::json!({"status":"known","ref":value.commit,"observed_at":observed_at})
}
Err(_) => serde_json::json!({"status":"unknown","observed_at":observed_at}),
};
let mut response =
serde_json::to_value(mr).map_err(|error| Error::InvalidInput(error.to_string()))?;
if let Some(object) = response.as_object_mut() {
object.insert("source".into(), source);
object.insert("target".into(), target);
}
Ok(Json(response))
}
async fn scoped_merge_request_readiness(
@@ -3815,9 +3762,15 @@ async fn scoped_merge_request_readiness(
let workspace_id = parse_workspace_id(&workspace_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 {
ticket_id,
expected_head_commit: None,
current_subject_ref,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: mr.repository_id,
@@ -3847,117 +3800,94 @@ async fn scoped_open_merge_request(
|| assignment.worker.worker_id != source.worker_id
{
return Err(Error::TicketAssignmentConflict(
"authenticated Worker is not the current Ticket assignee".into(),
"authenticated Worker is not current assignee".into(),
)
.into());
}
let (base_commit, source_commit, target) = validate_open_merge_request_evidence(
&api,
&ticket_id,
&input.repository_id,
&input.base_commit,
&input.head_commit,
)?;
if input.selector_to != target.selector {
let ticket = browser_ticket_backend(&api)?
.show(TicketIdOrSlug::Id(ticket_id.clone().into()))
.map_err(Error::from)?;
if ticket.meta.repository_id.as_deref() != Some(input.repository_id.as_str())
|| ticket.meta.ref_selector.as_deref() != Some(input.selector_to.as_str())
{
return Err(Error::InvalidInput(
"selector_to must match the authoritative Ticket target selector".into(),
"selectors must match the authoritative Ticket repository target".into(),
)
.into());
}
let observed_source = api
.repository_reader()
let reader = api.repository_reader();
reader
.observe_merge_target(&input.repository_id, Some(&input.selector_from))
.map_err(repository_merge_evidence_error)?;
if observed_source.commit != source_commit {
return Err(Error::InvalidInput(
"selector_from does not resolve to the nominated head commit".into(),
)
.into());
}
let mr = merge_request_store(&api, &workspace_id)?.open_merge_request(
merge_request::OpenMergeRequest {
merge_request_id: Uuid::now_v7().to_string(),
ticket_id,
repository_id: input.repository_id.clone(),
selector_from: input.selector_from,
selector_to: input.selector_to,
request: merge_request::RequestForReview {
base_commit,
head_commit: source_commit,
changed_paths: input.changed_paths,
reader
.observe_merge_target(&input.repository_id, Some(&input.selector_to))
.map_err(repository_merge_evidence_error)?;
Ok(Json(
merge_request_store(&api, &workspace_id)?.open_merge_request(
merge_request::OpenMergeRequest {
merge_request_id: Uuid::now_v7().to_string(),
ticket_id,
repository_id: input.repository_id.clone(),
selector_from: input.selector_from,
selector_to: input.selector_to,
summary: input.summary,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: input.repository_id,
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: assignment.assignment_id,
},
now: Utc::now(),
},
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: input.repository_id,
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: assignment.assignment_id,
},
now: Utc::now(),
},
)?;
Ok(Json(mr))
)?,
))
}
async fn scoped_request_merge_request_review(
async fn scoped_merge_request_thread(
State(api): State<WorkspaceApi>,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
Query(query): Query<MergeRequestThreadQuery>,
) -> ApiResult<Json<Vec<merge_request::MergeRequestThreadEvent>>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
Ok(Json(
merge_request_store(&api, &workspace_id)?.thread_page(
&workspace_id,
&ticket_id,
query.after,
query.limit.unwrap_or(100),
)?,
))
}
async fn scoped_repair_merge_request_selector(
State(api): State<WorkspaceApi>,
headers: HeaderMap,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
Json(input): Json<RequestMergeRequestReviewRequest>,
) -> ApiResult<Json<merge_request::RequestForReviewEvent>> {
Json(input): Json<RepairMergeRequestSelectorRequest>,
) -> ApiResult<Json<merge_request::MergeRequest>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
require_workspace_access(&workspace_id, &api)?;
let source = authenticate_worker_mutation_source(&api, &workspace_id, &headers)?;
let assignment = api
.store
.get_current_ticket_worker_assignment(&workspace_id, &ticket_id)?
.ok_or_else(|| {
Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into())
})?;
if assignment.worker.runtime_id != source.runtime_id
|| assignment.worker.worker_id != source.worker_id
{
return Err(Error::TicketAssignmentConflict(
"authenticated Worker is not the current Ticket assignee".into(),
)
.into());
reject_non_browser_reopen_auth(&headers)?;
let _actor = require_actor(&api, &headers).await?;
if !input.explicit_confirmation {
return Err(Error::BrowserReopenConfirmationRequired.into());
}
let store = merge_request_store(&api, &workspace_id)?;
let current = store.get(&workspace_id, &ticket_id)?;
let (base_commit, source_commit) = validate_revision_evidence(
&api,
&current.repository_id,
&input.base_commit,
&input.head_commit,
)?;
let observed_source = api
.repository_reader()
.observe_merge_target(&current.repository_id, Some(&current.selector_from))
let mr = store.get(&workspace_id, &ticket_id)?;
api.repository_reader()
.observe_merge_target(&mr.repository_id, Some(&input.selector_from))
.map_err(repository_merge_evidence_error)?;
if observed_source.commit != source_commit {
return Err(Error::InvalidInput(
"selector_from does not resolve to the nominated head commit".into(),
)
.into());
}
Ok(Json(store.request_review(
merge_request::RequestMergeRequestReview {
Ok(Json(store.repair_selector_from(
merge_request::RepairSelectorFrom {
workspace_id,
ticket_id,
expected_head_commit: input.expected_head_commit,
request: merge_request::RequestForReview {
base_commit,
head_commit: source_commit,
changed_paths: input.changed_paths,
summary: input.summary,
},
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: current.repository_id,
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: assignment.assignment_id,
selector_from: input.selector_from,
repaired_by: merge_request::WorkerIdentity {
runtime_id: "browser".into(),
worker_id: "authenticated-user".into(),
},
reason: input.reason,
now: Utc::now(),
},
)?))
@@ -4004,20 +3934,29 @@ async fn scoped_register_merge_request_review_capability(
|| assignment.worker.worker_id != source.worker_id
{
return Err(Error::TicketAssignmentConflict(
"authenticated Worker is not the current Ticket assignee".into(),
"authenticated Worker is not current assignee".into(),
)
.into());
}
let store = merge_request_store(&api, &workspace_id)?;
let current = store.get(&workspace_id, &ticket_id)?;
store.register_review_capability(merge_request::RegisterReviewCapability {
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 subject_ref = api
.repository_reader()
.observe_merge_target(&mr.repository_id, Some(selector))
.map_err(repository_merge_evidence_error)?
.commit;
store.request_review(merge_request::RequestMergeRequestReview {
ticket_id,
expected_head_commit: input.expected_head_commit,
subject_ref,
child_session_id: input.child_session_id,
capability_token: input.capability_token,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: current.repository_id,
repository_id: mr.repository_id,
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: assignment.assignment_id,
@@ -4033,19 +3972,73 @@ async fn scoped_submit_merge_request_review(
Json(input): Json<SubmitMergeRequestReviewRequest>,
) -> ApiResult<Json<merge_request::ReviewEvent>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
Ok(Json(
merge_request_store(&api, &workspace_id)?.submit_review(
merge_request::SubmitMergeRequestReview {
ticket_id,
expected_head_commit: input.expected_head_commit,
capability_token: input.capability_token,
decision: input.decision,
body: input.body,
findings: input.findings,
now: Utc::now(),
let store = merge_request_store(&api, &workspace_id)?;
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;
Ok(Json(store.submit_review(
merge_request::SubmitMergeRequestReview {
ticket_id,
current_subject_ref,
capability_token: input.capability_token,
decision: input.decision,
body: input.body,
findings: input.findings,
now: Utc::now(),
},
)?))
}
async fn scoped_revoke_merge_request_review(
State(api): State<WorkspaceApi>,
headers: HeaderMap,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
Json(input): Json<RevokeMergeRequestReviewRequest>,
) -> ApiResult<Json<merge_request::ReviewRevokedEvent>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
require_workspace_access(&workspace_id, &api)?;
if !input.explicit_confirmation {
return Err(Error::BrowserReopenConfirmationRequired.into());
}
let source = authenticate_worker_mutation_source(&api, &workspace_id, &headers)?;
let assignment = api
.store
.get_current_ticket_worker_assignment(&workspace_id, &ticket_id)?
.ok_or_else(|| {
Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into())
})?;
if assignment.worker.runtime_id != source.runtime_id
|| assignment.worker.worker_id != source.worker_id
{
return Err(Error::TicketAssignmentConflict(
"authenticated Worker is not current assignee".into(),
)
.into());
}
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
Ok(Json(store.revoke_review(
merge_request::RevokeMergeRequestReview {
ticket_id,
review_event_id: input.review_event_id,
reason: input.reason,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: mr.repository_id,
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: assignment.assignment_id,
},
)?,
))
now: Utc::now(),
},
)?))
}
async fn scoped_complete_merge_request(
@@ -4066,48 +4059,42 @@ async fn scoped_complete_merge_request(
})?;
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let current_request = mr.current_request().ok_or_else(|| {
Error::InvalidInput("Merge Request has no current RequestForReview event".into())
})?;
if current_request.head_commit != input.expected_head_commit
|| input.source_commit != current_request.head_commit
{
let selector = mr
.selector_from
.as_deref()
.ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?;
let repositories = api.repository_reader();
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)?;
if observed.commit != input.target_ref_before && observed.commit != input.target_ref_after {
return Err(Error::InvalidInput(
"source commit does not match the current review request".into(),
"target selector moved outside completion evidence".into(),
)
.into());
}
let repositories = api.repository_reader();
let observed_target = repositories
.observe_merge_target(&mr.repository_id, Some(&mr.selector_to))
.map_err(repository_merge_evidence_error)?;
if observed_target.commit != input.target_commit
&& observed_target.commit != input.result_commit
{
return Err(Error::InvalidInput(format!(
"Merge Request target moved: expected {}, observed {}",
input.target_commit, observed_target.commit
))
.into());
}
let target_was_already_updated = observed_target.commit == input.result_commit;
if !target_was_already_updated {
let already = observed.commit == input.target_ref_after;
if !already {
repositories
.update_merge_target(
&mr.repository_id,
&mr.selector_to,
&input.target_commit,
&input.result_commit,
&input.target_ref_before,
&input.target_ref_after,
)
.map_err(repository_merge_evidence_error)?;
.map_err(repository_merge_evidence_error)?
}
let completion = merge_request::CompleteMergeRequest {
ticket_id,
expected_head_commit: input.expected_head_commit,
operation_id: input.operation_id,
target_commit: input.target_commit.clone(),
source_commit: input.source_commit,
result_commit: input.result_commit.clone(),
approval_event_id: input.approval_event_id,
current_subject_ref: current_source_ref,
target_ref_before: input.target_ref_before.clone(),
target_ref_after: input.target_ref_after.clone(),
strategy: input.strategy,
resolution: input.resolution,
auth: merge_request::MergeRequestAuth {
@@ -4120,75 +4107,21 @@ async fn scoped_complete_merge_request(
now: Utc::now(),
};
match store.complete(completion) {
Ok(event) => Ok(Json(event)),
Err(error) => {
if !target_was_already_updated {
Ok(v) => Ok(Json(v)),
Err(e) => {
if !already {
let _ = repositories.update_merge_target(
&mr.repository_id,
&mr.selector_to,
&input.result_commit,
&input.target_commit,
&input.target_ref_after,
&input.target_ref_before,
);
}
Err(error.into())
Err(e.into())
}
}
}
async fn scoped_close_merge_request(
State(api): State<WorkspaceApi>,
headers: HeaderMap,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
Json(input): Json<MergeRequestStateRequest>,
) -> ApiResult<Json<merge_request::MergeRequest>> {
scoped_change_merge_request_state(api, headers, workspace_id, ticket_id, input, false).await
}
async fn scoped_reopen_merge_request(
State(api): State<WorkspaceApi>,
headers: HeaderMap,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
Json(input): Json<MergeRequestStateRequest>,
) -> ApiResult<Json<merge_request::MergeRequest>> {
scoped_change_merge_request_state(api, headers, workspace_id, ticket_id, input, true).await
}
async fn scoped_change_merge_request_state(
api: WorkspaceApi,
headers: HeaderMap,
workspace_id: String,
ticket_id: String,
input: MergeRequestStateRequest,
reopen: bool,
) -> ApiResult<Json<merge_request::MergeRequest>> {
let workspace_id = parse_workspace_id(&workspace_id)?;
require_workspace_access(&workspace_id, &api)?;
reject_non_browser_reopen_auth(&headers)?;
let _actor = require_actor(&api, &headers).await?;
if !input.explicit_confirmation {
return Err(Error::BrowserReopenConfirmationRequired.into());
}
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let operation = merge_request::ChangeMergeRequestState {
ticket_id,
body: input.body,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: mr.repository_id,
runtime_id: "browser".into(),
worker_id: "authenticated-user".into(),
assignment_id: String::new(),
},
now: Utc::now(),
};
Ok(Json(if reopen {
store.reopen(operation)?
} else {
store.close(operation)?
}))
}
fn reject_non_browser_reopen_auth(headers: &HeaderMap) -> Result<()> {
if headers.contains_key("authorization") {
return Err(Error::BrowserReopenConfirmationRequired);
@@ -12738,241 +12671,6 @@ mod tests {
));
}
#[tokio::test]
async fn merge_request_completion_endpoint_rejects_coder_and_accepts_orchestrator() {
let workspace = tempfile::tempdir().unwrap();
init_clean_git_workspace(workspace.path());
let git_value = |args: &[&str]| {
let output = std::process::Command::new("git")
.arg("-C")
.arg(workspace.path())
.args(args)
.output()
.unwrap();
assert!(output.status.success());
String::from_utf8(output.stdout).unwrap().trim().to_string()
};
let target_commit = git_value(&["rev-parse", "HEAD"]);
let target_ref = git_value(&["symbolic-ref", "HEAD"]);
std::fs::write(workspace.path().join("README.md"), "merge source\n").unwrap();
for args in [&["add", "README.md"][..], &["commit", "-m", "source"][..]] {
assert!(
std::process::Command::new("git")
.arg("-C")
.arg(workspace.path())
.args(args)
.status()
.unwrap()
.success()
);
}
let source_commit = git_value(&["rev-parse", "HEAD"]);
assert!(
std::process::Command::new("git")
.arg("-C")
.arg(workspace.path())
.args(["reset", "--hard", &target_commit])
.status()
.unwrap()
.success()
);
let api = test_api(workspace.path()).await;
let workspace_id = api.config.workspace_id.clone();
let backend = browser_ticket_backend(&api).unwrap();
let mut input = ticket::NewTicket::new("Orchestrator completion authority");
input.workflow_state = Some(TicketWorkflowState::InProgress);
let ticket = backend.create(input).unwrap();
let Json(coder) = create_workspace_worker(
State(api.clone()),
HeaderMap::new(),
Json(CreateWorkspaceWorkerRequest {
runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(),
display_name: "Assigned Coder".to_string(),
profile: Some("builtin:coder".to_string()),
ticket_assignment: Some(CreateWorkspaceWorkerTicketAssignmentRequest {
ticket_id: ticket.id.clone(),
operation_id: "completion-coder-assignment".to_string(),
}),
initial_submit: vec![Segment::Flow {
selector: "builtin:coder-review".to_string(),
}],
working_directory: None,
control_operation_id: None,
resolved_control_operation: None,
}),
)
.await
.unwrap();
let assignment = api
.store
.get_current_ticket_worker_assignment(&workspace_id, &ticket.id)
.unwrap()
.unwrap();
let mr_store = merge_request_store(&api, &workspace_id).unwrap();
mr_store
.open_merge_request(merge_request::OpenMergeRequest {
merge_request_id: "MR-server-completion".into(),
ticket_id: ticket.id.clone(),
repository_id: TEST_REPOSITORY_ID.into(),
selector_from: source_commit.clone(),
selector_to: target_ref.clone(),
request: merge_request::RequestForReview {
base_commit: target_commit.clone(),
head_commit: source_commit.clone(),
changed_paths: vec!["src/lib.rs".into()],
summary: "approved candidate".into(),
},
auth: merge_request::MergeRequestAuth {
workspace_id: workspace_id.clone(),
repository_id: TEST_REPOSITORY_ID.into(),
runtime_id: coder.worker_ref.runtime_id.clone(),
worker_id: coder.worker_ref.worker_id.clone(),
assignment_id: assignment.assignment_id.clone(),
},
now: Utc::now(),
})
.unwrap();
mr_store
.register_reviewer_child_session(merge_request::RegisterReviewerChildSession {
workspace_id: workspace_id.clone(),
parent_runtime_id: coder.worker_ref.runtime_id.clone(),
parent_worker_id: coder.worker_ref.worker_id.clone(),
child_session_id: "reviewer-child".into(),
reviewer_profile: "builtin:reviewer".into(),
now: Utc::now(),
})
.unwrap();
mr_store
.register_review_capability(merge_request::RegisterReviewCapability {
ticket_id: ticket.id.clone(),
expected_head_commit: source_commit.clone(),
child_session_id: "reviewer-child".into(),
capability_token: "review-token".into(),
auth: merge_request::MergeRequestAuth {
workspace_id: workspace_id.clone(),
repository_id: TEST_REPOSITORY_ID.into(),
runtime_id: coder.worker_ref.runtime_id.clone(),
worker_id: coder.worker_ref.worker_id.clone(),
assignment_id: assignment.assignment_id.clone(),
},
now: Utc::now(),
})
.unwrap();
mr_store
.submit_review(merge_request::SubmitMergeRequestReview {
ticket_id: ticket.id.clone(),
expected_head_commit: source_commit.clone(),
capability_token: "review-token".into(),
decision: merge_request::ReviewDecision::Approve,
body: "approved".into(),
findings: Vec::new(),
now: Utc::now(),
})
.unwrap();
let worker_headers = |worker: &RuntimeWorkerRef| {
let mut headers = HeaderMap::new();
headers.insert(
"x-yoi-runtime-id",
axum::http::HeaderValue::from_str(&worker.runtime_id).unwrap(),
);
headers.insert(
"x-yoi-worker-id",
axum::http::HeaderValue::from_str(&worker.worker_id).unwrap(),
);
headers
};
let request = || CompleteMergeRequestRequest {
operation_id: "complete-operation".into(),
expected_head_commit: source_commit.clone(),
target_commit: target_commit.clone(),
source_commit: source_commit.clone(),
result_commit: source_commit.clone(),
strategy: merge_request::MergeStrategy::FastForward,
resolution: merge_request::ConflictResolution::None,
};
let coder_error = scoped_complete_merge_request(
State(api.clone()),
worker_headers(&coder.worker_ref),
AxumPath((workspace_id.clone(), ticket.id.clone())),
Json(request()),
)
.await
.unwrap_err();
assert!(matches!(
coder_error.error,
Error::TicketAssignmentConflict(_)
));
let Json(started) = scoped_start_workspace_orchestrator(
State(api.clone()),
AxumPath(ScopedWorkspacePath {
workspace_id: workspace_id.clone(),
}),
)
.await
.unwrap();
let orchestrator = started.worker.unwrap().worker;
let Json(completed) = scoped_complete_merge_request(
State(api.clone()),
worker_headers(&orchestrator),
AxumPath((workspace_id, ticket.id.clone())),
Json(request()),
)
.await
.unwrap();
assert_eq!(completed.result_commit, source_commit);
assert_eq!(
backend
.show(ticket.id.clone().into())
.unwrap()
.meta
.workflow_state,
TicketWorkflowState::Done
);
assert_eq!(
api.repository_reader()
.observe_merge_target(TEST_REPOSITORY_ID, Some(&target_ref))
.unwrap()
.commit,
source_commit
);
let mr = mr_store.get(&api.config.workspace_id, &ticket.id).unwrap();
assert_eq!(mr.state, merge_request::MergeRequestState::Merged);
let merge = mr.thread.iter().find_map(|event| match event {
merge_request::MergeRequestThreadEvent::Merge(value) => Some(value),
_ => None,
});
assert_eq!(
merge.map(|value| value.result_commit.as_str()),
Some(source_commit.as_str())
);
let Json(replayed) = scoped_complete_merge_request(
State(api.clone()),
worker_headers(&orchestrator),
AxumPath((api.config.workspace_id.clone(), ticket.id.clone())),
Json(request()),
)
.await
.unwrap();
assert_eq!(replayed.result_commit, source_commit);
let conn = rusqlite::Connection::open(&api.config.database_path).unwrap();
let actor: String = conn
.query_row(
"SELECT author FROM typed_ticket_events WHERE workspace_id=?1 AND ticket_id=?2 AND kind='state_changed'",
rusqlite::params![api.config.workspace_id, ticket.id],
|row| row.get(0),
)
.unwrap();
assert_eq!(
actor,
format!(
"worker:{}:{}",
orchestrator.runtime_id, orchestrator.worker_id
)
);
}
#[tokio::test]
async fn production_profile_backend_launches_and_restores_workspace_orchestrator() {
let workspace = tempfile::tempdir().unwrap();