From 027f60d262f23dddaccdd439cc0988c3cb65b727 Mon Sep 17 00:00:00 2001 From: Hare Date: Sat, 15 Aug 2026 18:02:54 +0900 Subject: [PATCH] merge-request: add target and merge result authority --- crates/merge-request/src/lib.rs | 599 ++++++++++++++++-- crates/merge-request/tests/store.rs | 209 +++++- .../src/feature/builtin/merge_request.rs | 78 ++- crates/worker/src/spawn/tool.rs | 14 +- crates/worker/src/worker.rs | 10 + crates/workspace-server/src/repositories.rs | 240 +++++++ crates/workspace-server/src/server.rs | 310 ++++++++- .../tickets/[ticketId]/+page.svelte | 52 +- 8 files changed, 1439 insertions(+), 73 deletions(-) diff --git a/crates/merge-request/src/lib.rs b/crates/merge-request/src/lib.rs index 86662fb7..cd169cda 100644 --- a/crates/merge-request/src/lib.rs +++ b/crates/merge-request/src/lib.rs @@ -11,7 +11,7 @@ use std::path::{Path, PathBuf}; use std::time::Duration; use thiserror::Error; -const SCHEMA_VERSION: i64 = 8; +const SCHEMA_VERSION: i64 = 9; const REVIEWER_PROFILE: &str = "builtin:reviewer"; const MAX_SUMMARY_BYTES: usize = 16 * 1024; const MAX_REVIEW_BODY_BYTES: usize = 64 * 1024; @@ -51,6 +51,14 @@ pub enum MergeRequestError { NotOpen(String), #[error("completion operation id was reused with different input")] OperationConflict, + #[error("merge result operation id was reused with different input")] + MergeResultOperationConflict, + #[error("merge result {0} was not found for the current Merge Request revision")] + MergeResultNotFound(String), + #[error("merge result evidence is invalid: {0}")] + InvalidMergeResult(String), + #[error("Merge Request target is unknown and must be resolved explicitly")] + UnknownTarget, #[error("Ticket must be inprogress before Merge Request completion (current: {0})")] TicketStateConflict(String), #[error("only an authenticated user with explicit confirmation may merge")] @@ -117,13 +125,86 @@ pub enum ReviewStatus { ChangesRequested, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum MergeRequestTargetStatus { + Known, + Unknown, +} + +impl MergeRequestTargetStatus { + fn parse(value: &str) -> Self { + match value { + "known" => Self::Known, + _ => Self::Unknown, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum MergeStrategy { + FastForward, + Merge, +} + +impl MergeStrategy { + fn as_str(self) -> &'static str { + match self { + Self::FastForward => "fast_forward", + Self::Merge => "merge", + } + } + + fn parse(value: &str) -> Self { + match value { + "fast_forward" => Self::FastForward, + _ => Self::Merge, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum MergeResolution { + None, + Clean, + ConflictsResolved, +} + +impl MergeResolution { + fn as_str(self) -> &'static str { + match self { + Self::None => "none", + Self::Clean => "clean", + Self::ConflictsResolved => "conflicts_resolved", + } + } + + fn parse(value: &str) -> Self { + match value { + "none" => Self::None, + "conflicts_resolved" => Self::ConflictsResolved, + _ => Self::Clean, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum MergeResultTargetStatus { + Current, + Applied, + Stale, + Unknown, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct MergeRequestRevision { pub revision_id: String, pub ordinal: u64, pub base_commit: String, pub head_commit: String, - pub head_tree: String, pub diff_digest: String, pub changed_paths: Vec, pub summary: String, @@ -144,6 +225,8 @@ pub struct ReviewFinding { pub struct MergeRequestReview { pub attempt_id: String, pub revision_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub merge_result_id: Option, pub decision: ReviewDecision, pub body: String, pub findings: Vec, @@ -155,17 +238,46 @@ pub struct MergeRequestReview { pub submitted_at: String, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct MergeResult { + pub merge_result_id: String, + pub revision_id: String, + pub target_commit: String, + pub source_commit: String, + pub result_commit: String, + pub strategy: MergeStrategy, + pub resolution: MergeResolution, + pub created_by_runtime_id: String, + pub created_by_worker_id: String, + pub created_at: String, + pub operation_id: String, + pub validated_at: String, + pub target_status: MergeResultTargetStatus, + pub review_status: ReviewStatus, + pub current_review: Option, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct MergeRequest { pub merge_request_id: String, pub workspace_id: String, pub ticket_id: String, pub repository_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub target_ref_selector: Option, + pub target_status: MergeRequestTargetStatus, + #[serde(skip_serializing_if = "Option::is_none")] + pub observed_target_commit: Option, pub state: MergeRequestState, pub lifecycle_generation: u64, pub current_revision: MergeRequestRevision, pub review_status: ReviewStatus, pub current_review: Option, + pub merge_results: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + pub current_merge_result: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub applied_merge_result: Option, pub created_at: String, pub updated_at: String, pub merged_by_account_id: Option, @@ -177,6 +289,7 @@ pub struct OpenMergeRequest { pub merge_request_id: String, pub ticket_id: String, pub repository_id: String, + pub target_ref_selector: String, pub revision: MergeRequestRevision, pub authenticated_runtime_id: String, pub authenticated_worker_id: String, @@ -206,6 +319,7 @@ pub struct RegisterReviewAttempt { pub attempt_id: String, pub ticket_id: String, pub revision_id: String, + pub merge_result_id: Option, pub parent_assignment_id: String, pub parent_runtime_id: String, pub parent_worker_id: String, @@ -219,6 +333,7 @@ pub struct RegisterReviewAttempt { pub struct SubmitReview { pub ticket_id: String, pub revision_id: String, + pub merge_result_id: Option, pub capability_token: String, pub decision: ReviewDecision, pub body: String, @@ -226,6 +341,28 @@ pub struct SubmitReview { pub now: String, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RecordMergeResult { + pub merge_result_id: String, + pub ticket_id: String, + pub expected_revision_id: String, + pub target_commit: String, + pub source_commit: String, + pub result_commit: String, + pub strategy: MergeStrategy, + pub resolution: MergeResolution, + pub operation_id: String, + pub actor_runtime_id: String, + pub actor_worker_id: String, + pub created_at: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RecordMergeResultOutcome { + pub merge_result: MergeResult, + pub replayed: bool, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct CompleteMergeRequest { pub operation_id: String, @@ -251,8 +388,12 @@ pub struct MergeRequestReadiness { pub ticket_id: String, pub merge_request_id: String, pub revision_id: String, + pub target_ref_selector: Option, + pub observed_target_commit: Option, pub ready: bool, pub review_status: ReviewStatus, + pub merge_result_id: Option, + pub merge_result_review_status: Option, pub blockers: Vec, } @@ -327,33 +468,87 @@ impl SqliteMergeRequestStore { } pub fn show_for_ticket(&self, ticket_id: &str) -> Result> { + self.show_for_ticket_with_target(ticket_id, None) + } + + pub fn show_for_ticket_with_target( + &self, + ticket_id: &str, + observed_target_commit: Option<&str>, + ) -> Result> { nonempty("ticket_id", ticket_id)?; let conn = self.connect()?; verify(&conn)?; - load_merge_request(&conn, &self.workspace_id, ticket_id) + let mut mr = load_merge_request(&conn, &self.workspace_id, ticket_id)?; + if let Some(mr) = mr.as_mut() { + apply_target_observation(mr, observed_target_commit); + } + Ok(mr) } pub fn readiness_for_ticket(&self, ticket_id: &str) -> Result { + self.readiness_for_ticket_with_target(ticket_id, None) + } + + pub fn readiness_for_ticket_with_target( + &self, + ticket_id: &str, + observed_target_commit: Option<&str>, + ) -> Result { let mr = self - .show_for_ticket(ticket_id)? + .show_for_ticket_with_target(ticket_id, observed_target_commit)? .ok_or_else(|| MergeRequestError::NotFound(ticket_id.to_string()))?; let mut blockers = Vec::new(); if mr.state != MergeRequestState::Open { blockers.push(format!("merge request is {}", mr.state.as_str())); } match mr.review_status { - ReviewStatus::Pending => blockers.push("current revision has no review result".into()), + ReviewStatus::Pending => { + blockers.push("current source revision has no review result".into()) + } ReviewStatus::ChangesRequested => { - blockers.push("current revision has request_changes".into()) + blockers.push("current source revision has request_changes".into()) } ReviewStatus::Approved => {} } + if mr.target_status != MergeRequestTargetStatus::Known || observed_target_commit.is_none() { + blockers.push("merge request target is unknown or could not be resolved".into()); + } + let current_result = mr + .current_merge_result + .as_ref() + .or(mr.applied_merge_result.as_ref()); + if current_result.is_none() && observed_target_commit.is_some() { + if mr + .merge_results + .iter() + .any(|result| result.target_status == MergeResultTargetStatus::Stale) + { + blockers.push("target moved after the latest MergeResult was recorded".into()); + } else { + blockers.push( + "current source revision has no validated MergeResult for the current target" + .into(), + ); + } + } + if let Some(result) = current_result { + if matches!(result.strategy, MergeStrategy::Merge) + && result.review_status != ReviewStatus::Approved + { + blockers.push("non-fast-forward MergeResult is not independently approved".into()); + } + } Ok(MergeRequestReadiness { ticket_id: ticket_id.to_string(), merge_request_id: mr.merge_request_id, revision_id: mr.current_revision.revision_id, + target_ref_selector: mr.target_ref_selector, + observed_target_commit: mr.observed_target_commit, ready: blockers.is_empty(), review_status: mr.review_status, + merge_result_id: current_result.map(|result| result.merge_result_id.clone()), + merge_result_review_status: current_result.map(|result| result.review_status), blockers, }) } @@ -364,6 +559,7 @@ impl SqliteMergeRequestStore { ("merge_request_id", input.merge_request_id.as_str()), ("ticket_id", input.ticket_id.as_str()), ("repository_id", input.repository_id.as_str()), + ("target_ref_selector", input.target_ref_selector.as_str()), ("runtime_id", input.authenticated_runtime_id.as_str()), ("worker_id", input.authenticated_worker_id.as_str()), ] { @@ -379,8 +575,8 @@ impl SqliteMergeRequestStore { &input.authenticated_worker_id, )?; conn.execute( - "INSERT INTO merge_requests (workspace_id, merge_request_id, repository_id, state, lifecycle_generation, current_revision_id, created_at, updated_at) VALUES (?1, ?2, ?3, 'open', 1, ?4, ?5, ?5)", - params![self.workspace_id, input.merge_request_id, input.repository_id, input.revision.revision_id, input.now], + "INSERT INTO merge_requests (workspace_id, merge_request_id, repository_id, target_ref_selector, target_status, state, lifecycle_generation, current_revision_id, created_at, updated_at) VALUES (?1, ?2, ?3, ?4, 'known', 'open', 1, ?5, ?6, ?6)", + params![self.workspace_id, input.merge_request_id, input.repository_id, input.target_ref_selector, input.revision.revision_id, input.now], ).map_err(db)?; conn.execute( "INSERT INTO merge_request_ticket_relations (workspace_id,merge_request_id,ticket_id,relation_kind,created_at) VALUES (?1,?2,?3,'implements',?4)", @@ -414,13 +610,13 @@ impl SqliteMergeRequestStore { if input.revision.ordinal != current.current_revision.ordinal + 1 { return Err(MergeRequestError::RevisionConflict(input.revision.revision_id.clone())); } - let existing: Option<(String, String, String, String)> = conn.query_row( - "SELECT base_commit, head_commit, head_tree, diff_digest FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3", + let existing: Option<(String, String, String)> = conn.query_row( + "SELECT base_commit, head_commit, diff_digest FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3", params![self.workspace_id, current.merge_request_id, input.revision.revision_id], - |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)), + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), ).optional().map_err(db)?; if let Some(existing) = existing { - if existing == (input.revision.base_commit.clone(), input.revision.head_commit.clone(), input.revision.head_tree.clone(), input.revision.diff_digest.clone()) { + if existing == (input.revision.base_commit.clone(), input.revision.head_commit.clone(), input.revision.diff_digest.clone()) { return Ok(current); } return Err(MergeRequestError::RevisionConflict(input.revision.revision_id.clone())); @@ -468,7 +664,19 @@ impl SqliteMergeRequestStore { if mr.current_revision.revision_id != input.revision_id { return Err(MergeRequestError::StaleRevision { expected: input.revision_id.clone(), current: mr.current_revision.revision_id }); } - validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id, &input.parent_runtime_id, &input.parent_worker_id)?; + if let Some(merge_result_id) = input.merge_result_id.as_deref() { + let exists: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM merge_request_merge_results WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3 AND merge_result_id=?4)", + params![self.workspace_id,mr.merge_request_id,input.revision_id,merge_result_id], + |row| row.get(0), + ).map_err(db)?; + if !exists { + return Err(MergeRequestError::MergeResultNotFound(merge_result_id.into())); + } + validate_current_assignment_id(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id)?; + } else { + validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id, &input.parent_runtime_id, &input.parent_worker_id)?; + } let effective_profile: Option = conn.query_row( "SELECT effective_profile FROM merge_request_reviewer_child_sessions WHERE workspace_id=?1 AND child_session_id=?2 AND parent_runtime_id=?3 AND parent_worker_id=?4", params![self.workspace_id,input.child_session_id,input.parent_runtime_id,input.parent_worker_id], @@ -478,8 +686,8 @@ impl SqliteMergeRequestStore { return Err(MergeRequestError::InvalidReviewer); } conn.execute( - "INSERT INTO merge_request_review_attempts (workspace_id, attempt_id, merge_request_id, ticket_id, revision_id, lifecycle_generation, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, capability_token_sha256, status, created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,'open',?13)", - params![self.workspace_id, input.attempt_id, mr.merge_request_id, input.ticket_id, input.revision_id, mr.lifecycle_generation as i64, input.parent_assignment_id, input.parent_runtime_id, input.parent_worker_id, input.child_session_id, REVIEWER_PROFILE, token_hash(&input.capability_token), input.now], + "INSERT INTO merge_request_review_attempts (workspace_id, attempt_id, merge_request_id, ticket_id, revision_id, merge_result_id, lifecycle_generation, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, capability_token_sha256, status, created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,'open',?14)", + params![self.workspace_id, input.attempt_id, mr.merge_request_id, input.ticket_id, input.revision_id, input.merge_result_id, mr.lifecycle_generation as i64, input.parent_assignment_id, input.parent_runtime_id, input.parent_worker_id, input.child_session_id, REVIEWER_PROFILE, token_hash(&input.capability_token), input.now], ).map_err(|_| MergeRequestError::InvalidReviewAttempt)?; Ok(()) }) @@ -505,12 +713,12 @@ impl SqliteMergeRequestStore { validate_review_input(&input)?; self.write(|conn| { let token = token_hash(&input.capability_token); - let attempt: Option<(String,String,String,String,String,String,String,String,i64)> = conn.query_row( - "SELECT attempt_id, merge_request_id, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, status, lifecycle_generation FROM merge_request_review_attempts WHERE workspace_id=?1 AND ticket_id=?2 AND revision_id=?3 AND capability_token_sha256=?4", - params![self.workspace_id, input.ticket_id, input.revision_id, token], - |row| Ok((row.get(0)?,row.get(1)?,row.get(2)?,row.get(3)?,row.get(4)?,row.get(5)?,row.get(6)?,row.get(7)?,row.get(8)?)), + let attempt: Option<(String,String,Option,String,String,String,String,String,String,i64)> = conn.query_row( + "SELECT attempt_id, merge_request_id, merge_result_id, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, status, lifecycle_generation FROM merge_request_review_attempts WHERE workspace_id=?1 AND ticket_id=?2 AND revision_id=?3 AND ((?4 IS NULL AND merge_result_id IS NULL) OR merge_result_id=?4) AND capability_token_sha256=?5", + params![self.workspace_id, input.ticket_id, input.revision_id, input.merge_result_id, token], + |row| Ok((row.get(0)?,row.get(1)?,row.get(2)?,row.get(3)?,row.get(4)?,row.get(5)?,row.get(6)?,row.get(7)?,row.get(8)?,row.get(9)?)), ).optional().map_err(db)?; - let Some((attempt_id, mr_id, assignment_id, runtime_id, worker_id, child_session_id, effective_profile, status, lifecycle_generation)) = attempt else { + let Some((attempt_id, mr_id, merge_result_id, assignment_id, runtime_id, worker_id, child_session_id, effective_profile, status, lifecycle_generation)) = attempt else { return Err(MergeRequestError::InvalidReviewAttempt); }; if status != "open" || effective_profile != REVIEWER_PROFILE || child_session_id == worker_id { @@ -525,10 +733,14 @@ impl SqliteMergeRequestStore { if mr.current_revision.revision_id != input.revision_id { return Err(MergeRequestError::StaleRevision { expected: input.revision_id.clone(), current: mr.current_revision.revision_id }); } - validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &assignment_id, &runtime_id, &worker_id)?; + if merge_result_id.is_some() { + validate_current_assignment_id(conn, &self.workspace_id, &input.ticket_id, &assignment_id)?; + } else { + validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &assignment_id, &runtime_id, &worker_id)?; + } conn.execute( - "INSERT INTO merge_request_reviews (workspace_id, attempt_id, merge_request_id, revision_id, decision, body, submitted_at) VALUES (?1,?2,?3,?4,?5,?6,?7)", - params![self.workspace_id, attempt_id, mr_id, input.revision_id, input.decision.as_str(), input.body, input.now], + "INSERT INTO merge_request_reviews (workspace_id, attempt_id, merge_request_id, revision_id, merge_result_id, decision, body, submitted_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8)", + params![self.workspace_id, attempt_id, mr_id, input.revision_id, merge_result_id, input.decision.as_str(), input.body, input.now], ).map_err(|_| MergeRequestError::InvalidReviewAttempt)?; for (ordinal, finding) in input.findings.iter().enumerate() { nonempty("finding.body", &finding.body)?; @@ -545,6 +757,88 @@ impl SqliteMergeRequestStore { }) } + pub fn record_merge_result( + &self, + input: RecordMergeResult, + ) -> Result { + for (name, value) in [ + ("merge_result_id", input.merge_result_id.as_str()), + ("ticket_id", input.ticket_id.as_str()), + ("expected_revision_id", input.expected_revision_id.as_str()), + ("target_commit", input.target_commit.as_str()), + ("source_commit", input.source_commit.as_str()), + ("result_commit", input.result_commit.as_str()), + ("operation_id", input.operation_id.as_str()), + ("actor_runtime_id", input.actor_runtime_id.as_str()), + ("actor_worker_id", input.actor_worker_id.as_str()), + ] { + nonempty(name, value)?; + } + if matches!(input.strategy, MergeStrategy::FastForward) + && (input.result_commit != input.source_commit + || !matches!(input.resolution, MergeResolution::None)) + { + return Err(MergeRequestError::InvalidMergeResult( + "fast-forward result must equal the source commit and use resolution=none".into(), + )); + } + if matches!(input.strategy, MergeStrategy::Merge) + && matches!(input.resolution, MergeResolution::None) + { + return Err(MergeRequestError::InvalidMergeResult( + "merge strategy requires clean or conflicts_resolved resolution".into(), + )); + } + let fingerprint = merge_result_fingerprint(&input); + self.write(|conn| { + if let Some((stored, merge_result_id, generation)) = conn + .query_row( + "SELECT r.operation_fingerprint,r.merge_result_id,mr.lifecycle_generation FROM merge_request_merge_results r JOIN merge_requests mr ON mr.workspace_id=r.workspace_id AND mr.merge_request_id=r.merge_request_id WHERE r.workspace_id=?1 AND r.operation_id=?2", + params![self.workspace_id,input.operation_id], + |row| Ok((row.get::<_,String>(0)?,row.get::<_,String>(1)?,row.get::<_,i64>(2)?)), + ) + .optional() + .map_err(db)? + { + if stored != fingerprint { + return Err(MergeRequestError::MergeResultOperationConflict); + } + let merge_result = load_merge_result( + conn, + &self.workspace_id, + &merge_result_id, + generation as u64, + )? + .ok_or_else(|| MergeRequestError::MergeResultNotFound(merge_result_id.clone()))?; + return Ok(RecordMergeResultOutcome { merge_result, replayed: true }); + } + let mr = load_merge_request(conn, &self.workspace_id, &input.ticket_id)? + .ok_or_else(|| MergeRequestError::NotFound(input.ticket_id.clone()))?; + ensure_open(&mr)?; + if mr.target_status != MergeRequestTargetStatus::Known { + return Err(MergeRequestError::UnknownTarget); + } + if mr.current_revision.revision_id != input.expected_revision_id { + return Err(MergeRequestError::StaleRevision { + expected: input.expected_revision_id.clone(), + current: mr.current_revision.revision_id, + }); + } + if mr.current_revision.head_commit != input.source_commit { + return Err(MergeRequestError::InvalidMergeResult( + "source commit does not match the current source revision".into(), + )); + } + conn.execute( + "INSERT INTO merge_request_merge_results (workspace_id,merge_result_id,merge_request_id,ticket_id,revision_id,target_commit,source_commit,result_commit,strategy,resolution,created_by_runtime_id,created_by_worker_id,created_at,operation_id,operation_fingerprint,validated_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?13)", + params![self.workspace_id,input.merge_result_id,mr.merge_request_id,input.ticket_id,input.expected_revision_id,input.target_commit,input.source_commit,input.result_commit,input.strategy.as_str(),input.resolution.as_str(),input.actor_runtime_id,input.actor_worker_id,input.created_at,input.operation_id,fingerprint], + ).map_err(db)?; + let merge_result = load_merge_result(conn, &self.workspace_id, &input.merge_result_id, mr.lifecycle_generation)? + .ok_or_else(|| MergeRequestError::MergeResultNotFound(input.merge_result_id.clone()))?; + Ok(RecordMergeResultOutcome { merge_result, replayed: false }) + }) + } + pub fn complete(&self, input: CompleteMergeRequest) -> Result { for (name, value) in [ ("operation_id", input.operation_id.as_str()), @@ -714,6 +1008,9 @@ fn migrate_transaction(conn: &Connection) -> Result<()> { { migrate_completion_authority_v8(conn)?; } + if version < 9 && !column_exists(conn, "merge_requests", "target_ref_selector")? { + migrate_merge_result_authority_v9(conn)?; + } if version < 1 { conn.execute( "INSERT INTO merge_request_schema_migrations(version) VALUES (1)", @@ -764,6 +1061,54 @@ fn migrate_completion_authority_v8(conn: &Connection) -> Result<()> { .map_err(db) } +fn migrate_merge_result_authority_v9(conn: &Connection) -> Result<()> { + conn.pragma_update(None, "foreign_keys", "OFF") + .map_err(db)?; + let result = conn.execute_batch( + "ALTER TABLE merge_requests ADD COLUMN target_ref_selector TEXT; + ALTER TABLE merge_requests ADD COLUMN target_status TEXT NOT NULL DEFAULT 'unknown' CHECK(target_status IN ('known','unknown')); + + CREATE TABLE merge_request_revisions_v9 ( + workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL, + ordinal INTEGER NOT NULL, base_commit TEXT NOT NULL, head_commit TEXT NOT NULL, + diff_digest TEXT NOT NULL, summary TEXT NOT NULL, assignment_id TEXT NOT NULL, created_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,merge_request_id,revision_id), + UNIQUE(workspace_id,merge_request_id,ordinal), + FOREIGN KEY(workspace_id,merge_request_id) REFERENCES merge_requests(workspace_id,merge_request_id) ON DELETE CASCADE + ); + INSERT INTO merge_request_revisions_v9( + workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at + ) SELECT workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at + FROM merge_request_revisions; + DROP TABLE merge_request_revisions; + ALTER TABLE merge_request_revisions_v9 RENAME TO merge_request_revisions; + + CREATE TABLE merge_request_merge_results ( + workspace_id TEXT NOT NULL, merge_result_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, revision_id TEXT NOT NULL, target_commit TEXT NOT NULL, + source_commit TEXT NOT NULL, result_commit TEXT NOT NULL, + strategy TEXT NOT NULL CHECK(strategy IN ('fast_forward','merge')), + resolution TEXT NOT NULL CHECK(resolution IN ('none','clean','conflicts_resolved')), + created_by_runtime_id TEXT NOT NULL, created_by_worker_id TEXT NOT NULL, + created_at TEXT NOT NULL, operation_id TEXT NOT NULL, operation_fingerprint TEXT NOT NULL, + validated_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,merge_result_id), + UNIQUE(workspace_id,operation_id), + FOREIGN KEY(workspace_id,merge_request_id,revision_id) + REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id), + FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id) + ); + CREATE INDEX merge_request_merge_results_current_idx + ON merge_request_merge_results(workspace_id,merge_request_id,revision_id,target_commit,created_at); + + ALTER TABLE merge_request_review_attempts ADD COLUMN merge_result_id TEXT; + ALTER TABLE merge_request_reviews ADD COLUMN merge_result_id TEXT;", + ); + let restore = conn.pragma_update(None, "foreign_keys", "ON").map_err(db); + result.map_err(db)?; + restore +} + pub fn verify(conn: &Connection) -> Result<()> { let version: i64 = conn .query_row( @@ -786,6 +1131,7 @@ pub fn verify(conn: &Connection) -> Result<()> { "merge_request_review_attempts", "merge_request_reviews", "merge_request_review_findings", + "merge_request_merge_results", "merge_request_completion_operations", ] { let present: Option = conn @@ -809,6 +1155,8 @@ pub fn verify(conn: &Connection) -> Result<()> { "workspace_id", "merge_request_id", "repository_id", + "target_ref_selector", + "target_status", "state", "lifecycle_generation", "current_revision_id", @@ -832,7 +1180,6 @@ pub fn verify(conn: &Connection) -> Result<()> { "ordinal", "base_commit", "head_commit", - "head_tree", "diff_digest", "assignment_id", ] as &[_], @@ -855,6 +1202,7 @@ pub fn verify(conn: &Connection) -> Result<()> { "merge_request_id", "ticket_id", "revision_id", + "merge_result_id", "lifecycle_generation", "parent_assignment_id", "parent_runtime_id", @@ -872,10 +1220,29 @@ pub fn verify(conn: &Connection) -> Result<()> { "attempt_id", "merge_request_id", "revision_id", + "merge_result_id", "decision", "body", ] as &[_], ), + ( + "merge_request_merge_results", + &[ + "workspace_id", + "merge_result_id", + "merge_request_id", + "ticket_id", + "revision_id", + "target_commit", + "source_commit", + "result_commit", + "strategy", + "resolution", + "operation_id", + "operation_fingerprint", + "validated_at", + ] as &[_], + ), ( "merge_request_completion_operations", &[ @@ -991,6 +1358,8 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result "workspace_id", "merge_request_id", "repository_id", + "target_ref_selector", + "target_status", "state", "lifecycle_generation", "current_revision_id", @@ -1016,7 +1385,6 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result "ordinal", "base_commit", "head_commit", - "head_tree", "diff_digest", "assignment_id", ], @@ -1029,6 +1397,7 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result "merge_request_reviews", "merge_request_review_attempts", "merge_request_reviewer_child_sessions", + "merge_request_merge_results", "merge_request_completion_operations", "merge_request_revision_paths", "merge_request_ticket_relations", @@ -1111,14 +1480,16 @@ fn load_merge_request( workspace_id: &str, ticket_id: &str, ) -> Result> { - let row: Option<(String,String,String,String,i64,String,String,String,Option,Option)> = conn.query_row( - "SELECT mr.merge_request_id,rel.ticket_id,mr.repository_id,mr.state,mr.lifecycle_generation,mr.current_revision_id,mr.created_at,mr.updated_at,mr.merged_by_account_id,mr.merged_at FROM merge_requests mr JOIN merge_request_ticket_relations rel ON rel.workspace_id=mr.workspace_id AND rel.merge_request_id=mr.merge_request_id WHERE mr.workspace_id=?1 AND rel.ticket_id=?2 AND rel.relation_kind='implements' ORDER BY mr.updated_at DESC,mr.merge_request_id DESC LIMIT 1", - params![workspace_id,ticket_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?)), + let row: Option<(String,String,String,Option,String,String,i64,String,String,String,Option,Option)> = conn.query_row( + "SELECT mr.merge_request_id,rel.ticket_id,mr.repository_id,mr.target_ref_selector,mr.target_status,mr.state,mr.lifecycle_generation,mr.current_revision_id,mr.created_at,mr.updated_at,mr.merged_by_account_id,mr.merged_at FROM merge_requests mr JOIN merge_request_ticket_relations rel ON rel.workspace_id=mr.workspace_id AND rel.merge_request_id=mr.merge_request_id WHERE mr.workspace_id=?1 AND rel.ticket_id=?2 AND rel.relation_kind='implements' ORDER BY mr.updated_at DESC,mr.merge_request_id DESC LIMIT 1", + params![workspace_id,ticket_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?,r.get(10)?,r.get(11)?)), ).optional().map_err(db)?; let Some(( mr_id, ticket_id, repository_id, + target_ref_selector, + target_status, state, generation, revision_id, @@ -1137,16 +1508,24 @@ fn load_merge_request( Some(ReviewDecision::RequestChanges) => ReviewStatus::ChangesRequested, None => ReviewStatus::Pending, }; + let merge_results = + load_merge_results(conn, workspace_id, &mr_id, &revision_id, generation as u64)?; Ok(Some(MergeRequest { merge_request_id: mr_id, workspace_id: workspace_id.into(), ticket_id, repository_id, + target_ref_selector, + target_status: MergeRequestTargetStatus::parse(&target_status), + observed_target_commit: None, state: MergeRequestState::parse(&state), lifecycle_generation: generation as u64, current_revision: revision, review_status, current_review, + merge_results, + current_merge_result: None, + applied_merge_result: None, created_at, updated_at, merged_by_account_id, @@ -1161,8 +1540,8 @@ fn load_revision( revision_id: &str, ) -> Result { let mut revision: MergeRequestRevision = conn.query_row( - "SELECT revision_id,ordinal,base_commit,head_commit,head_tree,diff_digest,summary,assignment_id,created_at FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3", - params![workspace_id,mr_id,revision_id], |r| Ok(MergeRequestRevision { revision_id:r.get(0)?, ordinal:r.get::<_,i64>(1)? as u64, base_commit:r.get(2)?, head_commit:r.get(3)?, head_tree:r.get(4)?, diff_digest:r.get(5)?, changed_paths:Vec::new(), summary:r.get(6)?, assignment_id:r.get(7)?, created_at:r.get(8)? }), + "SELECT revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3", + params![workspace_id,mr_id,revision_id], |r| Ok(MergeRequestRevision { revision_id:r.get(0)?, ordinal:r.get::<_,i64>(1)? as u64, base_commit:r.get(2)?, head_commit:r.get(3)?, diff_digest:r.get(4)?, changed_paths:Vec::new(), summary:r.get(5)?, assignment_id:r.get(6)?, created_at:r.get(7)? }), ).map_err(db)?; let mut statement = conn.prepare("SELECT path FROM merge_request_revision_paths WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3 ORDER BY ordinal").map_err(db)?; revision.changed_paths = statement @@ -1179,7 +1558,7 @@ fn insert_revision( mr_id: &str, revision: &MergeRequestRevision, ) -> Result<()> { - conn.execute("INSERT INTO merge_request_revisions (workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,head_tree,diff_digest,summary,assignment_id,created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11)", params![workspace_id,mr_id,revision.revision_id,revision.ordinal as i64,revision.base_commit,revision.head_commit,revision.head_tree,revision.diff_digest,revision.summary,revision.assignment_id,revision.created_at]).map_err(db)?; + conn.execute("INSERT INTO merge_request_revisions (workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10)", params![workspace_id,mr_id,revision.revision_id,revision.ordinal as i64,revision.base_commit,revision.head_commit,revision.diff_digest,revision.summary,revision.assignment_id,revision.created_at]).map_err(db)?; for (ordinal, path) in revision.changed_paths.iter().enumerate() { conn.execute("INSERT INTO merge_request_revision_paths (workspace_id,merge_request_id,revision_id,ordinal,path) VALUES (?1,?2,?3,?4,?5)", params![workspace_id,mr_id,revision.revision_id,ordinal as i64,path]).map_err(db)?; } @@ -1193,7 +1572,7 @@ fn load_latest_review( revision_id: &str, generation: i64, ) -> Result> { - let attempt: Option = conn.query_row("SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_request_id=?2 AND r.revision_id=?3 AND a.lifecycle_generation=?4 ORDER BY r.submitted_at DESC, r.attempt_id DESC LIMIT 1", params![workspace_id,mr_id,revision_id,generation], |r| r.get(0)).optional().map_err(db)?; + let attempt: Option = conn.query_row("SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_request_id=?2 AND r.revision_id=?3 AND r.merge_result_id IS NULL AND a.lifecycle_generation=?4 ORDER BY r.submitted_at DESC, r.attempt_id DESC LIMIT 1", params![workspace_id,mr_id,revision_id,generation], |r| r.get(0)).optional().map_err(db)?; match attempt { Some(id) => load_review(conn, workspace_id, &id), None => Ok(None), @@ -1205,12 +1584,13 @@ fn load_review( workspace_id: &str, attempt_id: &str, ) -> Result> { - let row: Option<(String,String,String,String,String,String,String,String,String)> = conn.query_row( - "SELECT r.revision_id,r.decision,r.body,a.parent_assignment_id,a.parent_runtime_id,a.parent_worker_id,a.child_session_id,a.child_effective_profile,r.submitted_at FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.attempt_id=?2", - params![workspace_id,attempt_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?)), + let row: Option<(String,Option,String,String,String,String,String,String,String,String)> = conn.query_row( + "SELECT r.revision_id,r.merge_result_id,r.decision,r.body,a.parent_assignment_id,a.parent_runtime_id,a.parent_worker_id,a.child_session_id,a.child_effective_profile,r.submitted_at FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.attempt_id=?2", + params![workspace_id,attempt_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?)), ).optional().map_err(db)?; let Some(( revision_id, + merge_result_id, decision, body, assignment, @@ -1240,6 +1620,7 @@ fn load_review( Ok(Some(MergeRequestReview { attempt_id: attempt_id.into(), revision_id, + merge_result_id, decision: ReviewDecision::parse(&decision), body, findings, @@ -1272,6 +1653,15 @@ fn validate_current_implementation_assignment( Ok(()) } +fn validate_current_assignment_id( + conn: &Connection, + workspace_id: &str, + ticket_id: &str, + assignment_id: &str, +) -> Result<()> { + validate_current_implementation_assignment(conn, workspace_id, ticket_id, assignment_id) +} + fn validate_current_assignment( conn: &Connection, workspace_id: &str, @@ -1318,7 +1708,6 @@ fn validate_revision(revision: &MergeRequestRevision) -> Result<()> { ("revision_id", revision.revision_id.as_str()), ("base_commit", revision.base_commit.as_str()), ("head_commit", revision.head_commit.as_str()), - ("head_tree", revision.head_tree.as_str()), ("diff_digest", revision.diff_digest.as_str()), ("assignment_id", revision.assignment_id.as_str()), ] { @@ -1354,6 +1743,144 @@ fn validate_revision(revision: &MergeRequestRevision) -> Result<()> { Ok(()) } +fn merge_result_fingerprint(input: &RecordMergeResult) -> String { + token_hash(&format!( + "{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}", + input.ticket_id, + input.expected_revision_id, + input.target_commit, + input.source_commit, + input.result_commit, + input.strategy.as_str(), + input.resolution.as_str(), + input.actor_runtime_id, + input.actor_worker_id, + )) +} + +fn apply_target_observation(mr: &mut MergeRequest, observed_target_commit: Option<&str>) { + mr.observed_target_commit = observed_target_commit.map(str::to_owned); + for result in &mut mr.merge_results { + result.target_status = match observed_target_commit { + Some(commit) if result.target_commit == commit => MergeResultTargetStatus::Current, + Some(commit) if result.result_commit == commit => MergeResultTargetStatus::Applied, + Some(_) => MergeResultTargetStatus::Stale, + None => MergeResultTargetStatus::Unknown, + }; + } + mr.current_merge_result = observed_target_commit.and_then(|commit| { + mr.merge_results + .iter() + .rev() + .find(|result| result.target_commit == commit) + .cloned() + }); + mr.applied_merge_result = observed_target_commit.and_then(|commit| { + mr.merge_results + .iter() + .rev() + .find(|result| result.result_commit == commit) + .cloned() + }); +} + +fn load_merge_result( + conn: &Connection, + workspace_id: &str, + merge_result_id: &str, + generation: u64, +) -> Result> { + let row: Option<(String,String,String,String,String,String,String,String,String,String,String)> = conn.query_row( + "SELECT revision_id,target_commit,source_commit,result_commit,strategy,resolution,created_by_runtime_id,created_by_worker_id,created_at,operation_id,validated_at FROM merge_request_merge_results WHERE workspace_id=?1 AND merge_result_id=?2", + params![workspace_id,merge_result_id], + |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?,r.get(10)?)), + ).optional().map_err(db)?; + let Some(( + revision_id, + target_commit, + source_commit, + result_commit, + strategy, + resolution, + created_by_runtime_id, + created_by_worker_id, + created_at, + operation_id, + validated_at, + )) = row + else { + return Ok(None); + }; + let current_review = + load_latest_merge_result_review(conn, workspace_id, merge_result_id, generation)?; + let review_status = + current_review + .as_ref() + .map_or(ReviewStatus::Pending, |review| match review.decision { + ReviewDecision::Approve => ReviewStatus::Approved, + ReviewDecision::RequestChanges => ReviewStatus::ChangesRequested, + }); + Ok(Some(MergeResult { + merge_result_id: merge_result_id.into(), + revision_id, + target_commit, + source_commit, + result_commit, + strategy: MergeStrategy::parse(&strategy), + resolution: MergeResolution::parse(&resolution), + created_by_runtime_id, + created_by_worker_id, + created_at, + operation_id, + validated_at, + target_status: MergeResultTargetStatus::Unknown, + review_status, + current_review, + })) +} + +fn load_merge_results( + conn: &Connection, + workspace_id: &str, + merge_request_id: &str, + revision_id: &str, + generation: u64, +) -> Result> { + let mut statement = conn.prepare( + "SELECT merge_result_id FROM merge_request_merge_results WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3 ORDER BY created_at,merge_result_id", + ).map_err(db)?; + let ids = statement + .query_map( + params![workspace_id, merge_request_id, revision_id], + |row| row.get::<_, String>(0), + ) + .map_err(db)? + .collect::, _>>() + .map_err(db)?; + ids.into_iter() + .map(|id| { + load_merge_result(conn, workspace_id, &id, generation)? + .ok_or(MergeRequestError::MergeResultNotFound(id)) + }) + .collect() +} + +fn load_latest_merge_result_review( + conn: &Connection, + workspace_id: &str, + merge_result_id: &str, + generation: u64, +) -> Result> { + let attempt: Option = conn.query_row( + "SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_result_id=?2 AND a.lifecycle_generation=?3 ORDER BY r.submitted_at DESC,r.attempt_id DESC LIMIT 1", + params![workspace_id,merge_result_id,generation as i64], |row| row.get(0), + ).optional().map_err(db)?; + attempt + .map(|attempt| load_review(conn, workspace_id, &attempt)) + .transpose() + .map(|review| review.flatten()) +} + fn validate_review_input(input: &SubmitReview) -> Result<()> { if input.body.len() > MAX_REVIEW_BODY_BYTES { return Err(MergeRequestError::TooLarge { diff --git a/crates/merge-request/tests/store.rs b/crates/merge-request/tests/store.rs index 4cf3c9a4..8997879f 100644 --- a/crates/merge-request/tests/store.rs +++ b/crates/merge-request/tests/store.rs @@ -44,7 +44,6 @@ fn revision(id: &str, ordinal: u64, head: &str) -> MergeRequestRevision { ordinal, base_commit: "base".into(), head_commit: head.into(), - head_tree: format!("tree-{head}"), diff_digest: format!("sha256:diff-{head}"), changed_paths: vec!["src/lib.rs".into()], summary: format!("revision {id}"), @@ -58,6 +57,7 @@ fn open(store: &SqliteMergeRequestStore) { merge_request_id: "MR1".into(), ticket_id: "T1".into(), repository_id: "repo".into(), + target_ref_selector: "develop".into(), revision: revision("V1", 1, "h1"), authenticated_runtime_id: "R1".into(), authenticated_worker_id: "W1".into(), @@ -79,6 +79,7 @@ fn attempt(store: &SqliteMergeRequestStore, id: &str, revision: &str, token: &st attempt_id: id.into(), ticket_id: "T1".into(), revision_id: revision.into(), + merge_result_id: None, parent_assignment_id: "A1".into(), parent_runtime_id: "R1".into(), parent_worker_id: "W1".into(), @@ -97,6 +98,7 @@ fn review( store.submit_review(SubmitReview { ticket_id: "T1".into(), revision_id: revision.into(), + merge_result_id: None, capability_token: token.into(), decision, body: "evidence".into(), @@ -105,6 +107,34 @@ fn review( }) } +fn record_result( + store: &SqliteMergeRequestStore, + revision: &str, + target: &str, + source: &str, + result: &str, + strategy: MergeStrategy, + resolution: MergeResolution, + operation: &str, +) -> RecordMergeResultOutcome { + store + .record_merge_result(RecordMergeResult { + merge_result_id: format!("result-{operation}"), + ticket_id: "T1".into(), + expected_revision_id: revision.into(), + target_commit: target.into(), + source_commit: source.into(), + result_commit: result.into(), + strategy, + resolution, + operation_id: operation.into(), + actor_runtime_id: "runtime-orchestrator".into(), + actor_worker_id: "workspace-orchestrator".into(), + created_at: "2026-07-26T00:00:02Z".into(), + }) + .unwrap() +} + #[test] fn storage_allows_multiple_merge_requests_for_one_ticket() { let (_dir, store) = setup(); @@ -114,6 +144,7 @@ fn storage_allows_multiple_merge_requests_for_one_ticket() { merge_request_id: "MR2".into(), ticket_id: "T1".into(), repository_id: "repo".into(), + target_ref_selector: "develop".into(), revision: revision("V2", 1, "h2"), authenticated_runtime_id: "R1".into(), authenticated_worker_id: "W1".into(), @@ -133,6 +164,162 @@ fn storage_allows_multiple_merge_requests_for_one_ticket() { ); } +#[test] +fn merge_result_is_idempotent_target_fenced_and_non_ff_reviewed_independently() { + let (_tmp, store) = setup(); + open(&store); + attempt( + &store, + "attempt-source", + "V1", + "token-source", + "child-source", + ); + review(&store, "V1", "token-source", ReviewDecision::Approve).unwrap(); + + let first = record_result( + &store, + "V1", + "T0", + "h1", + "h1", + MergeStrategy::FastForward, + MergeResolution::None, + "op-ff", + ); + assert!(!first.replayed); + let replay = record_result( + &store, + "V1", + "T0", + "h1", + "h1", + MergeStrategy::FastForward, + MergeResolution::None, + "op-ff", + ); + assert!(replay.replayed); + assert_eq!( + first.merge_result.merge_result_id, + replay.merge_result.merge_result_id + ); + + let ready = store + .readiness_for_ticket_with_target("T1", Some("T0")) + .unwrap(); + assert!(ready.ready, "{:?}", ready.blockers); + assert_eq!(ready.merge_result_id.as_deref(), Some("result-op-ff")); + + let stale = store + .readiness_for_ticket_with_target("T1", Some("T1")) + .unwrap(); + assert!(!stale.ready); + assert!( + stale + .blockers + .iter() + .any(|blocker| blocker.contains("target moved")) + ); + + let changed = store.record_merge_result(RecordMergeResult { + merge_result_id: "result-conflict".into(), + ticket_id: "T1".into(), + expected_revision_id: "V1".into(), + target_commit: "different".into(), + source_commit: "h1".into(), + result_commit: "h1".into(), + strategy: MergeStrategy::FastForward, + resolution: MergeResolution::None, + operation_id: "op-ff".into(), + actor_runtime_id: "runtime-orchestrator".into(), + actor_worker_id: "workspace-orchestrator".into(), + created_at: "2026-07-26T00:00:03Z".into(), + }); + assert!(matches!( + changed, + Err(MergeRequestError::MergeResultOperationConflict) + )); + + record_result( + &store, + "V1", + "T1", + "h1", + "M1", + MergeStrategy::Merge, + MergeResolution::ConflictsResolved, + "op-merge", + ); + let pending = store + .readiness_for_ticket_with_target("T1", Some("T1")) + .unwrap(); + assert!(!pending.ready); + assert_eq!( + pending.merge_result_review_status, + Some(ReviewStatus::Pending) + ); + + store + .register_reviewer_child_session(RegisterReviewerChildSession { + parent_runtime_id: "runtime-orchestrator".into(), + parent_worker_id: "workspace-orchestrator".into(), + child_session_id: "child-merge".into(), + now: "2026-07-26T00:00:04Z".into(), + }) + .unwrap(); + store + .register_review_attempt(RegisterReviewAttempt { + attempt_id: "attempt-merge".into(), + ticket_id: "T1".into(), + revision_id: "V1".into(), + merge_result_id: Some("result-op-merge".into()), + parent_assignment_id: "A1".into(), + parent_runtime_id: "runtime-orchestrator".into(), + parent_worker_id: "workspace-orchestrator".into(), + child_session_id: "child-merge".into(), + capability_token: "token-merge".into(), + now: "2026-07-26T00:00:04Z".into(), + }) + .unwrap(); + store + .submit_review(SubmitReview { + ticket_id: "T1".into(), + revision_id: "V1".into(), + merge_result_id: Some("result-op-merge".into()), + capability_token: "token-merge".into(), + decision: ReviewDecision::Approve, + body: "merge evidence is valid".into(), + findings: vec![], + now: "2026-07-26T00:00:05Z".into(), + }) + .unwrap(); + let approved = store + .readiness_for_ticket_with_target("T1", Some("T1")) + .unwrap(); + assert!(approved.ready, "{:?}", approved.blockers); + assert_eq!( + approved.merge_result_review_status, + Some(ReviewStatus::Approved) + ); + let applied = store + .show_for_ticket_with_target("T1", Some("M1")) + .unwrap() + .unwrap(); + assert_eq!( + applied + .applied_merge_result + .as_ref() + .map(|result| result.target_status), + Some(MergeResultTargetStatus::Applied) + ); + assert!( + store + .readiness_for_ticket_with_target("T1", Some("M1")) + .unwrap() + .ready + ); +} + #[test] fn bounded_context_rejects_oversized_revision_evidence() { let (_dir, store) = setup(); @@ -142,6 +329,7 @@ fn bounded_context_rejects_oversized_revision_evidence() { merge_request_id: "MR1".into(), ticket_id: "T1".into(), repository_id: "repo".into(), + target_ref_selector: "develop".into(), revision: oversized, authenticated_runtime_id: "R1".into(), authenticated_worker_id: "W1".into(), @@ -311,7 +499,23 @@ fn v7_completion_operations_are_preserved_as_legacy_assigned_coder_authority() { |row| row.get(0), ) .unwrap(); - assert_eq!(version, 8); + assert_eq!(version, 9); + let head_tree_columns: i64 = conn + .query_row( + "SELECT COUNT(*) FROM pragma_table_info('merge_request_revisions') WHERE name='head_tree'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(head_tree_columns, 0); + let merge_result_table: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='merge_request_merge_results'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(merge_result_table, 1); } #[test] @@ -431,6 +635,7 @@ fn spoof_self_approval_replay_and_cross_workspace_are_rejected() { attempt_id: "bad".into(), ticket_id: "T1".into(), revision_id: "V1".into(), + merge_result_id: None, parent_assignment_id: "A1".into(), parent_runtime_id: "R1".into(), parent_worker_id: "W1".into(), diff --git a/crates/worker/src/feature/builtin/merge_request.rs b/crates/worker/src/feature/builtin/merge_request.rs index 104c4db5..6a9aa923 100644 --- a/crates/worker/src/feature/builtin/merge_request.rs +++ b/crates/worker/src/feature/builtin/merge_request.rs @@ -12,6 +12,7 @@ pub const MERGE_REQUEST_COMMON_TOOL_NAMES: &[&str] = &[ "MergeRequestReadinessCheck", "MergeRequestOpen", "MergeRequestAddRevision", + "MergeRequestRecordMergeResult", "MergeRequestComplete", ]; pub const MERGE_REQUEST_REVIEW_TOOL_NAME: &str = "MergeRequestReviewSubmit"; @@ -21,6 +22,7 @@ enum Kind { Readiness, Open, AddRevision, + RecordMergeResult, Complete, Review, } @@ -41,7 +43,6 @@ struct OpenInput { revision_id: String, base_commit: String, head_commit: String, - head_tree: String, diff_digest: String, #[serde(default)] changed_paths: Vec, @@ -55,7 +56,6 @@ struct AddRevisionInput { revision_id: String, base_commit: String, head_commit: String, - head_tree: String, diff_digest: String, #[serde(default)] changed_paths: Vec, @@ -63,6 +63,30 @@ struct AddRevisionInput { summary: String, } #[derive(Debug, Deserialize, JsonSchema)] +struct RecordMergeResultInput { + ticket: String, + expected_current_revision_id: String, + operation_id: String, + target_commit: String, + source_commit: String, + result_commit: String, + strategy: MergeStrategyInput, + resolution: MergeResolutionInput, +} +#[derive(Debug, Deserialize, JsonSchema)] +#[serde(rename_all = "snake_case")] +enum MergeStrategyInput { + FastForward, + Merge, +} +#[derive(Debug, Deserialize, JsonSchema)] +#[serde(rename_all = "snake_case")] +enum MergeResolutionInput { + None, + Clean, + ConflictsResolved, +} +#[derive(Debug, Deserialize, JsonSchema)] struct CompleteInput { ticket: String, operation_id: String, @@ -101,6 +125,7 @@ impl Kind { Self::Readiness => "MergeRequestReadinessCheck", Self::Open => "MergeRequestOpen", Self::AddRevision => "MergeRequestAddRevision", + Self::RecordMergeResult => "MergeRequestRecordMergeResult", Self::Complete => "MergeRequestComplete", Self::Review => "MergeRequestReviewSubmit", } @@ -113,6 +138,7 @@ impl Kind { Self::Show | Self::Readiness => json!(schemars::schema_for!(ShowInput)), Self::Open => json!(schemars::schema_for!(OpenInput)), Self::AddRevision => json!(schemars::schema_for!(AddRevisionInput)), + Self::RecordMergeResult => json!(schemars::schema_for!(RecordMergeResultInput)), Self::Complete => json!(schemars::schema_for!(CompleteInput)), Self::Review => json!(schemars::schema_for!(ReviewInput)), } @@ -158,7 +184,7 @@ impl Tool for MergeRequestTool { WorkspaceRequestMethod::Post, format!("/api/w/{workspace_id}/tickets/{}/merge-request", v.ticket), Some( - json!({"repository_id":v.repository_id,"revision_id":v.revision_id,"base_commit":v.base_commit,"head_commit":v.head_commit,"head_tree":v.head_tree,"diff_digest":v.diff_digest,"changed_paths":v.changed_paths,"summary":v.summary}), + json!({"repository_id":v.repository_id,"revision_id":v.revision_id,"base_commit":v.base_commit,"head_commit":v.head_commit,"diff_digest":v.diff_digest,"changed_paths":v.changed_paths,"summary":v.summary}), ), ) } @@ -172,7 +198,30 @@ impl Tool for MergeRequestTool { v.ticket ), Some( - json!({"expected_current_revision_id":v.expected_current_revision_id,"revision_id":v.revision_id,"base_commit":v.base_commit,"head_commit":v.head_commit,"head_tree":v.head_tree,"diff_digest":v.diff_digest,"changed_paths":v.changed_paths,"summary":v.summary}), + json!({"expected_current_revision_id":v.expected_current_revision_id,"revision_id":v.revision_id,"base_commit":v.base_commit,"head_commit":v.head_commit,"diff_digest":v.diff_digest,"changed_paths":v.changed_paths,"summary":v.summary}), + ), + ) + } + Kind::RecordMergeResult => { + let v: RecordMergeResultInput = parse(input)?; + nonempty(&v.ticket)?; + let strategy = match v.strategy { + MergeStrategyInput::FastForward => "fast_forward", + MergeStrategyInput::Merge => "merge", + }; + let resolution = match v.resolution { + MergeResolutionInput::None => "none", + MergeResolutionInput::Clean => "clean", + MergeResolutionInput::ConflictsResolved => "conflicts_resolved", + }; + ( + WorkspaceRequestMethod::Post, + format!( + "/api/w/{workspace_id}/tickets/{}/merge-request/merge-results", + v.ticket + ), + Some( + json!({"expected_current_revision_id":v.expected_current_revision_id,"operation_id":v.operation_id,"target_commit":v.target_commit,"source_commit":v.source_commit,"result_commit":v.result_commit,"strategy":strategy,"resolution":resolution}), ), ) } @@ -261,6 +310,7 @@ pub fn common_tools(client: Arc) -> Vec { definition(client.clone(), Kind::Readiness), definition(client.clone(), Kind::Open), definition(client.clone(), Kind::AddRevision), + definition(client.clone(), Kind::RecordMergeResult), definition(client, Kind::Complete), ] } @@ -288,6 +338,9 @@ pub fn description(name: &str) -> Option<&'static str> { "MergeRequestAddRevision" => { Some("Append an immutable revision; prior approval cannot carry to the new revision.") } + "MergeRequestRecordMergeResult" => Some( + "Record validated immutable integration evidence for the current source revision and target tip.", + ), "MergeRequestComplete" => { Some("CAS-complete an approved revision with operation-id replay and crash fencing.") } @@ -297,3 +350,20 @@ pub fn description(name: &str) -> Option<&'static str> { _ => None, } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn merge_request_tool_contract_omits_tree_hashes_and_exposes_merge_result() { + let open = serde_json::to_string(&schemars::schema_for!(OpenInput)).unwrap(); + let add = serde_json::to_string(&schemars::schema_for!(AddRevisionInput)).unwrap(); + let result = serde_json::to_string(&schemars::schema_for!(RecordMergeResultInput)).unwrap(); + assert!(!open.contains("head_tree")); + assert!(!add.contains("head_tree")); + assert!(result.contains("fast_forward")); + assert!(result.contains("conflicts_resolved")); + assert!(MERGE_REQUEST_COMMON_TOOL_NAMES.contains(&"MergeRequestRecordMergeResult")); + } +} diff --git a/crates/worker/src/spawn/tool.rs b/crates/worker/src/spawn/tool.rs index 5f29ef5f..c1fb6379 100644 --- a/crates/worker/src/spawn/tool.rs +++ b/crates/worker/src/spawn/tool.rs @@ -68,6 +68,9 @@ struct SubWorkerSpawnInput { struct ReviewerHandoffInput { ticket_id: String, revision_id: String, + /// Optional immutable MergeResult subject. Omit for source-revision review. + #[serde(default)] + merge_result_id: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -422,6 +425,7 @@ impl Tool for SubWorkerSpawnTool { ( review.ticket_id.clone(), review.revision_id.clone(), + review.merge_result_id.clone(), uuid::Uuid::now_v7().to_string(), format!( "{}{}", @@ -431,7 +435,9 @@ impl Tool for SubWorkerSpawnTool { ) }); let child_workspace_context = - if let Some((ticket_id, revision_id, _, capability_token)) = &reviewer_attempt { + if let Some((ticket_id, revision_id, merge_result_id, _, capability_token)) = + &reviewer_attempt + { let workspace_id = self.workspace_context .workspace_id() @@ -453,6 +459,7 @@ impl Tool for SubWorkerSpawnTool { ReviewerAttemptContext { ticket_id: ticket_id.clone(), revision_id: revision_id.clone(), + merge_result_id: merge_result_id.clone(), }, capability_token.clone(), )); @@ -547,7 +554,9 @@ impl Tool for SubWorkerSpawnTool { } }; - if let Some((ticket_id, revision_id, attempt_id, capability_token)) = &reviewer_attempt { + if let Some((ticket_id, revision_id, merge_result_id, attempt_id, capability_token)) = + &reviewer_attempt + { let workspace_id = self.workspace_context.workspace_id().ok_or_else(|| { ToolError::ExecutionFailed("reviewer attempt lost Workspace identity".to_string()) })?; @@ -579,6 +588,7 @@ impl Tool for SubWorkerSpawnTool { let body = serde_json::json!({ "attempt_id": attempt_id, "revision_id": revision_id, + "merge_result_id": merge_result_id, "child_session_id": child_session_id, "capability_token": capability_token, }); diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index a9b5bd6d..bda08079 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -249,6 +249,7 @@ pub trait WorkspaceClient: std::fmt::Debug + Send + Sync { pub struct ReviewerAttemptContext { pub ticket_id: String, pub revision_id: String, + pub merge_result_id: Option, } #[derive(Debug)] @@ -310,6 +311,14 @@ impl WorkspaceClient for ReviewerChildWorkspaceClient { "revision_id".to_string(), serde_json::Value::String(self.context.revision_id.clone()), ); + object.insert( + "merge_result_id".to_string(), + self.context + .merge_result_id + .clone() + .map(serde_json::Value::String) + .unwrap_or(serde_json::Value::Null), + ); object.insert( "capability_token".to_string(), serde_json::Value::String(self.capability_token.clone()), @@ -453,6 +462,7 @@ mod reviewer_client_tests { ReviewerAttemptContext { ticket_id: "T1".into(), revision_id: "V1".into(), + merge_result_id: None, }, "secret".into(), ); diff --git a/crates/workspace-server/src/repositories.rs b/crates/workspace-server/src/repositories.rs index 9dfaea4c..a55ad387 100644 --- a/crates/workspace-server/src/repositories.rs +++ b/crates/workspace-server/src/repositories.rs @@ -83,10 +83,27 @@ pub struct GitCommitSummary { pub refs: Vec, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MergeTargetObservation { + pub selector: RepositorySelector, + pub commit: String, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct CommitObservation { + pub commit: String, + pub parents: Vec, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub enum RepositoryLookupError { UnknownRepository { id: RepositoryId }, UnsupportedProvider { id: RepositoryId, provider: String }, + MissingDefaultSelector { id: RepositoryId }, + InvalidSelector { id: RepositoryId, selector: String }, + CommitNotFound { id: RepositoryId, commit: String }, + InvalidCommitRelation { id: RepositoryId, detail: String }, + ProviderFailure { id: RepositoryId, operation: String }, } #[derive(Debug, Clone)] @@ -167,6 +184,121 @@ impl RepositoryRegistryReader { }) } + pub fn observe_merge_target( + &self, + id: &str, + requested_selector: Option<&str>, + ) -> Result { + let repository = self.merge_repository(id)?; + let selector = requested_selector + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(str::to_owned) + .or_else(|| repository.default_selector.clone()) + .ok_or_else(|| RepositoryLookupError::MissingDefaultSelector { id: id.to_string() })?; + if selector.starts_with('-') || selector.as_bytes().contains(&0) { + return Err(RepositoryLookupError::InvalidSelector { + id: id.to_string(), + selector, + }); + } + let spec = format!("{selector}^{{commit}}"); + let commit = merge_git_stdout( + repository, + "resolve target", + &["rev-parse", "--verify", "--end-of-options", &spec], + )? + .lines() + .next() + .unwrap_or_default() + .trim() + .to_owned(); + if commit.is_empty() { + return Err(RepositoryLookupError::InvalidSelector { + id: id.to_string(), + selector, + }); + } + Ok(MergeTargetObservation { selector, commit }) + } + + pub fn observe_commit( + &self, + id: &str, + commit: &str, + ) -> Result { + let repository = self.merge_repository(id)?; + let commit = commit.trim(); + if commit.is_empty() || commit.starts_with('-') { + return Err(RepositoryLookupError::CommitNotFound { + id: id.to_string(), + commit: commit.into(), + }); + } + let line = merge_git_stdout( + repository, + "read commit", + &[ + "show", + "--no-patch", + "--format=%H %P", + "--end-of-options", + commit, + ], + )?; + let mut parts = line.split_whitespace(); + let canonical = parts.next().unwrap_or_default().to_owned(); + if canonical.is_empty() { + return Err(RepositoryLookupError::CommitNotFound { + id: id.to_string(), + commit: commit.into(), + }); + } + Ok(CommitObservation { + commit: canonical, + parents: parts.map(str::to_owned).collect(), + }) + } + + pub fn ensure_ancestor( + &self, + id: &str, + ancestor: &str, + descendant: &str, + ) -> Result<(), RepositoryLookupError> { + let repository = self.merge_repository(id)?; + let status = Command::new("git") + .arg("-C") + .arg(&repository.path) + .args(["merge-base", "--is-ancestor", ancestor, descendant]) + .status() + .map_err(|_| RepositoryLookupError::ProviderFailure { + id: id.into(), + operation: "check commit ancestry".into(), + })?; + if status.success() { + Ok(()) + } else { + Err(RepositoryLookupError::InvalidCommitRelation { + id: id.into(), + detail: format!("commit {ancestor} is not an ancestor of {descendant}"), + }) + } + } + + fn merge_repository(&self, id: &str) -> Result<&ConfiguredRepository, RepositoryLookupError> { + let repository = self + .find(id) + .ok_or_else(|| RepositoryLookupError::UnknownRepository { id: id.into() })?; + if repository.provider != "git" { + return Err(RepositoryLookupError::UnsupportedProvider { + id: id.into(), + provider: repository.provider.clone(), + }); + } + Ok(repository) + } + fn find(&self, id: &str) -> Option<&ConfiguredRepository> { self.repositories .iter() @@ -256,6 +388,19 @@ impl RepositoryRegistryReader { } } +fn merge_git_stdout( + repository: &ConfiguredRepository, + operation: &str, + args: &[&str], +) -> Result { + git_stdout(&repository.path, args.iter().copied()).map_err(|_| { + RepositoryLookupError::ProviderFailure { + id: repository.id.clone(), + operation: operation.into(), + } + }) +} + fn git_stdout<'a, I>(repository_path: &PathBuf, args: I) -> Result where I: IntoIterator, @@ -422,6 +567,101 @@ mod tests { assert_eq!(projection.diagnostics[0].code, "repository_config_empty"); } + #[test] + fn merge_evidence_is_resolved_by_repository_identity() { + let temp = tempfile::tempdir().unwrap(); + let path = temp.path(); + assert!( + Command::new("git") + .args(["init", "-b", "main"]) + .arg(path) + .status() + .unwrap() + .success() + ); + for args in [ + vec!["config", "user.email", "test@example.com"], + vec!["config", "user.name", "Test"], + ] { + assert!( + Command::new("git") + .arg("-C") + .arg(path) + .args(args) + .status() + .unwrap() + .success() + ); + } + std::fs::write(path.join("file.txt"), "base\n").unwrap(); + assert!( + Command::new("git") + .arg("-C") + .arg(path) + .args(["add", "file.txt"]) + .status() + .unwrap() + .success() + ); + assert!( + Command::new("git") + .arg("-C") + .arg(path) + .args(["commit", "-m", "base"]) + .status() + .unwrap() + .success() + ); + let base = git_stdout(&path.to_path_buf(), ["rev-parse", "HEAD"]) + .unwrap() + .trim() + .to_string(); + assert!( + Command::new("git") + .arg("-C") + .arg(path) + .args(["checkout", "-b", "feature"]) + .status() + .unwrap() + .success() + ); + std::fs::write(path.join("file.txt"), "base\nfeature\n").unwrap(); + assert!( + Command::new("git") + .arg("-C") + .arg(path) + .args(["commit", "-am", "feature"]) + .status() + .unwrap() + .success() + ); + let source = git_stdout(&path.to_path_buf(), ["rev-parse", "HEAD"]) + .unwrap() + .trim() + .to_string(); + + let reader = RepositoryRegistryReader::new(vec![ConfiguredRepository { + id: "main".into(), + display_name: Some("Main".into()), + provider: "git".into(), + path: path.to_path_buf(), + uri: path.display().to_string(), + default_selector: Some("main".into()), + }]); + let target = reader.observe_merge_target("main", None).unwrap(); + assert_eq!(target.selector, "main"); + assert_eq!(target.commit, base); + assert_eq!( + reader.observe_commit("main", &source).unwrap().parents, + vec![base.clone()] + ); + reader.ensure_ancestor("main", &base, &source).unwrap(); + assert!(matches!( + reader.ensure_ancestor("main", &source, &base), + Err(RepositoryLookupError::InvalidCommitRelation { .. }) + )); + } + #[test] fn unknown_repository_is_not_resolved_from_fallback() { let reader = RepositoryRegistryReader::new(Vec::new()); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 997e8ddb..8a6ea694 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -94,8 +94,8 @@ use crate::observation::{ use crate::profile_settings::UpdateWorkspaceMetadataRequest; use crate::records::{ObjectiveDetail, ProjectRecordList, TicketDetail}; use crate::repositories::{ - ConfiguredRepository, RepositoryListProjection, RepositoryLogRead, RepositoryLookupError, - RepositoryRegistryReader, RepositorySummary, + ConfiguredRepository, MergeTargetObservation, RepositoryListProjection, RepositoryLogRead, + RepositoryLookupError, RepositoryRegistryReader, RepositorySummary, }; use crate::resource_broker::BackendResourceBroker; use crate::runtime_subscription::RuntimeSubscriptionBroker; @@ -1290,6 +1290,10 @@ pub fn build_router(api: WorkspaceApi) -> Router { "/api/w/{workspace_id}/tickets/{id}/merge-request/revisions", post(scoped_add_merge_request_revision), ) + .route( + "/api/w/{workspace_id}/tickets/{id}/merge-request/merge-results", + post(scoped_record_merge_request_result), + ) .route( "/api/w/{workspace_id}/internal/reviewer-child-sessions", post(scoped_register_reviewer_child_session), @@ -3521,7 +3525,6 @@ struct OpenMergeRequestRequest { revision_id: String, base_commit: String, head_commit: String, - head_tree: String, diff_digest: String, #[serde(default)] changed_paths: Vec, @@ -3535,7 +3538,6 @@ struct AddMergeRequestRevisionRequest { revision_id: String, base_commit: String, head_commit: String, - head_tree: String, diff_digest: String, #[serde(default)] changed_paths: Vec, @@ -3543,6 +3545,17 @@ struct AddMergeRequestRevisionRequest { summary: String, } +#[derive(Debug, serde::Deserialize)] +struct RecordMergeResultRequest { + expected_current_revision_id: String, + operation_id: String, + target_commit: String, + source_commit: String, + result_commit: String, + strategy: merge_request::MergeStrategy, + resolution: merge_request::MergeResolution, +} + #[derive(Debug, serde::Deserialize)] struct RegisterReviewerChildSessionRequest { child_session_id: String, @@ -3552,6 +3565,8 @@ struct RegisterReviewerChildSessionRequest { struct RegisterMergeRequestReviewAttemptRequest { attempt_id: String, revision_id: String, + #[serde(default)] + merge_result_id: Option, child_session_id: String, capability_token: String, } @@ -3559,6 +3574,8 @@ struct RegisterMergeRequestReviewAttemptRequest { #[derive(Debug, serde::Deserialize)] struct SubmitMergeRequestReviewRequest { revision_id: String, + #[serde(default)] + merge_result_id: Option, capability_token: String, decision: merge_request::ReviewDecision, #[serde(default)] @@ -3612,14 +3629,89 @@ fn merge_request_store( .map_err(Into::into) } +fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError { + Error::InvalidInput(format!( + "repository merge evidence validation failed: {error:?}" + )) + .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)) +} + +fn observe_merge_request_target( + api: &WorkspaceApi, + mr: &merge_request::MergeRequest, +) -> Option { + let selector = mr.target_ref_selector.as_deref()?; + api.repository_reader() + .observe_merge_target(&mr.repository_id, Some(selector)) + .ok() + .map(|target| target.commit) +} + async fn scoped_show_merge_request( State(api): State, AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>, ) -> ApiResult> { let workspace_id = parse_workspace_id(&workspace_id)?; let store = merge_request_store(&api, &workspace_id)?; + let current = store.show_for_ticket(&ticket_id)?.ok_or_else(|| { + Error::from(merge_request::MergeRequestError::NotFound( + ticket_id.clone(), + )) + })?; + let target_commit = observe_merge_request_target(&api, ¤t); let value = store - .show_for_ticket(&ticket_id)? + .show_for_ticket_with_target(&ticket_id, target_commit.as_deref())? .ok_or_else(|| Error::from(merge_request::MergeRequestError::NotFound(ticket_id)))?; Ok(Json(value)) } @@ -3629,9 +3721,17 @@ async fn scoped_merge_request_readiness( AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>, ) -> ApiResult> { let workspace_id = parse_workspace_id(&workspace_id)?; - Ok(Json( - merge_request_store(&api, &workspace_id)?.readiness_for_ticket(&ticket_id)?, - )) + let store = merge_request_store(&api, &workspace_id)?; + let current = store.show_for_ticket(&ticket_id)?.ok_or_else(|| { + Error::from(merge_request::MergeRequestError::NotFound( + ticket_id.clone(), + )) + })?; + let target_commit = observe_merge_request_target(&api, ¤t); + Ok(Json(store.readiness_for_ticket_with_target( + &ticket_id, + target_commit.as_deref(), + )?)) } async fn scoped_open_merge_request( @@ -3657,13 +3757,19 @@ async fn scoped_open_merge_request( ) .into()); } + let (base_commit, source_commit, target) = validate_open_merge_request_evidence( + &api, + &ticket_id, + &input.repository_id, + &input.base_commit, + &input.head_commit, + )?; let now = Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true); let revision = merge_request::MergeRequestRevision { revision_id: input.revision_id, ordinal: 1, - base_commit: input.base_commit, - head_commit: input.head_commit, - head_tree: input.head_tree, + base_commit, + head_commit: source_commit, diff_digest: input.diff_digest, changed_paths: input.changed_paths, summary: input.summary, @@ -3675,6 +3781,7 @@ async fn scoped_open_merge_request( merge_request_id: format!("mr_{}", Uuid::now_v7().simple()), ticket_id, repository_id: input.repository_id, + target_ref_selector: target.selector, revision, authenticated_runtime_id: source.runtime_id, authenticated_worker_id: source.worker_id, @@ -3714,6 +3821,12 @@ async fn scoped_add_merge_request_revision( ticket_id.clone(), )) })?; + let (base_commit, head_commit) = validate_revision_evidence( + &api, + ¤t.repository_id, + &input.base_commit, + &input.head_commit, + )?; let now = Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true); let mr = merge_request_store(&api, &workspace_id)?.add_revision(merge_request::AddRevision { @@ -3722,9 +3835,8 @@ async fn scoped_add_merge_request_revision( revision: merge_request::MergeRequestRevision { revision_id: input.revision_id, ordinal: current.current_revision.ordinal + 1, - base_commit: input.base_commit, - head_commit: input.head_commit, - head_tree: input.head_tree, + base_commit, + head_commit, diff_digest: input.diff_digest, changed_paths: input.changed_paths, summary: input.summary, @@ -3738,6 +3850,105 @@ async fn scoped_add_merge_request_revision( Ok(Json(mr)) } +async fn scoped_record_merge_request_result( + State(api): State, + AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(input): Json, +) -> ApiResult> { + require_workspace_access(&workspace_id, &api)?; + let source = authenticate_worker_mutation_source(&api, &workspace_id, &headers)?; + require_online_workspace_orchestrator_source(&api, &source)?; + let store = merge_request_store(&api, &workspace_id)?; + let mr = store.show_for_ticket(&ticket_id)?.ok_or_else(|| { + Error::from(merge_request::MergeRequestError::NotFound( + ticket_id.clone(), + )) + })?; + if mr.current_revision.revision_id != input.expected_current_revision_id { + return Err( + Error::from(merge_request::MergeRequestError::StaleRevision { + expected: input.expected_current_revision_id, + current: mr.current_revision.revision_id, + }) + .into(), + ); + } + let target = api + .repository_reader() + .observe_merge_target(&mr.repository_id, mr.target_ref_selector.as_deref()) + .map_err(repository_merge_evidence_error)?; + let supplied_target = api + .repository_reader() + .observe_commit(&mr.repository_id, &input.target_commit) + .map_err(repository_merge_evidence_error)?; + let supplied_source = api + .repository_reader() + .observe_commit(&mr.repository_id, &input.source_commit) + .map_err(repository_merge_evidence_error)?; + let supplied_result = api + .repository_reader() + .observe_commit(&mr.repository_id, &input.result_commit) + .map_err(repository_merge_evidence_error)?; + if target.commit != supplied_target.commit { + return Err(Error::InvalidInput( + "MergeResult target_commit is not the current target tip".into(), + ) + .into()); + } + if mr.current_revision.head_commit != supplied_source.commit { + return Err(Error::InvalidInput( + "MergeResult source_commit is not the current source revision".into(), + ) + .into()); + } + match input.strategy { + merge_request::MergeStrategy::FastForward => { + if input.resolution != merge_request::MergeResolution::None + || supplied_result.commit != supplied_source.commit + { + return Err(Error::InvalidInput( + "fast-forward MergeResult must use the source commit and resolution=none" + .into(), + ) + .into()); + } + api.repository_reader() + .ensure_ancestor(&mr.repository_id, &target.commit, &supplied_source.commit) + .map_err(repository_merge_evidence_error)?; + } + merge_request::MergeStrategy::Merge => { + if input.resolution == merge_request::MergeResolution::None { + return Err(Error::InvalidInput( + "merge MergeResult requires clean or conflicts_resolved resolution".into(), + ) + .into()); + } + if supplied_result.parents.len() != 2 + || !supplied_result.parents.contains(&target.commit) + || !supplied_result.parents.contains(&supplied_source.commit) + { + return Err(Error::InvalidInput("merge result commit must have exactly the target and source commits as parents".into()).into()); + } + } + } + let outcome = store.record_merge_result(merge_request::RecordMergeResult { + merge_result_id: format!("MRG-{}", Uuid::new_v4()), + ticket_id, + expected_revision_id: mr.current_revision.revision_id, + target_commit: target.commit, + source_commit: supplied_source.commit, + result_commit: supplied_result.commit, + strategy: input.strategy, + resolution: input.resolution, + operation_id: input.operation_id, + actor_runtime_id: source.runtime_id, + actor_worker_id: source.worker_id, + created_at: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), + })?; + Ok(Json(outcome)) +} + async fn scoped_register_reviewer_child_session( State(api): State, headers: HeaderMap, @@ -3773,7 +3984,9 @@ async fn scoped_register_merge_request_review_attempt( .ok_or_else(|| { Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; - if assignment.worker.runtime_id != source.runtime_id + if input.merge_result_id.is_some() { + require_online_workspace_orchestrator_source(&api, &source)?; + } else if assignment.worker.runtime_id != source.runtime_id || assignment.worker.worker_id != source.worker_id { return Err(Error::TicketAssignmentConflict( @@ -3786,6 +3999,7 @@ async fn scoped_register_merge_request_review_attempt( attempt_id: input.attempt_id, ticket_id, revision_id: input.revision_id, + merge_result_id: input.merge_result_id, parent_assignment_id: assignment.assignment_id, parent_runtime_id: source.runtime_id, parent_worker_id: source.worker_id, @@ -3807,6 +4021,7 @@ async fn scoped_submit_merge_request_review( merge_request_store(&api, &workspace_id)?.submit_review(merge_request::SubmitReview { ticket_id, revision_id: input.revision_id, + merge_result_id: input.merge_result_id, capability_token: input.capability_token, decision: input.decision, body: input.body, @@ -3884,16 +4099,41 @@ async fn scoped_confirm_merge_request( require_workspace_access(&workspace_id, &api)?; reject_non_browser_merge_auth(&headers)?; let actor = require_actor(&api, &headers).await?; - let mr = merge_request_store(&api, &workspace_id)?.confirm_merge( - merge_request::MergeConfirmation { - ticket_id, - expected_revision_id: input.expected_revision_id, - authenticated_account_id: actor.account_id, - actor_kind: "user".to_string(), - explicit_confirmation: input.explicit_confirmation, - now: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), - }, - )?; + let store = merge_request_store(&api, &workspace_id)?; + let current = store.show_for_ticket(&ticket_id)?.ok_or_else(|| { + Error::from(merge_request::MergeRequestError::NotFound( + ticket_id.clone(), + )) + })?; + let target_commit = observe_merge_request_target(&api, ¤t); + let observed = store + .show_for_ticket_with_target(&ticket_id, target_commit.as_deref())? + .ok_or_else(|| { + Error::from(merge_request::MergeRequestError::NotFound( + ticket_id.clone(), + )) + })?; + if observed.applied_merge_result.is_none() { + return Err(Error::InvalidInput( + "Merge Request result is not the current target tip; target update is a separate prerequisite operation".into(), + ).into()); + } + let readiness = store.readiness_for_ticket_with_target(&ticket_id, target_commit.as_deref())?; + if !readiness.ready { + return Err(Error::InvalidInput(format!( + "Merge Request is not integration-ready: {}", + readiness.blockers.join("; ") + )) + .into()); + } + let mr = store.confirm_merge(merge_request::MergeConfirmation { + ticket_id, + expected_revision_id: input.expected_revision_id, + authenticated_account_id: actor.account_id, + actor_kind: "user".to_string(), + explicit_confirmation: input.explicit_confirmation, + now: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), + })?; Ok(Json(mr)) } @@ -11128,6 +11368,21 @@ fn repository_lookup(result: std::result::Result) - }], ) } + other => { + let message = format!("repository evidence validation failed: {other:?}"); + ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: "workspace-repository-registry".to_string(), + code: "repository_evidence_invalid".to_string(), + message: message.clone(), + }, + vec![RuntimeDiagnostic { + code: "repository_evidence_invalid".to_string(), + severity: DiagnosticSeverity::Error, + message, + }], + ) + } }) } @@ -12193,12 +12448,13 @@ mod tests { merge_request_id: "MR-server-completion".into(), ticket_id: ticket.id.clone(), repository_id: TEST_REPOSITORY_ID.into(), + target_ref_selector: "develop".into(), revision: merge_request::MergeRequestRevision { revision_id: "V1".into(), ordinal: 1, base_commit: "base".into(), head_commit: "head".into(), - head_tree: "tree".into(), + diff_digest: "sha256:diff".into(), changed_paths: vec!["src/lib.rs".into()], summary: "approved revision".into(), @@ -12223,6 +12479,7 @@ mod tests { attempt_id: "attempt".into(), ticket_id: ticket.id.clone(), revision_id: "V1".into(), + merge_result_id: None, parent_assignment_id: assignment.assignment_id.clone(), parent_runtime_id: coder.worker_ref.runtime_id.clone(), parent_worker_id: coder.worker_ref.worker_id.clone(), @@ -12235,6 +12492,7 @@ mod tests { .submit_review(merge_request::SubmitReview { ticket_id: ticket.id.clone(), revision_id: "V1".into(), + merge_result_id: None, capability_token: "review-token".into(), decision: merge_request::ReviewDecision::Approve, body: "approved".into(), diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte index 0f808103..06b469f4 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte @@ -20,8 +20,31 @@ type MergeRequestDetail = { state: "draft" | "open" | "closed" | "merged"; review_status: "pending" | "approved" | "changes_requested"; - current_revision: { revision_id: string; head_commit: string; head_tree: string; diff_digest: string; changed_paths: string[]; summary: string }; + target_ref_selector?: string | null; + target_status: "known" | "unknown"; + observed_target_commit?: string | null; + current_revision: { revision_id: string; head_commit: string; diff_digest: string; changed_paths: string[]; summary: string }; current_review?: { decision: string; body: string; reviewer_effective_profile: string } | null; + current_merge_result?: { + merge_result_id: string; + target_commit: string; + source_commit: string; + result_commit: string; + strategy: "fast_forward" | "merge"; + resolution: "none" | "clean" | "conflicts_resolved"; + target_status: "current" | "applied" | "stale" | "unknown"; + review_status: "pending" | "approved" | "changes_requested"; + } | null; + applied_merge_result?: { + merge_result_id: string; + target_commit: string; + source_commit: string; + result_commit: string; + strategy: "fast_forward" | "merge"; + resolution: "none" | "clean" | "conflicts_resolved"; + target_status: "current" | "applied" | "stale" | "unknown"; + review_status: "pending" | "approved" | "changes_requested"; + } | null; merged_at?: string | null; }; @@ -368,15 +391,38 @@

{data.mergeRequest.error}

{:else if mergeRequest}

{mergeRequest.state} · {mergeRequest.review_status}

+

Target {mergeRequest.target_ref_selector ?? "unknown"} · {mergeRequest.target_status}

+ {#if mergeRequest.observed_target_commit}

Target tip {mergeRequest.observed_target_commit}

{/if}

{mergeRequest.current_revision.revision_id}

Head {mergeRequest.current_revision.head_commit}

+ {#if mergeRequest.current_merge_result} +

+ MergeResult {mergeRequest.current_merge_result.merge_result_id} · + {mergeRequest.current_merge_result.strategy} / {mergeRequest.current_merge_result.resolution} · + {mergeRequest.current_merge_result.target_status} / {mergeRequest.current_merge_result.review_status} +

+

Result {mergeRequest.current_merge_result.result_commit}

+ {:else if mergeRequest.applied_merge_result} +

+ Applied MergeResult {mergeRequest.applied_merge_result.merge_result_id} · + {mergeRequest.applied_merge_result.strategy} / {mergeRequest.applied_merge_result.resolution} · + {mergeRequest.applied_merge_result.review_status} +

+

Result {mergeRequest.applied_merge_result.result_commit} is the current target tip.

+ {:else if mergeRequest.state === "open"} +

No current MergeResult has been recorded for this target tip.

+ {/if} {#if mergeRequest.current_revision.summary}

{mergeRequest.current_revision.summary}

{/if} {#if mergeRequest.current_review}

{mergeRequest.current_review.decision} by {mergeRequest.current_review.reviewer_effective_profile}

{#if mergeRequest.current_review.body}{/if} {/if} - {#if mergeRequest.state === "open" && mergeRequest.review_status === "approved"} - + {#if mergeRequest.state === "open" + && mergeRequest.review_status === "approved" + && mergeRequest.applied_merge_result?.target_status === "applied" + && (mergeRequest.applied_merge_result.strategy === "fast_forward" + || mergeRequest.applied_merge_result.review_status === "approved")} + {/if} {:else}