diff --git a/crates/merge-request/src/lib.rs b/crates/merge-request/src/lib.rs index 8f987a19..4f98e85d 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 = 10; +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,18 +51,8 @@ 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 is not the current final integration candidate")] - MergeResultNotFinal, - #[error("current final merge result is missing")] - FinalMergeResultMissing, - #[error("current target does not equal the final merge result commit")] - FinalMergeResultNotApplied, - #[error("merge result evidence is invalid: {0}")] - InvalidMergeResult(String), + #[error("completion merge outcome is invalid: {0}")] + InvalidMergeOutcome(String), #[error("Merge Request target is unknown and must be resolved explicitly")] UnknownTarget, #[error("Ticket must be inprogress before Merge Request completion (current: {0})")] @@ -194,15 +184,6 @@ impl MergeResolution { } } -#[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, @@ -229,8 +210,6 @@ 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, @@ -242,25 +221,6 @@ 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, @@ -277,14 +237,22 @@ pub struct MergeRequest { pub current_revision: MergeRequestRevision, pub review_status: ReviewStatus, pub current_review: Option, - pub merge_results: Vec, - /// The one explicitly selected final integration candidate. Historical - /// candidates remain in `merge_results` but never compete with this pointer. - #[serde(skip_serializing_if = "Option::is_none")] - pub final_merge_result: Option, pub created_at: String, pub updated_at: String, - pub merged_by_account_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merged_revision_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merged_target_commit: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merged_result_commit: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merge_strategy: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merge_resolution: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merged_by_runtime_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub merged_by_worker_id: Option, pub merged_at: Option, } @@ -323,7 +291,6 @@ 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, @@ -337,7 +304,6 @@ 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, @@ -346,8 +312,8 @@ pub struct SubmitReview { } #[derive(Debug, Clone, PartialEq, Eq)] -pub struct RecordMergeResult { - pub merge_result_id: String, +pub struct CompleteMergeRequest { + pub operation_id: String, pub ticket_id: String, pub expected_revision_id: String, pub target_commit: String, @@ -355,25 +321,6 @@ pub struct RecordMergeResult { 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, - pub ticket_id: String, - pub expected_revision_id: String, - pub expected_merge_result_id: String, - pub observed_target_commit: String, pub implementation_assignment_id: String, pub completion_actor_runtime_id: String, pub completion_actor_worker_id: String, @@ -398,8 +345,6 @@ pub struct MergeRequestReadiness { 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, } @@ -510,26 +455,6 @@ impl SqliteMergeRequestStore { 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 final_result = mr.final_merge_result.as_ref(); - match final_result { - None if observed_target_commit.is_some() => { - blockers.push("current source revision has no final validated MergeResult".into()) - } - Some(result) if result.target_status == MergeResultTargetStatus::Stale => { - blockers.push("target moved after the final MergeResult was recorded".into()) - } - Some(result) if result.target_status == MergeResultTargetStatus::Unknown => { - blockers.push("final MergeResult target state could not be resolved".into()) - } - Some(result) - if matches!(result.strategy, MergeStrategy::Merge) - && result.review_status != ReviewStatus::Approved => - { - blockers - .push("non-fast-forward final MergeResult is not independently approved".into()) - } - _ => {} - } Ok(MergeRequestReadiness { ticket_id: ticket_id.to_string(), merge_request_id: mr.merge_request_id, @@ -538,8 +463,6 @@ impl SqliteMergeRequestStore { observed_target_commit: mr.observed_target_commit, ready: blockers.is_empty(), review_status: mr.review_status, - merge_result_id: final_result.map(|result| result.merge_result_id.clone()), - merge_result_review_status: final_result.map(|result| result.review_status), blockers, }) } @@ -617,10 +540,6 @@ impl SqliteMergeRequestStore { "UPDATE merge_requests SET current_revision_id=?3, updated_at=?4 WHERE workspace_id=?1 AND merge_request_id=?2 AND current_revision_id=?5", params![self.workspace_id, current.merge_request_id, input.revision.revision_id, input.now, input.expected_current_revision_id], ).map_err(db)?; - conn.execute( - "DELETE FROM merge_request_final_results WHERE workspace_id=?1 AND merge_request_id=?2", - params![self.workspace_id,current.merge_request_id], - ).map_err(db)?; load_merge_request(conn, &self.workspace_id, &input.ticket_id)?.ok_or_else(|| MergeRequestError::NotFound(input.ticket_id.clone())) }) } @@ -659,27 +578,7 @@ 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 }); } - 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())); - } - let is_final: bool = conn.query_row( - "SELECT EXISTS(SELECT 1 FROM merge_request_final_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 !is_final { - return Err(MergeRequestError::MergeResultNotFinal); - } - 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)?; - } + 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], @@ -689,8 +588,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, 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], + "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], ).map_err(|_| MergeRequestError::InvalidReviewAttempt)?; Ok(()) }) @@ -716,12 +615,12 @@ impl SqliteMergeRequestStore { validate_review_input(&input)?; self.write(|conn| { let token = token_hash(&input.capability_token); - 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)?)), + 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)?)), ).optional().map_err(db)?; - 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 { + let Some((attempt_id, mr_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 { @@ -736,22 +635,10 @@ 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 }); } - if let Some(merge_result_id) = merge_result_id.as_deref() { - let is_final: bool = conn.query_row( - "SELECT EXISTS(SELECT 1 FROM merge_request_final_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 !is_final { - return Err(MergeRequestError::MergeResultNotFinal); - } - 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)?; - } + 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, 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], + "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], ).map_err(|_| MergeRequestError::InvalidReviewAttempt)?; for (ordinal, finding) in input.findings.iter().enumerate() { nonempty("finding.body", &finding.body)?; @@ -768,102 +655,14 @@ 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)?; - conn.execute( - "INSERT INTO merge_request_final_results (workspace_id,merge_request_id,revision_id,merge_result_id,selected_at) VALUES (?1,?2,?3,?4,?5) ON CONFLICT(workspace_id,merge_request_id) DO UPDATE SET revision_id=excluded.revision_id,merge_result_id=excluded.merge_result_id,selected_at=excluded.selected_at", - params![self.workspace_id,mr.merge_request_id,input.expected_revision_id,input.merge_result_id,input.created_at], - ).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()), ("ticket_id", input.ticket_id.as_str()), ("revision_id", input.expected_revision_id.as_str()), - ("merge_result_id", input.expected_merge_result_id.as_str()), - ( - "observed_target_commit", - input.observed_target_commit.as_str(), - ), + ("target_commit", input.target_commit.as_str()), + ("source_commit", input.source_commit.as_str()), + ("result_commit", input.result_commit.as_str()), ( "implementation_assignment_id", input.implementation_assignment_id.as_str(), @@ -879,6 +678,7 @@ impl SqliteMergeRequestStore { ] { nonempty(name, value)?; } + validate_completion_outcome(&input)?; let fingerprint = completion_fingerprint(&input); self.write(|conn| { if let Some((stored, status, state)) = conn.query_row( @@ -892,51 +692,36 @@ impl SqliteMergeRequestStore { } } else { conn.execute( - "INSERT INTO merge_request_completion_operations (workspace_id, operation_id, ticket_id, revision_id, merge_result_id, authority_kind, implementation_assignment_id, completion_actor_runtime_id, completion_actor_worker_id, fingerprint, status, created_at, updated_at) VALUES (?1,?2,?3,?4,?5,'workspace_orchestrator',?6,?7,?8,?9,'pending',?10,?10)", - params![self.workspace_id, input.operation_id, input.ticket_id, input.expected_revision_id, input.expected_merge_result_id, input.implementation_assignment_id, input.completion_actor_runtime_id, input.completion_actor_worker_id, fingerprint, input.now], + "INSERT INTO merge_request_completion_operations (workspace_id, operation_id, ticket_id, revision_id, authority_kind, implementation_assignment_id, completion_actor_runtime_id, completion_actor_worker_id, target_commit, source_commit, result_commit, strategy, resolution, fingerprint, status, created_at, updated_at) VALUES (?1,?2,?3,?4,'workspace_orchestrator',?5,?6,?7,?8,?9,?10,?11,?12,?13,'pending',?14,?14)", + params![self.workspace_id, input.operation_id, input.ticket_id, input.expected_revision_id, input.implementation_assignment_id, input.completion_actor_runtime_id, input.completion_actor_worker_id, input.target_commit, input.source_commit, input.result_commit, input.strategy.as_str(), input.resolution.as_str(), fingerprint, input.now], ).map_err(db)?; } let mr = load_merge_request(conn, &self.workspace_id, &input.ticket_id)? .ok_or_else(|| MergeRequestError::NotFound(input.ticket_id.clone()))?; - validate_current_implementation_assignment( - conn, - &self.workspace_id, - &input.ticket_id, - &input.implementation_assignment_id, - )?; + validate_current_implementation_assignment(conn, &self.workspace_id, &input.ticket_id, &input.implementation_assignment_id)?; ensure_open(&mr)?; 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::InvalidMergeOutcome("source commit does not match the current approved revision".into())); + } if mr.review_status != ReviewStatus::Approved { return Err(MergeRequestError::NotApproved); } - let final_result = mr.final_merge_result.as_ref().ok_or(MergeRequestError::FinalMergeResultMissing)?; - if final_result.merge_result_id != input.expected_merge_result_id { - return Err(MergeRequestError::MergeResultNotFinal); - } - if final_result.result_commit != input.observed_target_commit { - return Err(MergeRequestError::FinalMergeResultNotApplied); - } - if matches!(final_result.strategy, MergeStrategy::Merge) - && final_result.review_status != ReviewStatus::Approved - { - return Err(MergeRequestError::NotApproved); - } let current_state: String = conn.query_row( "SELECT workflow_state FROM typed_tickets WHERE workspace_id=?1 AND ticket_id=?2", params![self.workspace_id, input.ticket_id], |row| row.get(0), ).optional().map_err(db)?.ok_or_else(|| MergeRequestError::NotFound(input.ticket_id.clone()))?; - if current_state != "inprogress" { - return Err(MergeRequestError::TicketStateConflict(current_state)); - } + if current_state != "inprogress" { return Err(MergeRequestError::TicketStateConflict(current_state)); } let changed = conn.execute( "UPDATE typed_tickets SET workflow_state='done', workflow_state_explicit=1, updated_at=?3 WHERE workspace_id=?1 AND ticket_id=?2 AND workflow_state='inprogress'", params![self.workspace_id, input.ticket_id, input.now], ).map_err(db)?; if changed != 1 { return Err(MergeRequestError::TicketStateConflict("concurrent_change".into())); } - conn.execute( - "UPDATE merge_requests SET state='merged',merged_at=?3,updated_at=?3 WHERE workspace_id=?1 AND merge_request_id=?2 AND current_revision_id=?4 AND state='open'", - params![self.workspace_id,mr.merge_request_id,input.now,input.expected_revision_id], + let merged = conn.execute( + "UPDATE merge_requests SET state='merged', merged_revision_id=?3, merged_target_commit=?4, merged_result_commit=?5, merge_strategy=?6, merge_resolution=?7, merged_by_runtime_id=?8, merged_by_worker_id=?9, merged_at=?10, updated_at=?10 WHERE workspace_id=?1 AND merge_request_id=?2 AND state='open' AND current_revision_id=?3", + params![self.workspace_id, mr.merge_request_id, input.expected_revision_id, input.target_commit, input.result_commit, input.strategy.as_str(), input.resolution.as_str(), input.completion_actor_runtime_id, input.completion_actor_worker_id, input.now], ).map_err(db)?; + if merged != 1 { return Err(MergeRequestError::OperationConflict); } append_completion_event(conn, &self.workspace_id, &input)?; conn.execute( "UPDATE merge_request_completion_operations SET status='completed', result_ticket_state='done', updated_at=?3 WHERE workspace_id=?1 AND operation_id=?2 AND status='pending'", @@ -1030,14 +815,13 @@ fn migrate_locked(conn: &Connection, force_failure_after_v9_ddl: bool) -> Result if !marker_exists { if has_merge_request_domain_tables(conn)? { return Err(MergeRequestError::Database( - "unsupported unversioned legacy merge request schema; automatic migration requires a fresh database or exact version 9" + "unsupported unversioned legacy merge request schema; automatic migration requires a fresh database or exact version 8" .into(), )); } conn.execute_batch(MIGRATION_TABLE_SQL).map_err(db)?; conn.execute_batch(SCHEMA_V9).map_err(db)?; - migrate_v9_to_v10(conn)?; - verify_schema_v10(conn)?; + verify_schema_shape(conn, SCHEMA_V9, "v9")?; ensure_foreign_key_integrity(conn)?; replace_schema_marker(conn, SCHEMA_VERSION)?; return verify(conn); @@ -1049,50 +833,113 @@ fn migrate_locked(conn: &Connection, force_failure_after_v9_ddl: bool) -> Result verify_marker_state(conn, SCHEMA_VERSION)?; verify(conn) } - 9 => { - verify_marker_state(conn, 9)?; - verify_schema_shape(conn, SCHEMA_V9, "v9").map_err(|_| { + 8 => { + verify_marker_state(conn, 8)?; + if verify_schema_shape(conn, SCHEMA_V9, "v9").is_ok() { + ensure_foreign_key_integrity(conn)?; + replace_schema_marker(conn, SCHEMA_VERSION)?; + return verify(conn); + } + verify_schema_shape(conn, SCHEMA_V8, "v8").map_err(|_| { MergeRequestError::Database( - "schema drift at merge request version 9; automatic migration requires the exact v9 shape" + "schema drift at merge request version 8; automatic migration requires the exact v8 shape or a complete v9 shape for marker repair" .into(), ) })?; - migrate_v9_to_v10(conn)?; + migrate_v8_to_v9(conn)?; if force_failure_after_v9_ddl { return Err(MergeRequestError::Database( - "forced v9 to v10 migration failure after DDL".into(), + "forced v8 to v9 migration failure after DDL and data copy".into(), )); } - verify_schema_v10(conn)?; + verify_schema_shape(conn, SCHEMA_V9, "v9")?; ensure_foreign_key_integrity(conn)?; replace_schema_marker(conn, SCHEMA_VERSION)?; verify(conn) } - 0..=8 => Err(MergeRequestError::Database(format!( - "unsupported legacy merge request schema version {version}; automatic migration only supports exact v9 to v10" + 0..=7 => Err(MergeRequestError::Database(format!( + "unsupported legacy merge request schema version {version}; automatic migration only supports exact v8 to v9" ))), other => Err(MergeRequestError::Database(format!( - "unsupported merge request schema version {other}; expected version 9 or {SCHEMA_VERSION}" + "unsupported merge request schema version {other}; expected version 8 or {SCHEMA_VERSION}" ))), } } -fn migrate_v9_to_v10(conn: &Connection) -> Result<()> { - conn.execute_batch(MIGRATE_V9_TO_V10_SQL).map_err(db) -} +fn migrate_v8_to_v9(conn: &Connection) -> 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')); + ALTER TABLE merge_requests ADD COLUMN merged_revision_id TEXT; + ALTER TABLE merge_requests ADD COLUMN merged_target_commit TEXT; + ALTER TABLE merge_requests ADD COLUMN merged_result_commit TEXT; + ALTER TABLE merge_requests ADD COLUMN merge_strategy TEXT CHECK(merge_strategy IN ('fast_forward','merge')); + ALTER TABLE merge_requests ADD COLUMN merge_resolution TEXT CHECK(merge_resolution IN ('none','clean','conflicts_resolved')); + ALTER TABLE merge_requests ADD COLUMN merged_by_runtime_id TEXT; + ALTER TABLE merge_requests ADD COLUMN merged_by_worker_id TEXT; -fn verify_schema_v10(conn: &Connection) -> Result<()> { - let expected = Connection::open_in_memory().map_err(db)?; - expected.execute_batch(SCHEMA_V9).map_err(db)?; - expected.execute_batch(MIGRATE_V9_TO_V10_SQL).map_err(db)?; - let expected_shape = domain_schema_shape(&expected)?; - let actual_shape = domain_schema_shape(conn)?; - if actual_shape != expected_shape { - return Err(MergeRequestError::Database( - "schema drift: merge request v10 shape mismatch".into(), - )); - } - Ok(()) + 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_review_attempts_v9 ( + workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, ticket_id TEXT NOT NULL, + revision_id TEXT NOT NULL, lifecycle_generation INTEGER NOT NULL, + parent_assignment_id TEXT NOT NULL, parent_runtime_id TEXT NOT NULL, parent_worker_id TEXT NOT NULL, + child_session_id TEXT NOT NULL, child_effective_profile TEXT NOT NULL CHECK(child_effective_profile='builtin:reviewer'), + capability_token_sha256 TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('open','submitted','revoked')), + created_at TEXT NOT NULL, consumed_at TEXT, + PRIMARY KEY(workspace_id,attempt_id), UNIQUE(workspace_id,capability_token_sha256), UNIQUE(workspace_id,child_session_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,parent_assignment_id) REFERENCES ticket_worker_assignments(workspace_id,ticket_id,assignment_id), + FOREIGN KEY(workspace_id,child_session_id) REFERENCES merge_request_reviewer_child_sessions(workspace_id,child_session_id) + ); + INSERT INTO merge_request_review_attempts_v9( + 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,consumed_at + ) SELECT 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,consumed_at + FROM merge_request_review_attempts; + DROP TABLE merge_request_review_attempts; + ALTER TABLE merge_request_review_attempts_v9 RENAME TO merge_request_review_attempts; + + CREATE TABLE merge_request_completion_operations_v9 ( + workspace_id TEXT NOT NULL, operation_id TEXT NOT NULL, ticket_id TEXT NOT NULL, revision_id TEXT NOT NULL, + authority_kind TEXT NOT NULL CHECK(authority_kind IN ('workspace_orchestrator','legacy_assigned_coder')), + implementation_assignment_id TEXT NOT NULL, completion_actor_runtime_id TEXT, completion_actor_worker_id TEXT, + target_commit TEXT, source_commit TEXT, result_commit TEXT, + strategy TEXT CHECK(strategy IN ('fast_forward','merge')), + resolution TEXT CHECK(resolution IN ('none','clean','conflicts_resolved')), + fingerprint TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('pending','completed')), + result_ticket_state TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,operation_id), + FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id), + FOREIGN KEY(workspace_id,ticket_id,implementation_assignment_id) + REFERENCES ticket_worker_assignments(workspace_id,ticket_id,assignment_id) + ); + INSERT INTO merge_request_completion_operations_v9( + workspace_id,operation_id,ticket_id,revision_id,authority_kind,implementation_assignment_id, + completion_actor_runtime_id,completion_actor_worker_id,fingerprint,status,result_ticket_state,created_at,updated_at + ) SELECT workspace_id,operation_id,ticket_id,revision_id,authority_kind,implementation_assignment_id, + completion_actor_runtime_id,completion_actor_worker_id,fingerprint,status,result_ticket_state,created_at,updated_at + FROM merge_request_completion_operations; + DROP TABLE merge_request_completion_operations; + ALTER TABLE merge_request_completion_operations_v9 RENAME TO merge_request_completion_operations;", + ) + .map_err(db) } pub fn verify(conn: &Connection) -> Result<()> { @@ -1108,7 +955,7 @@ pub fn verify(conn: &Connection) -> Result<()> { ))); } verify_marker_state(conn, SCHEMA_VERSION)?; - verify_schema_v10(conn) + verify_schema_shape(conn, SCHEMA_V9, "v9") } fn schema_version(conn: &Connection) -> Result { @@ -1458,6 +1305,73 @@ fn column_exists(conn: &Connection, table: &str, column: &str) -> Result { const MIGRATION_TABLE: &str = "merge_request_schema_migrations"; const MIGRATION_TABLE_SQL: &str = "CREATE TABLE merge_request_schema_migrations (version INTEGER PRIMARY KEY, applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);"; +const SCHEMA_V8: &str = r#" +CREATE TABLE merge_requests ( + workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, + repository_id TEXT NOT NULL, state TEXT NOT NULL CHECK(state IN ('draft','open','closed','merged')), + lifecycle_generation INTEGER NOT NULL, current_revision_id TEXT NOT NULL, + created_at TEXT NOT NULL, updated_at TEXT NOT NULL, merged_by_account_id TEXT, merged_at TEXT, + PRIMARY KEY(workspace_id,merge_request_id), + FOREIGN KEY(workspace_id,repository_id) REFERENCES repositories(workspace_id,repository_id) +); +CREATE TABLE merge_request_ticket_relations ( + workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, ticket_id TEXT NOT NULL, + relation_kind TEXT NOT NULL CHECK(relation_kind='implements'), created_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,merge_request_id,ticket_id), + FOREIGN KEY(workspace_id,merge_request_id) REFERENCES merge_requests(workspace_id,merge_request_id) ON DELETE CASCADE, + FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id) ON DELETE CASCADE +); +CREATE TABLE merge_request_revisions ( + 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, head_tree 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 +); +CREATE TABLE merge_request_revision_paths ( + workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL, ordinal INTEGER NOT NULL, path TEXT NOT NULL, + PRIMARY KEY(workspace_id,merge_request_id,revision_id,ordinal), + FOREIGN KEY(workspace_id,merge_request_id,revision_id) REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id) ON DELETE CASCADE +); +CREATE TABLE merge_request_reviewer_child_sessions ( + workspace_id TEXT NOT NULL, child_session_id TEXT NOT NULL, parent_runtime_id TEXT NOT NULL, + parent_worker_id TEXT NOT NULL, effective_profile TEXT NOT NULL CHECK(effective_profile='builtin:reviewer'), registered_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,child_session_id) +); +CREATE TABLE merge_request_review_attempts ( + workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, ticket_id TEXT NOT NULL, + revision_id TEXT NOT NULL, lifecycle_generation INTEGER NOT NULL, + parent_assignment_id TEXT NOT NULL, parent_runtime_id TEXT NOT NULL, parent_worker_id TEXT NOT NULL, + child_session_id TEXT NOT NULL, child_effective_profile TEXT NOT NULL CHECK(child_effective_profile='builtin:reviewer'), + capability_token_sha256 TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('open','submitted','revoked')), + created_at TEXT NOT NULL, consumed_at TEXT, + PRIMARY KEY(workspace_id,attempt_id), UNIQUE(workspace_id,capability_token_sha256), UNIQUE(workspace_id,child_session_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,parent_assignment_id) REFERENCES ticket_worker_assignments(workspace_id,ticket_id,assignment_id) +); +CREATE TABLE merge_request_reviews ( + workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL, + decision TEXT NOT NULL CHECK(decision IN ('approve','request_changes')), body TEXT NOT NULL, submitted_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,attempt_id), + FOREIGN KEY(workspace_id,attempt_id) REFERENCES merge_request_review_attempts(workspace_id,attempt_id), + FOREIGN KEY(workspace_id,merge_request_id,revision_id) REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id) +); +CREATE TABLE merge_request_review_findings ( + workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, ordinal INTEGER NOT NULL, severity TEXT NOT NULL, + code TEXT, path TEXT, line INTEGER, body TEXT NOT NULL, PRIMARY KEY(workspace_id,attempt_id,ordinal), + FOREIGN KEY(workspace_id,attempt_id) REFERENCES merge_request_reviews(workspace_id,attempt_id) ON DELETE CASCADE +); +CREATE TABLE merge_request_completion_operations ( + workspace_id TEXT NOT NULL, operation_id TEXT NOT NULL, ticket_id TEXT NOT NULL, revision_id TEXT NOT NULL, + authority_kind TEXT NOT NULL CHECK(authority_kind IN ('workspace_orchestrator','legacy_assigned_coder')), + implementation_assignment_id TEXT NOT NULL, completion_actor_runtime_id TEXT, completion_actor_worker_id TEXT, + fingerprint TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('pending','completed')), + result_ticket_state TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, + PRIMARY KEY(workspace_id,operation_id), + FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id) +); +"#; + const SCHEMA_V9: &str = r#" CREATE TABLE merge_requests ( workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, @@ -1466,6 +1380,10 @@ CREATE TABLE merge_requests ( created_at TEXT NOT NULL, updated_at TEXT NOT NULL, merged_by_account_id TEXT, merged_at TEXT, target_ref_selector TEXT, target_status TEXT NOT NULL DEFAULT 'unknown' CHECK(target_status IN ('known','unknown')), + merged_revision_id TEXT, merged_target_commit TEXT, merged_result_commit TEXT, + merge_strategy TEXT CHECK(merge_strategy IN ('fast_forward','merge')), + merge_resolution TEXT CHECK(merge_resolution IN ('none','clean','conflicts_resolved')), + merged_by_runtime_id TEXT, merged_by_worker_id TEXT, PRIMARY KEY(workspace_id,merge_request_id), FOREIGN KEY(workspace_id,repository_id) REFERENCES repositories(workspace_id,repository_id) ); @@ -1499,48 +1417,31 @@ CREATE TABLE merge_request_review_attempts ( parent_assignment_id TEXT NOT NULL, parent_runtime_id TEXT NOT NULL, parent_worker_id TEXT NOT NULL, child_session_id TEXT NOT NULL, child_effective_profile TEXT NOT NULL CHECK(child_effective_profile='builtin:reviewer'), capability_token_sha256 TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('open','submitted','revoked')), - created_at TEXT NOT NULL, consumed_at TEXT, merge_result_id TEXT, + created_at TEXT NOT NULL, consumed_at TEXT, PRIMARY KEY(workspace_id,attempt_id), UNIQUE(workspace_id,capability_token_sha256), UNIQUE(workspace_id,child_session_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,parent_assignment_id) REFERENCES ticket_worker_assignments(workspace_id,ticket_id,assignment_id), - FOREIGN KEY(workspace_id,child_session_id) REFERENCES merge_request_reviewer_child_sessions(workspace_id,child_session_id), - FOREIGN KEY(workspace_id,merge_result_id) REFERENCES merge_request_merge_results(workspace_id,merge_result_id) + FOREIGN KEY(workspace_id,child_session_id) REFERENCES merge_request_reviewer_child_sessions(workspace_id,child_session_id) ); CREATE TABLE merge_request_reviews ( workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL, decision TEXT NOT NULL CHECK(decision IN ('approve','request_changes')), body TEXT NOT NULL, submitted_at TEXT NOT NULL, - merge_result_id TEXT, PRIMARY KEY(workspace_id,attempt_id), FOREIGN KEY(workspace_id,attempt_id) REFERENCES merge_request_review_attempts(workspace_id,attempt_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,merge_result_id) REFERENCES merge_request_merge_results(workspace_id,merge_result_id) + FOREIGN KEY(workspace_id,merge_request_id,revision_id) REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id) ); CREATE TABLE merge_request_review_findings ( workspace_id TEXT NOT NULL, attempt_id TEXT NOT NULL, ordinal INTEGER NOT NULL, severity TEXT NOT NULL, code TEXT, path TEXT, line INTEGER, body TEXT NOT NULL, PRIMARY KEY(workspace_id,attempt_id,ordinal), FOREIGN KEY(workspace_id,attempt_id) REFERENCES merge_request_reviews(workspace_id,attempt_id) ON DELETE CASCADE ); -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); CREATE TABLE merge_request_completion_operations ( workspace_id TEXT NOT NULL, operation_id TEXT NOT NULL, ticket_id TEXT NOT NULL, revision_id TEXT NOT NULL, authority_kind TEXT NOT NULL CHECK(authority_kind IN ('workspace_orchestrator','legacy_assigned_coder')), implementation_assignment_id TEXT NOT NULL, completion_actor_runtime_id TEXT, completion_actor_worker_id TEXT, + target_commit TEXT, source_commit TEXT, result_commit TEXT, + strategy TEXT CHECK(strategy IN ('fast_forward','merge')), + resolution TEXT CHECK(resolution IN ('none','clean','conflicts_resolved')), fingerprint TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('pending','completed')), result_ticket_state TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, PRIMARY KEY(workspace_id,operation_id), @@ -1550,38 +1451,14 @@ CREATE TABLE merge_request_completion_operations ( ); "#; -const MIGRATE_V9_TO_V10_SQL: &str = r#" -CREATE TABLE merge_request_final_results ( - workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL, - merge_result_id TEXT NOT NULL, selected_at TEXT NOT NULL, - PRIMARY KEY(workspace_id,merge_request_id), - FOREIGN KEY(workspace_id,merge_request_id,revision_id) - REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id) ON DELETE CASCADE, - FOREIGN KEY(workspace_id,merge_result_id) - REFERENCES merge_request_merge_results(workspace_id,merge_result_id) -); -ALTER TABLE merge_request_completion_operations ADD COLUMN merge_result_id TEXT; -"#; - #[cfg(test)] mod migration_tests { use super::*; const SUPPORT_SCHEMA: &str = r#" -CREATE TABLE repositories( - workspace_id TEXT NOT NULL, repository_id TEXT NOT NULL, - PRIMARY KEY(workspace_id,repository_id) -); -CREATE TABLE typed_tickets( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, workflow_state TEXT NOT NULL, - workflow_state_explicit INTEGER NOT NULL DEFAULT 1, updated_at TEXT NOT NULL, - PRIMARY KEY(workspace_id,ticket_id) -); -CREATE TABLE ticket_worker_assignments( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, assignment_id TEXT NOT NULL, - runtime_id TEXT NOT NULL, worker_id TEXT NOT NULL, - PRIMARY KEY(workspace_id,ticket_id,assignment_id) -); +CREATE TABLE repositories(workspace_id TEXT NOT NULL,repository_id TEXT NOT NULL,PRIMARY KEY(workspace_id,repository_id)); +CREATE TABLE typed_tickets(workspace_id TEXT NOT NULL,ticket_id TEXT NOT NULL,workflow_state TEXT NOT NULL,workflow_state_explicit INTEGER NOT NULL DEFAULT 1,updated_at TEXT NOT NULL,PRIMARY KEY(workspace_id,ticket_id)); +CREATE TABLE ticket_worker_assignments(workspace_id TEXT NOT NULL,ticket_id TEXT NOT NULL,assignment_id TEXT NOT NULL,runtime_id TEXT NOT NULL,worker_id TEXT NOT NULL,PRIMARY KEY(workspace_id,ticket_id,assignment_id)); "#; fn fresh_connection() -> Connection { @@ -1590,27 +1467,29 @@ CREATE TABLE ticket_worker_assignments( conn } - fn exact_v9_connection() -> Connection { + fn exact_v8_connection() -> Connection { let conn = fresh_connection(); conn.execute_batch(MIGRATION_TABLE_SQL).unwrap(); conn.execute( - "INSERT INTO merge_request_schema_migrations(version) VALUES(9)", + "INSERT INTO merge_request_schema_migrations(version) VALUES(8)", [], ) .unwrap(); - conn.execute_batch(SCHEMA_V9).unwrap(); + conn.execute_batch(SCHEMA_V8).unwrap(); conn.execute_batch( "INSERT INTO repositories VALUES('ws','repo'); INSERT INTO typed_tickets VALUES('ws','T1','inprogress',1,'t0'); INSERT INTO ticket_worker_assignments VALUES('ws','T1','A1','R1','W1'); - INSERT INTO merge_requests VALUES('ws','MR1','repo','open',3,'V1','t0','t1',NULL,NULL,'refs/heads/develop','known'); + INSERT INTO merge_requests VALUES('ws','MR1','repo','open',3,'V1','t0','t1',NULL,NULL); INSERT INTO merge_request_ticket_relations VALUES('ws','MR1','T1','implements','t0'); - INSERT INTO merge_request_revisions VALUES('ws','MR1','V1',1,'base','head','digest','summary','A1','t0'); - INSERT INTO merge_request_merge_results VALUES('ws','R1','MR1','T1','V1','base','head','head','fast_forward','none','R1','W1','t1','OPR1','fp1','t1'); - INSERT INTO merge_request_merge_results VALUES('ws','R2','MR1','T1','V1','base','head','merge','merge','clean','R1','W1','t2','OPR2','fp2','t2'); - INSERT INTO merge_request_completion_operations VALUES('ws','OP1','T1','V1','workspace_orchestrator','A1','OR','OW','fingerprint','pending',NULL,'t0','t1');", - ) - .unwrap(); + INSERT INTO merge_request_revisions VALUES('ws','MR1','V1',1,'base','head','legacy-tree','digest','summary','A1','t0'); + INSERT INTO merge_request_revision_paths VALUES('ws','MR1','V1',0,'src/lib.rs'); + INSERT INTO merge_request_reviewer_child_sessions VALUES('ws','C1','R1','W1','builtin:reviewer','t0'); + INSERT INTO merge_request_review_attempts VALUES('ws','AT1','MR1','T1','V1',3,'A1','R1','W1','C1','builtin:reviewer','token','submitted','t0','t1'); + INSERT INTO merge_request_reviews VALUES('ws','AT1','MR1','V1','approve','approved','t1'); + INSERT INTO merge_request_review_findings VALUES('ws','AT1',0,'warning','C','src/lib.rs',7,'finding'); + INSERT INTO merge_request_completion_operations VALUES('ws','OP1','T1','V1','workspace_orchestrator','A1','R1','W1','fp','pending',NULL,'t0','t1');", + ).unwrap(); conn } @@ -1624,20 +1503,64 @@ CREATE TABLE ticket_worker_assignments( } #[test] - fn fresh_database_materializes_latest_v10_contract() { + fn fresh_database_materializes_final_merge_evidence_contract() { let conn = fresh_connection(); migrate(&conn).unwrap(); verify(&conn).unwrap(); - assert_eq!(marker_version(&conn), 10); - assert!(table_exists(&conn, "merge_request_final_results").unwrap()); - assert!( - column_exists( - &conn, - "merge_request_completion_operations", - "merge_result_id" + assert_eq!(marker_version(&conn), 9); + for column in [ + "target_ref_selector", + "merged_revision_id", + "merged_target_commit", + "merged_result_commit", + "merge_strategy", + "merge_resolution", + "merged_by_runtime_id", + "merged_by_worker_id", + ] { + assert!( + column_exists(&conn, "merge_requests", column).unwrap(), + "missing {column}" + ); + } + assert!(!column_exists(&conn, "merge_request_revisions", "head_tree").unwrap()); + assert!(!table_exists(&conn, "merge_request_merge_results").unwrap()); + } + + #[test] + fn exact_v8_migrates_preserving_review_and_operation_evidence() { + let conn = exact_v8_connection(); + migrate(&conn).unwrap(); + verify(&conn).unwrap(); + assert_eq!(marker_version(&conn), 9); + assert_eq!( + conn.query_row( + "SELECT head_commit FROM merge_request_revisions WHERE revision_id='V1'", + [], + |row| row.get::<_, String>(0) ) - .unwrap() + .unwrap(), + "head" ); + assert_eq!( + conn.query_row( + "SELECT decision FROM merge_request_reviews WHERE attempt_id='AT1'", + [], + |row| row.get::<_, String>(0) + ) + .unwrap(), + "approve" + ); + assert_eq!( + conn.query_row( + "SELECT target_ref_selector FROM merge_requests WHERE merge_request_id='MR1'", + [], + |row| row.get::<_, Option>(0) + ) + .unwrap(), + None + ); + assert_eq!(conn.query_row("SELECT result_commit FROM merge_request_completion_operations WHERE operation_id='OP1'", [], |row| row.get::<_,Option>(0)).unwrap(), None); assert_eq!( conn.query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |row| { row.get::<_, i64>(0) @@ -1648,85 +1571,47 @@ CREATE TABLE ticket_worker_assignments( } #[test] - fn exact_v9_migrates_without_guessing_a_final_candidate() { - let conn = exact_v9_connection(); - migrate(&conn).unwrap(); - verify(&conn).unwrap(); - assert_eq!(marker_version(&conn), 10); - assert_eq!( - conn.query_row( - "SELECT COUNT(*) FROM merge_request_merge_results", - [], - |row| row.get::<_, i64>(0) - ) - .unwrap(), - 2 - ); - assert_eq!( - conn.query_row( - "SELECT COUNT(*) FROM merge_request_final_results", - [], - |row| row.get::<_, i64>(0) - ) - .unwrap(), - 0, - "migration must not guess which historical candidate is final" - ); - assert_eq!( - conn.query_row("SELECT merge_result_id FROM merge_request_completion_operations WHERE operation_id='OP1'", [], |row| row.get::<_, Option>(0)).unwrap(), - None - ); - } - - #[test] - fn v9_to_v10_failure_rolls_back_schema_and_marker() { - let conn = exact_v9_connection(); + fn v8_to_v9_failure_rolls_back_schema_data_and_marker() { + let conn = exact_v8_connection(); let error = migrate_with_failpoint(&conn, true).unwrap_err(); assert!( error .to_string() - .contains("forced v9 to v10 migration failure") + .contains("forced v8 to v9 migration failure") ); - assert_eq!(marker_version(&conn), 9); - assert!(!table_exists(&conn, "merge_request_final_results").unwrap()); - assert!( - !column_exists( - &conn, - "merge_request_completion_operations", - "merge_result_id" - ) - .unwrap() - ); - verify_schema_shape(&conn, SCHEMA_V9, "v9 after rollback").unwrap(); + assert_eq!(marker_version(&conn), 8); + assert!(column_exists(&conn, "merge_request_revisions", "head_tree").unwrap()); + assert!(!column_exists(&conn, "merge_requests", "merged_result_commit").unwrap()); + verify_schema_shape(&conn, SCHEMA_V8, "v8 after rollback").unwrap(); } #[test] - fn drifted_v9_fails_closed_without_mutation() { - let conn = exact_v9_connection(); + fn drifted_v8_fails_closed_without_mutation() { + let conn = exact_v8_connection(); conn.execute_batch("ALTER TABLE merge_requests ADD COLUMN drift TEXT;") .unwrap(); let error = migrate(&conn).unwrap_err(); assert!( error .to_string() - .contains("schema drift at merge request version 9") + .contains("schema drift at merge request version 8") ); - assert_eq!(marker_version(&conn), 9); - assert!(!table_exists(&conn, "merge_request_final_results").unwrap()); + assert_eq!(marker_version(&conn), 8); + assert!(!column_exists(&conn, "merge_requests", "merged_result_commit").unwrap()); } #[test] - fn versions_older_than_v9_are_rejected() { + fn versions_older_than_v8_are_rejected() { let conn = fresh_connection(); conn.execute_batch(MIGRATION_TABLE_SQL).unwrap(); conn.execute( - "INSERT INTO merge_request_schema_migrations(version) VALUES(8)", + "INSERT INTO merge_request_schema_migrations(version) VALUES(7)", [], ) .unwrap(); let error = migrate(&conn).unwrap_err(); - assert!(error.to_string().contains("only supports exact v9 to v10")); - assert_eq!(marker_version(&conn), 8); + assert!(error.to_string().contains("only supports exact v8 to v9")); + assert_eq!(marker_version(&conn), 7); } } @@ -1735,9 +1620,30 @@ fn load_merge_request( workspace_id: &str, ticket_id: &str, ) -> Result> { - 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)?)), + type Row = ( + String, + String, + String, + Option, + String, + String, + i64, + String, + String, + String, + Option, + Option, + Option, + Option, + Option, + Option, + Option, + Option, + ); + let row: 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_revision_id,mr.merged_target_commit,mr.merged_result_commit,mr.merge_strategy,mr.merge_resolution,mr.merged_by_runtime_id,mr.merged_by_worker_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)?,r.get(12)?,r.get(13)?,r.get(14)?,r.get(15)?,r.get(16)?,r.get(17)?)), ).optional().map_err(db)?; let Some(( mr_id, @@ -1750,7 +1656,13 @@ fn load_merge_request( revision_id, created_at, updated_at, - merged_by_account_id, + merged_revision_id, + merged_target_commit, + merged_result_commit, + merge_strategy, + merge_resolution, + merged_by_runtime_id, + merged_by_worker_id, merged_at, )) = row else { @@ -1763,22 +1675,6 @@ 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)?; - let final_merge_result_id: Option = conn - .query_row( - "SELECT merge_result_id FROM merge_request_final_results WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3", - params![workspace_id,mr_id,revision_id], - |row| row.get(0), - ) - .optional() - .map_err(db)?; - let final_merge_result = final_merge_result_id.and_then(|id| { - merge_results - .iter() - .find(|result| result.merge_result_id == id) - .cloned() - }); Ok(Some(MergeRequest { merge_request_id: mr_id, workspace_id: workspace_id.into(), @@ -1792,11 +1688,15 @@ fn load_merge_request( current_revision: revision, review_status, current_review, - merge_results, - final_merge_result, created_at, updated_at, - merged_by_account_id, + merged_revision_id, + merged_target_commit, + merged_result_commit, + merge_strategy: merge_strategy.as_deref().map(MergeStrategy::parse), + merge_resolution: merge_resolution.as_deref().map(MergeResolution::parse), + merged_by_runtime_id, + merged_by_worker_id, merged_at, })) } @@ -1840,7 +1740,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 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)?; + 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)?; match attempt { Some(id) => load_review(conn, workspace_id, &id), None => Ok(None), @@ -1852,13 +1752,12 @@ fn load_review( workspace_id: &str, attempt_id: &str, ) -> Result> { - 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)?)), + 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)?)), ).optional().map_err(db)?; let Some(( revision_id, - merge_result_id, decision, body, assignment, @@ -1888,7 +1787,6 @@ fn load_review( Ok(Some(MergeRequestReview { attempt_id: attempt_id.into(), revision_id, - merge_result_id, decision: ReviewDecision::parse(&decision), body, findings, @@ -1921,15 +1819,6 @@ 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, @@ -2011,138 +1900,29 @@ 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 validate_completion_outcome(input: &CompleteMergeRequest) -> Result<()> { + match input.strategy { + MergeStrategy::FastForward + if input.result_commit == input.source_commit + && input.resolution == MergeResolution::None => {} + MergeStrategy::Merge if input.resolution != MergeResolution::None => {} + MergeStrategy::FastForward => { + return Err(MergeRequestError::InvalidMergeOutcome( + "fast-forward result must equal the approved source commit and use resolution=none" + .into(), + )); + } + MergeStrategy::Merge => { + return Err(MergeRequestError::InvalidMergeOutcome( + "merge strategy requires clean or conflicts_resolved resolution".into(), + )); + } + } + Ok(()) } 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, - }; - } - let final_id = mr - .final_merge_result - .as_ref() - .map(|result| result.merge_result_id.clone()); - mr.final_merge_result = final_id.and_then(|id| { - mr.merge_results - .iter() - .find(|result| result.merge_result_id == id) - .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<()> { @@ -2200,11 +1980,14 @@ fn token_hash(token: &str) -> String { } fn completion_fingerprint(input: &CompleteMergeRequest) -> String { token_hash(&format!( - "workspace_orchestrator\0{}\0{}\0{}\0{}\0{}\0{}\0{}", + "workspace_orchestrator\0{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}", input.ticket_id, input.expected_revision_id, - input.expected_merge_result_id, - input.observed_target_commit, + input.target_commit, + input.source_commit, + input.result_commit, + input.strategy.as_str(), + input.resolution.as_str(), input.implementation_assignment_id, input.completion_actor_runtime_id, input.completion_actor_worker_id diff --git a/crates/merge-request/tests/store.rs b/crates/merge-request/tests/store.rs index 0caed23c..b6e46edb 100644 --- a/crates/merge-request/tests/store.rs +++ b/crates/merge-request/tests/store.rs @@ -1,5 +1,7 @@ use merge_request::*; use rusqlite::{Connection, params}; +use std::sync::{Arc, Barrier}; +use std::thread; use tempfile::TempDir; fn setup() -> (TempDir, SqliteMergeRequestStore) { @@ -38,6 +40,7 @@ fn setup() -> (TempDir, SqliteMergeRequestStore) { let store = SqliteMergeRequestStore::open(&path, "ws-a").unwrap(); (dir, store) } + fn revision(id: &str, ordinal: u64, head: &str) -> MergeRequestRevision { MergeRequestRevision { revision_id: id.into(), @@ -51,653 +54,199 @@ fn revision(id: &str, ordinal: u64, head: &str) -> MergeRequestRevision { created_at: format!("t{ordinal}"), } } + fn open(store: &SqliteMergeRequestStore) { store .open_merge_request(OpenMergeRequest { merge_request_id: "MR1".into(), ticket_id: "T1".into(), repository_id: "repo".into(), - target_ref_selector: "develop".into(), - revision: revision("V1", 1, "h1"), + target_ref_selector: "refs/heads/develop".into(), + revision: revision("V1", 1, "head"), authenticated_runtime_id: "R1".into(), authenticated_worker_id: "W1".into(), now: "t1".into(), }) .unwrap(); } -fn attempt(store: &SqliteMergeRequestStore, id: &str, revision: &str, token: &str, child: &str) { + +fn attempt(store: &SqliteMergeRequestStore, revision: &str, token: &str) { + let child = format!("child-{revision}"); store .register_reviewer_child_session(RegisterReviewerChildSession { parent_runtime_id: "R1".into(), parent_worker_id: "W1".into(), - child_session_id: child.into(), - now: "t".into(), - }) - .unwrap(); - store - .register_review_attempt(RegisterReviewAttempt { - 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(), - child_session_id: child.into(), - capability_token: token.into(), - now: "t".into(), - }) - .unwrap(); -} -fn review( - store: &SqliteMergeRequestStore, - revision: &str, - token: &str, - decision: ReviewDecision, -) -> Result { - store.submit_review(SubmitReview { - ticket_id: "T1".into(), - revision_id: revision.into(), - merge_result_id: None, - capability_token: token.into(), - decision, - body: "evidence".into(), - findings: vec![], - now: "tr".into(), - }) -} - -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(); - open(&store); - store - .open_merge_request(OpenMergeRequest { - 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(), + child_session_id: child.clone(), now: "t2".into(), }) .unwrap(); - let conn = Connection::open(store.db_path()).unwrap(); - let count:i64=conn.query_row("SELECT COUNT(*) FROM merge_request_ticket_relations WHERE workspace_id='ws-a' AND ticket_id='T1'",[],|row|row.get(0)).unwrap(); - assert_eq!(count, 2); - assert_eq!( - store - .show_for_ticket("T1") - .unwrap() - .unwrap() - .merge_request_id, - "MR2" - ); -} - -#[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) - )); - - // Multiple valid candidates for the same target are retained as history. The - // most recently recorded candidate is the one explicit final result. - record_result( - &store, - "V1", - "T1", - "h1", - "h1", - MergeStrategy::FastForward, - MergeResolution::None, - "op-same-target-old", - ); - 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) - ); - let candidates = store - .show_for_ticket_with_target("T1", Some("T1")) - .unwrap() - .unwrap(); - assert_eq!(candidates.merge_results.len(), 3); - assert_eq!( - candidates - .final_merge_result - .as_ref() - .map(|result| result.merge_result_id.as_str()), - Some("result-op-merge") - ); - store - .register_reviewer_child_session(RegisterReviewerChildSession { - parent_runtime_id: "runtime-orchestrator".into(), - parent_worker_id: "workspace-orchestrator".into(), - child_session_id: "child-old-candidate".into(), - now: "2026-07-26T00:00:04Z".into(), - }) - .unwrap(); - let old_candidate_review = store.register_review_attempt(RegisterReviewAttempt { - attempt_id: "attempt-old-candidate".into(), - ticket_id: "T1".into(), - revision_id: "V1".into(), - merge_result_id: Some("result-op-same-target-old".into()), - parent_assignment_id: "A1".into(), - parent_runtime_id: "runtime-orchestrator".into(), - parent_worker_id: "workspace-orchestrator".into(), - child_session_id: "child-old-candidate".into(), - capability_token: "token-old-candidate".into(), - now: "2026-07-26T00:00:04Z".into(), - }); - assert!(matches!( - old_candidate_review, - Err(MergeRequestError::MergeResultNotFinal) - )); - - 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(), + attempt_id: format!("attempt-{revision}"), ticket_id: "T1".into(), - revision_id: "V1".into(), - merge_result_id: Some("result-op-merge".into()), + revision_id: revision.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(), + parent_runtime_id: "R1".into(), + parent_worker_id: "W1".into(), + child_session_id: child, + capability_token: token.into(), + now: "t2".into(), }) .unwrap(); +} + +fn approve(store: &SqliteMergeRequestStore, revision: &str, token: &str) { + attempt(store, revision, token); 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(), + revision_id: revision.into(), + capability_token: token.into(), decision: ReviewDecision::Approve, - body: "merge evidence is valid".into(), + body: "approved".into(), findings: vec![], - now: "2026-07-26T00:00:05Z".into(), + now: "t3".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 - .final_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(); - let mut oversized = revision("V1", 1, "h1"); - oversized.changed_paths = (0..=1_000).map(|i| format!("src/{i}.rs")).collect(); - let result = store.open_merge_request(OpenMergeRequest { - merge_request_id: "MR1".into(), +fn completion(operation_id: &str) -> CompleteMergeRequest { + CompleteMergeRequest { + operation_id: operation_id.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(), - now: "t".into(), - }); - assert!(matches!( - result, - Err(MergeRequestError::TooLarge { - field: "revision.changed_paths", - .. - }) - )); + expected_revision_id: "V1".into(), + target_commit: "base".into(), + source_commit: "head".into(), + result_commit: "head".into(), + strategy: MergeStrategy::FastForward, + resolution: MergeResolution::None, + implementation_assignment_id: "A1".into(), + completion_actor_runtime_id: "OR".into(), + completion_actor_worker_id: "OW".into(), + now: "t4".into(), + } } #[test] -fn v6_legacy_schema_fails_closed_without_archiving() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("legacy.db"); - let conn = Connection::open(&path).unwrap(); - conn.execute_batch( - "CREATE TABLE merge_request_schema_migrations(version INTEGER PRIMARY KEY,name TEXT NOT NULL,applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);\ - INSERT INTO merge_request_schema_migrations(version,name) VALUES(6,'rejected_merge_request_v6');\ - CREATE TABLE merge_requests(workspace_id TEXT NOT NULL,merge_request_id TEXT NOT NULL,repository_id TEXT NOT NULL,state TEXT NOT NULL,lifecycle_generation INTEGER NOT NULL,current_revision_id TEXT NOT NULL,created_at TEXT NOT NULL,updated_at TEXT NOT NULL,PRIMARY KEY(workspace_id,merge_request_id));", - ).unwrap(); - drop(conn); - - let error = SqliteMergeRequestStore::open(&path, "ws-a").unwrap_err(); - assert!( - error - .to_string() - .contains("unsupported legacy merge request schema version 6") - ); - let conn = Connection::open(&path).unwrap(); - let version: i64 = conn - .query_row( - "SELECT MAX(version) FROM merge_request_schema_migrations", - [], - |row| row.get(0), - ) - .unwrap(); - assert_eq!(version, 6); - let original: i64 = conn - .query_row( - "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='merge_requests'", - [], - |row| row.get(0), - ) - .unwrap(); - let archived: i64 = conn - .query_row( - "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name LIKE 'legacy_v6_%'", - [], - |row| row.get(0), - ) - .unwrap(); - assert_eq!(original, 1); - assert_eq!(archived, 0); -} - -#[test] -fn v7_schema_is_rejected_without_mutating_completion_evidence() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("v7.db"); - let conn = Connection::open(&path).unwrap(); - conn.execute_batch( - "CREATE TABLE merge_request_schema_migrations(version INTEGER PRIMARY KEY,name TEXT NOT NULL,applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);\ - INSERT INTO merge_request_schema_migrations(version,name) VALUES(7,'fresh_bounded_context_authority');\ - CREATE TABLE typed_tickets(workspace_id TEXT NOT NULL,ticket_id TEXT NOT NULL,workflow_state TEXT NOT NULL,workflow_state_explicit INTEGER NOT NULL DEFAULT 1,updated_at TEXT NOT NULL,PRIMARY KEY(workspace_id,ticket_id));\ - INSERT INTO typed_tickets VALUES('ws-a','T1','done',1,'t');\ - CREATE TABLE merge_request_completion_operations(workspace_id TEXT NOT NULL,operation_id TEXT NOT NULL,ticket_id TEXT NOT NULL,revision_id TEXT NOT NULL,assignment_id TEXT NOT NULL,fingerprint TEXT NOT NULL,status TEXT NOT NULL CHECK(status IN ('pending','completed')),result_ticket_state TEXT,created_at TEXT NOT NULL,updated_at TEXT NOT NULL,PRIMARY KEY(workspace_id,operation_id),FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id));\ - INSERT INTO merge_request_completion_operations VALUES('ws-a','legacy-op','T1','V1','A1','legacy-fingerprint','completed','done','t','t');", - ).unwrap(); - drop(conn); - - let error = SqliteMergeRequestStore::open(&path, "ws-a").unwrap_err(); - assert!( - error - .to_string() - .contains("unsupported legacy merge request schema version 7") - ); - let conn = Connection::open(&path).unwrap(); - let row: (String, String) = conn.query_row( - "SELECT assignment_id,fingerprint FROM merge_request_completion_operations WHERE operation_id='legacy-op'", - [], - |row| Ok((row.get(0)?, row.get(1)?)), - ).unwrap(); - assert_eq!(row, ("A1".into(), "legacy-fingerprint".into())); - let version: i64 = conn - .query_row( - "SELECT MAX(version) FROM merge_request_schema_migrations", - [], - |row| row.get(0), - ) - .unwrap(); - assert_eq!(version, 7); -} - -#[test] -fn request_changes_new_revision_resets_and_exact_completion_replay_converges() { +fn target_movement_does_not_invalidate_source_revision_approval() { let (_dir, store) = setup(); open(&store); - attempt(&store, "AT1", "V1", "tok1", "child1"); - review(&store, "V1", "tok1", ReviewDecision::RequestChanges).unwrap(); - assert_eq!( - store.show_for_ticket("T1").unwrap().unwrap().review_status, - ReviewStatus::ChangesRequested - ); + approve(&store, "V1", "token-v1"); + for target in ["base", "advanced-target"] { + let readiness = store + .readiness_for_ticket_with_target("T1", Some(target)) + .unwrap(); + assert!( + readiness.ready, + "target movement must not invalidate source approval" + ); + assert_eq!(readiness.review_status, ReviewStatus::Approved); + assert_eq!(readiness.observed_target_commit.as_deref(), Some(target)); + } store .add_revision(AddRevision { ticket_id: "T1".into(), expected_current_revision_id: "V1".into(), - revision: revision("V2", 2, "h2"), + revision: revision("V2", 2, "head2"), authenticated_runtime_id: "R1".into(), authenticated_worker_id: "W1".into(), - now: "t2".into(), + now: "t5".into(), }) .unwrap(); assert_eq!( - store.show_for_ticket("T1").unwrap().unwrap().review_status, + store.readiness_for_ticket("T1").unwrap().review_status, ReviewStatus::Pending ); - assert!(review(&store, "V1", "tok1", ReviewDecision::Approve).is_err()); - attempt(&store, "AT2", "V2", "tok2", "child2"); - review(&store, "V2", "tok2", ReviewDecision::Approve).unwrap(); - let missing_result = store.complete(CompleteMergeRequest { - operation_id: "OP-missing-result".into(), - ticket_id: "T1".into(), - expected_revision_id: "V2".into(), - expected_merge_result_id: "missing".into(), - observed_target_commit: "h2".into(), - implementation_assignment_id: "A1".into(), - completion_actor_runtime_id: "OR".into(), - completion_actor_worker_id: "OW".into(), - now: "tc".into(), - }); - assert!(matches!( - missing_result, - Err(MergeRequestError::FinalMergeResultMissing) - )); - record_result( - &store, - "V2", - "base", - "h2", - "h2", - MergeStrategy::FastForward, - MergeResolution::None, - "final-v2", - ); - let not_applied = store.complete(CompleteMergeRequest { - operation_id: "OP-not-applied".into(), - ticket_id: "T1".into(), - expected_revision_id: "V2".into(), - expected_merge_result_id: "result-final-v2".into(), - observed_target_commit: "base".into(), - implementation_assignment_id: "A1".into(), - completion_actor_runtime_id: "OR".into(), - completion_actor_worker_id: "OW".into(), - now: "tc".into(), - }); - assert!(matches!( - not_applied, - Err(MergeRequestError::FinalMergeResultNotApplied) - )); - let input = CompleteMergeRequest { - operation_id: "OP1".into(), - ticket_id: "T1".into(), - expected_revision_id: "V2".into(), - expected_merge_result_id: "result-final-v2".into(), - observed_target_commit: "h2".into(), - implementation_assignment_id: "A1".into(), - completion_actor_runtime_id: "OR".into(), - completion_actor_worker_id: "OW".into(), - now: "tc".into(), - }; - let first = store.complete(input.clone()).unwrap(); +} + +#[test] +fn completion_records_one_final_merge_outcome_and_replays_idempotently() { + let (_dir, store) = setup(); + open(&store); + approve(&store, "V1", "token-v1"); + let first = store.complete(completion("OP1")).unwrap(); assert!(!first.replayed); + assert_eq!(first.ticket_state, "done"); + let merged = store.show_for_ticket("T1").unwrap().unwrap(); + assert_eq!(merged.state, MergeRequestState::Merged); + assert_eq!(merged.merged_revision_id.as_deref(), Some("V1")); + assert_eq!(merged.merged_target_commit.as_deref(), Some("base")); + assert_eq!(merged.merged_result_commit.as_deref(), Some("head")); + assert_eq!(merged.merge_strategy, Some(MergeStrategy::FastForward)); + assert_eq!(merged.merge_resolution, Some(MergeResolution::None)); + assert_eq!(merged.merged_by_runtime_id.as_deref(), Some("OR")); + assert_eq!(merged.merged_by_worker_id.as_deref(), Some("OW")); + assert!(store.complete(completion("OP1")).unwrap().replayed); + let mut conflicting = completion("OP1"); + conflicting.target_commit = "other".into(); + assert!(matches!( + store.complete(conflicting), + Err(MergeRequestError::OperationConflict) + )); +} + +#[test] +fn completion_rejects_invalid_or_non_current_source_outcomes_without_side_effects() { + let (_dir, store) = setup(); + open(&store); + approve(&store, "V1", "token-v1"); + let mut invalid_ff = completion("bad-ff"); + invalid_ff.result_commit = "different".into(); + assert!(matches!( + store.complete(invalid_ff), + Err(MergeRequestError::InvalidMergeOutcome(_)) + )); + let mut invalid_merge = completion("bad-merge"); + invalid_merge.strategy = MergeStrategy::Merge; + assert!(matches!( + store.complete(invalid_merge), + Err(MergeRequestError::InvalidMergeOutcome(_)) + )); + let mut wrong_source = completion("wrong-source"); + wrong_source.source_commit = "not-approved".into(); + wrong_source.result_commit = "not-approved".into(); + assert!(matches!( + store.complete(wrong_source), + Err(MergeRequestError::InvalidMergeOutcome(_)) + )); assert_eq!( store.show_for_ticket("T1").unwrap().unwrap().state, - MergeRequestState::Merged + MergeRequestState::Open ); - assert_eq!( - store - .show_for_ticket("T1") - .unwrap() - .unwrap() - .merged_at - .as_deref(), - Some("tc") - ); - let replay = store.complete(input).unwrap(); - assert!(replay.replayed); - assert_eq!(replay.ticket_state, "done"); let conn = Connection::open(store.db_path()).unwrap(); assert_eq!( conn.query_row( "SELECT workflow_state FROM typed_tickets WHERE workspace_id='ws-a' AND ticket_id='T1'", [], - |r| r.get::<_, String>(0) + |row| row.get::<_, String>(0) ) .unwrap(), - "done" + "inprogress" ); - assert_eq!( - conn.query_row( - "SELECT COUNT(*) FROM typed_ticket_events WHERE workspace_id='ws-a' AND ticket_id='T1'", - [], - |r| r.get::<_, i64>(0) - ) - .unwrap(), - 1 - ); - assert_eq!( - conn.query_row( - "SELECT authority_kind || ':' || implementation_assignment_id || ':' || completion_actor_runtime_id || ':' || completion_actor_worker_id FROM merge_request_completion_operations WHERE workspace_id='ws-a' AND operation_id='OP1'", - [], - |r| r.get::<_, String>(0) - ) - .unwrap(), - "workspace_orchestrator:A1:OR:OW" - ); - assert_eq!( - conn.query_row( - "SELECT author FROM typed_ticket_events WHERE workspace_id='ws-a' AND ticket_id='T1' AND kind='state_changed'", - [], - |r| r.get::<_, String>(0) - ) - .unwrap(), - "worker:OR:OW" - ); - let authority: String = conn - .query_row( - "SELECT value FROM typed_ticket_event_attributes WHERE workspace_id='ws-a' AND ticket_id='T1' AND key='completion_authority'", - [], - |r| r.get(0), - ) - .unwrap(); - assert_eq!(authority, "workspace_orchestrator"); } #[test] -fn spoof_self_approval_replay_and_cross_workspace_are_rejected() { +fn concurrent_completion_converges_on_one_operation() { let (_dir, store) = setup(); open(&store); - let mut bad = RegisterReviewAttempt { - 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(), - child_session_id: "W1".into(), - capability_token: "bad".into(), - now: "t".into(), - }; - assert!(matches!( - store.register_review_attempt(bad.clone()), - Err(MergeRequestError::SelfApproval) - )); - bad.child_session_id = "child".into(); - assert!(matches!( - store.register_review_attempt(bad), - Err(MergeRequestError::InvalidReviewer) - )); - attempt(&store, "AT", "V1", "secret", "child"); - assert!(review(&store, "V1", "spoof", ReviewDecision::Approve).is_err()); - review(&store, "V1", "secret", ReviewDecision::Approve).unwrap(); - assert!(review(&store, "V1", "secret", ReviewDecision::Approve).is_err()); - let other = SqliteMergeRequestStore::open_verified(store.db_path(), "ws-b").unwrap(); - assert!(other.show_for_ticket("T1").unwrap().is_none()); -} - -#[test] -fn reopen_resets_approval() { - let (_dir, store) = setup(); - open(&store); - attempt(&store, "AT", "V1", "token", "child"); - review(&store, "V1", "token", ReviewDecision::Approve).unwrap(); - store.close("T1", "V1", "tc").unwrap(); - let reopened = store.reopen("T1", "V1", "tr").unwrap(); - assert_eq!(reopened.review_status, ReviewStatus::Pending); -} - -#[test] -fn concurrent_exact_completion_replays_commit_one_ticket_side_effect() { - let (_dir, store) = setup(); - open(&store); - attempt(&store, "AT", "V1", "token", "child"); - review(&store, "V1", "token", ReviewDecision::Approve).unwrap(); - record_result( - &store, - "V1", - "base", - "h1", - "h1", - MergeStrategy::FastForward, - MergeResolution::None, - "final-concurrent", - ); - let input = CompleteMergeRequest { - operation_id: "OP-concurrent".into(), - ticket_id: "T1".into(), - expected_revision_id: "V1".into(), - expected_merge_result_id: "result-final-concurrent".into(), - observed_target_commit: "h1".into(), - implementation_assignment_id: "A1".into(), - completion_actor_runtime_id: "OR".into(), - completion_actor_worker_id: "OW".into(), - now: "t".into(), - }; - let left_store = store.clone(); - let left_input = input.clone(); - let left = std::thread::spawn(move || left_store.complete(left_input)); - let right_store = store.clone(); - let right = std::thread::spawn(move || right_store.complete(input)); - let outcomes = [ - left.join().unwrap().unwrap(), - right.join().unwrap().unwrap(), - ]; + approve(&store, "V1", "token-v1"); + let path = store.db_path().to_path_buf(); + let barrier = Arc::new(Barrier::new(3)); + let mut handles = Vec::new(); + for _ in 0..2 { + let path = path.clone(); + let barrier = barrier.clone(); + handles.push(thread::spawn(move || { + let store = SqliteMergeRequestStore::open_verified(path, "ws-a").unwrap(); + barrier.wait(); + store.complete(completion("OP-concurrent")) + })); + } + barrier.wait(); + let outcomes: Vec<_> = handles + .into_iter() + .map(|handle| handle.join().unwrap().unwrap()) + .collect(); assert_eq!( outcomes.iter().filter(|outcome| !outcome.replayed).count(), 1 @@ -706,75 +255,44 @@ fn concurrent_exact_completion_replays_commit_one_ticket_side_effect() { outcomes.iter().filter(|outcome| outcome.replayed).count(), 1 ); - let conn = Connection::open(store.db_path()).unwrap(); - let events: i64 = conn - .query_row( - "SELECT COUNT(*) FROM typed_ticket_events WHERE workspace_id='ws-a' AND ticket_id='T1'", - [], - |row| row.get(0), - ) - .unwrap(); - assert_eq!(events, 1); } #[test] -fn operation_key_mismatch_and_actor_or_assignment_change_are_fenced() { +fn reviewer_attempt_is_bound_to_direct_child_and_current_assignment() { let (_dir, store) = setup(); open(&store); - attempt(&store, "AT", "V1", "token", "child"); - review(&store, "V1", "token", ReviewDecision::Approve).unwrap(); - record_result( - &store, - "V1", - "base", - "h1", - "h1", - MergeStrategy::FastForward, - MergeResolution::None, - "final-operation", - ); - let mut input = CompleteMergeRequest { - operation_id: "OP".into(), + store + .register_reviewer_child_session(RegisterReviewerChildSession { + parent_runtime_id: "R1".into(), + parent_worker_id: "W1".into(), + child_session_id: "child".into(), + now: "t2".into(), + }) + .unwrap(); + store + .register_review_attempt(RegisterReviewAttempt { + attempt_id: "attempt".into(), + ticket_id: "T1".into(), + revision_id: "V1".into(), + parent_assignment_id: "A1".into(), + parent_runtime_id: "R1".into(), + parent_worker_id: "W1".into(), + child_session_id: "child".into(), + capability_token: "token".into(), + now: "t2".into(), + }) + .unwrap(); + let wrong_token = store.submit_review(SubmitReview { ticket_id: "T1".into(), - expected_revision_id: "V1".into(), - expected_merge_result_id: "result-final-operation".into(), - observed_target_commit: "h1".into(), - implementation_assignment_id: "A1".into(), - completion_actor_runtime_id: "OR".into(), - completion_actor_worker_id: "OW".into(), - now: "t".into(), - }; - let conn = Connection::open(store.db_path()).unwrap(); - conn.execute( - "UPDATE ticket_current_worker_assignments SET assignment_id='A2',runtime_id='R2',worker_id='W2' WHERE workspace_id='ws-a' AND ticket_id='T1'", - [], - ) - .unwrap(); + revision_id: "V1".into(), + capability_token: "wrong".into(), + decision: ReviewDecision::Approve, + body: "approved".into(), + findings: vec![], + now: "t3".into(), + }); assert!(matches!( - store.complete(input.clone()), - Err(MergeRequestError::AssignmentMismatch) - )); - conn.execute( - "UPDATE ticket_current_worker_assignments SET assignment_id='A1',runtime_id='R1',worker_id='W1' WHERE workspace_id='ws-a' AND ticket_id='T1'", - [], - ) - .unwrap(); - store.complete(input.clone()).unwrap(); - input.completion_actor_worker_id = "other".into(); - assert!(matches!( - store.complete(input.clone()), - Err(MergeRequestError::OperationConflict) - )); - input.completion_actor_worker_id = "OW".into(); - input.implementation_assignment_id = "A2".into(); - assert!(matches!( - store.complete(input.clone()), - Err(MergeRequestError::OperationConflict) - )); - input.implementation_assignment_id = "A1".into(); - input.expected_revision_id = "other".into(); - assert!(matches!( - store.complete(input), - Err(MergeRequestError::OperationConflict) + wrong_token, + Err(MergeRequestError::InvalidReviewAttempt) )); } diff --git a/crates/worker/src/feature/builtin/merge_request.rs b/crates/worker/src/feature/builtin/merge_request.rs index 6a9aa923..89edd913 100644 --- a/crates/worker/src/feature/builtin/merge_request.rs +++ b/crates/worker/src/feature/builtin/merge_request.rs @@ -12,7 +12,6 @@ pub const MERGE_REQUEST_COMMON_TOOL_NAMES: &[&str] = &[ "MergeRequestReadinessCheck", "MergeRequestOpen", "MergeRequestAddRevision", - "MergeRequestRecordMergeResult", "MergeRequestComplete", ]; pub const MERGE_REQUEST_REVIEW_TOOL_NAME: &str = "MergeRequestReviewSubmit"; @@ -22,7 +21,6 @@ enum Kind { Readiness, Open, AddRevision, - RecordMergeResult, Complete, Review, } @@ -63,10 +61,10 @@ struct AddRevisionInput { summary: String, } #[derive(Debug, Deserialize, JsonSchema)] -struct RecordMergeResultInput { +struct CompleteInput { ticket: String, - expected_current_revision_id: String, operation_id: String, + expected_revision_id: String, target_commit: String, source_commit: String, result_commit: String, @@ -87,12 +85,6 @@ enum MergeResolutionInput { ConflictsResolved, } #[derive(Debug, Deserialize, JsonSchema)] -struct CompleteInput { - ticket: String, - operation_id: String, - expected_revision_id: String, -} -#[derive(Debug, Deserialize, JsonSchema)] struct ReviewInput { decision: ReviewDecisionInput, #[serde(default)] @@ -125,7 +117,6 @@ impl Kind { Self::Readiness => "MergeRequestReadinessCheck", Self::Open => "MergeRequestOpen", Self::AddRevision => "MergeRequestAddRevision", - Self::RecordMergeResult => "MergeRequestRecordMergeResult", Self::Complete => "MergeRequestComplete", Self::Review => "MergeRequestReviewSubmit", } @@ -138,7 +129,6 @@ 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)), } @@ -202,8 +192,8 @@ impl Tool for MergeRequestTool { ), ) } - Kind::RecordMergeResult => { - let v: RecordMergeResultInput = parse(input)?; + Kind::Complete => { + let v: CompleteInput = parse(input)?; nonempty(&v.ticket)?; let strategy = match v.strategy { MergeStrategyInput::FastForward => "fast_forward", @@ -214,20 +204,6 @@ impl Tool for MergeRequestTool { 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}), - ), - ) - } - Kind::Complete => { - let v: CompleteInput = parse(input)?; - nonempty(&v.ticket)?; ( WorkspaceRequestMethod::Post, format!( @@ -235,7 +211,7 @@ impl Tool for MergeRequestTool { v.ticket ), Some( - json!({"operation_id":v.operation_id,"expected_revision_id":v.expected_revision_id}), + json!({"operation_id":v.operation_id,"expected_revision_id":v.expected_revision_id,"target_commit":v.target_commit,"source_commit":v.source_commit,"result_commit":v.result_commit,"strategy":strategy,"resolution":resolution}), ), ) } @@ -310,7 +286,6 @@ 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), ] } @@ -338,9 +313,6 @@ 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.") } @@ -356,14 +328,14 @@ mod tests { use super::*; #[test] - fn merge_request_tool_contract_omits_tree_hashes_and_exposes_merge_result() { + fn merge_request_tool_contract_omits_tree_hashes_and_candidate_result_tool() { 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(); + let complete = serde_json::to_string(&schemars::schema_for!(CompleteInput)).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")); + assert!(complete.contains("result_commit")); + assert!(complete.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 c1fb6379..5f29ef5f 100644 --- a/crates/worker/src/spawn/tool.rs +++ b/crates/worker/src/spawn/tool.rs @@ -68,9 +68,6 @@ 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)] @@ -425,7 +422,6 @@ impl Tool for SubWorkerSpawnTool { ( review.ticket_id.clone(), review.revision_id.clone(), - review.merge_result_id.clone(), uuid::Uuid::now_v7().to_string(), format!( "{}{}", @@ -435,9 +431,7 @@ impl Tool for SubWorkerSpawnTool { ) }); let child_workspace_context = - if let Some((ticket_id, revision_id, merge_result_id, _, capability_token)) = - &reviewer_attempt - { + if let Some((ticket_id, revision_id, _, capability_token)) = &reviewer_attempt { let workspace_id = self.workspace_context .workspace_id() @@ -459,7 +453,6 @@ 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(), )); @@ -554,9 +547,7 @@ impl Tool for SubWorkerSpawnTool { } }; - if let Some((ticket_id, revision_id, merge_result_id, attempt_id, capability_token)) = - &reviewer_attempt - { + if let Some((ticket_id, revision_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()) })?; @@ -588,7 +579,6 @@ 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 bda08079..a9b5bd6d 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -249,7 +249,6 @@ 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)] @@ -311,14 +310,6 @@ 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()), @@ -462,7 +453,6 @@ 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 a55ad387..aa273fe1 100644 --- a/crates/workspace-server/src/repositories.rs +++ b/crates/workspace-server/src/repositories.rs @@ -97,13 +97,38 @@ pub struct CommitObservation { #[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 }, + 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, + }, + TargetMoved { + id: RepositoryId, + selector: String, + expected: String, + observed: Option, + }, + ProviderFailure { + id: RepositoryId, + operation: String, + }, } #[derive(Debug, Clone)] @@ -286,6 +311,48 @@ impl RepositoryRegistryReader { } } + pub fn update_merge_target( + &self, + id: &str, + selector: &str, + expected_target: &str, + result_commit: &str, + ) -> Result<(), RepositoryLookupError> { + let repository = self.merge_repository(id)?; + if !selector.starts_with("refs/heads/") + || selector.starts_with('-') + || selector.as_bytes().contains(&0) + { + return Err(RepositoryLookupError::InvalidSelector { + id: id.into(), + selector: selector.into(), + }); + } + self.observe_commit(id, result_commit)?; + let status = Command::new("git") + .arg("-C") + .arg(&repository.path) + .args(["update-ref", selector, result_commit, expected_target]) + .status() + .map_err(|_| RepositoryLookupError::ProviderFailure { + id: id.into(), + operation: "guarded target update".into(), + })?; + if status.success() { + return Ok(()); + } + let observed = self + .observe_merge_target(id, Some(selector)) + .ok() + .map(|target| target.commit); + Err(RepositoryLookupError::TargetMoved { + id: id.into(), + selector: selector.into(), + expected: expected_target.into(), + observed, + }) + } + fn merge_repository(&self, id: &str) -> Result<&ConfiguredRepository, RepositoryLookupError> { let repository = self .find(id) @@ -656,6 +723,20 @@ mod tests { vec![base.clone()] ); reader.ensure_ancestor("main", &base, &source).unwrap(); + reader + .update_merge_target("main", "refs/heads/main", &base, &source) + .unwrap(); + assert_eq!( + reader + .observe_merge_target("main", Some("refs/heads/main")) + .unwrap() + .commit, + source + ); + assert!(matches!( + reader.update_merge_target("main", "refs/heads/main", &base, &base), + Err(RepositoryLookupError::TargetMoved { .. }) + )); assert!(matches!( reader.ensure_ancestor("main", &source, &base), Err(RepositoryLookupError::InvalidCommitRelation { .. }) diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 6e28c8ae..27f6de80 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1290,10 +1290,6 @@ 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), @@ -3541,17 +3537,6 @@ 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, @@ -3561,8 +3546,6 @@ struct RegisterReviewerChildSessionRequest { struct RegisterMergeRequestReviewAttemptRequest { attempt_id: String, revision_id: String, - #[serde(default)] - merge_result_id: Option, child_session_id: String, capability_token: String, } @@ -3570,8 +3553,6 @@ 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)] @@ -3584,6 +3565,11 @@ struct SubmitMergeRequestReviewRequest { struct CompleteMergeRequestRequest { operation_id: String, expected_revision_id: String, + target_commit: String, + source_commit: String, + result_commit: String, + strategy: merge_request::MergeStrategy, + resolution: merge_request::MergeResolution, } #[derive(Debug, serde::Deserialize)] @@ -3840,105 +3826,6 @@ 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, @@ -3974,9 +3861,7 @@ async fn scoped_register_merge_request_review_attempt( .ok_or_else(|| { Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; - if input.merge_result_id.is_some() { - require_online_workspace_orchestrator_source(&api, &source)?; - } else if assignment.worker.runtime_id != source.runtime_id + if assignment.worker.runtime_id != source.runtime_id || assignment.worker.worker_id != source.worker_id { return Err(Error::TicketAssignmentConflict( @@ -3989,7 +3874,6 @@ 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, @@ -4011,7 +3895,6 @@ 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, @@ -4038,47 +3921,141 @@ async fn scoped_complete_merge_request( Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; 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(), - )) - })?; - let final_result = observed - .final_merge_result - .as_ref() - .ok_or(merge_request::MergeRequestError::FinalMergeResultMissing)?; - if final_result.target_status != merge_request::MergeResultTargetStatus::Applied { - return Err(merge_request::MergeRequestError::FinalMergeResultNotApplied.into()); + let mr = store + .show_for_ticket(&ticket_id)? + .ok_or_else(|| merge_request::MergeRequestError::NotFound(ticket_id.clone()))?; + if mr.current_revision.revision_id != input.expected_revision_id { + return Err(merge_request::MergeRequestError::StaleRevision { + expected: input.expected_revision_id, + current: mr.current_revision.revision_id, + } + .into()); } - let readiness = store.readiness_for_ticket_with_target(&ticket_id, target_commit.as_deref())?; - if !readiness.ready { + if mr.review_status != merge_request::ReviewStatus::Approved { + return Err(merge_request::MergeRequestError::NotApproved.into()); + } + if input.source_commit != mr.current_revision.head_commit { + return Err(merge_request::MergeRequestError::InvalidMergeOutcome( + "source commit does not match the current approved revision".into(), + ) + .into()); + } + if mr.state == merge_request::MergeRequestState::Merged { + return Ok(Json(store.complete( + merge_request::CompleteMergeRequest { + operation_id: input.operation_id, + ticket_id, + expected_revision_id: mr.current_revision.revision_id, + target_commit: input.target_commit, + source_commit: input.source_commit, + result_commit: input.result_commit, + strategy: input.strategy, + resolution: input.resolution, + implementation_assignment_id: assignment.assignment_id, + completion_actor_runtime_id: source.runtime_id, + completion_actor_worker_id: source.worker_id, + now: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), + }, + )?)); + } + let selector = mr + .target_ref_selector + .as_deref() + .ok_or(merge_request::MergeRequestError::UnknownTarget)?; + let repositories = api.repository_reader(); + let observed_target = repositories + .observe_merge_target(&mr.repository_id, Some(selector)) + .map_err(repository_merge_evidence_error)?; + let source_commit = repositories + .observe_commit(&mr.repository_id, &input.source_commit) + .map_err(repository_merge_evidence_error)?; + if source_commit.commit != input.source_commit { + return Err(Error::InvalidInput("source commit must be canonical".into()).into()); + } + let result_commit = repositories + .observe_commit(&mr.repository_id, &input.result_commit) + .map_err(repository_merge_evidence_error)?; + if result_commit.commit != input.result_commit { + return Err(Error::InvalidInput("result commit must be canonical".into()).into()); + } + match input.strategy { + merge_request::MergeStrategy::FastForward => { + if input.resolution != merge_request::MergeResolution::None + || input.result_commit != input.source_commit + { + return Err(merge_request::MergeRequestError::InvalidMergeOutcome( + "fast-forward result must equal the approved source and use resolution=none" + .into(), + ) + .into()); + } + repositories + .ensure_ancestor( + &mr.repository_id, + &input.target_commit, + &input.source_commit, + ) + .map_err(repository_merge_evidence_error)?; + } + merge_request::MergeStrategy::Merge => { + if input.resolution == merge_request::MergeResolution::None + || result_commit.parents + != vec![input.target_commit.clone(), input.source_commit.clone()] + { + return Err(merge_request::MergeRequestError::InvalidMergeOutcome( + "merge result must have the expected target and approved source as its two ordered parents" + .into(), + ) + .into()); + } + } + } + let target_was_already_updated = observed_target.commit == input.result_commit; + if observed_target.commit != input.target_commit && !target_was_already_updated { return Err(Error::InvalidInput(format!( - "Merge Request is not completion-ready: {}", - readiness.blockers.join("; ") + "Merge Request target moved: expected {}, observed {}", + input.target_commit, observed_target.commit )) .into()); } - let outcome = store.complete(merge_request::CompleteMergeRequest { + if !target_was_already_updated { + repositories + .update_merge_target( + &mr.repository_id, + selector, + &input.target_commit, + &input.result_commit, + ) + .map_err(repository_merge_evidence_error)?; + } + let completion = merge_request::CompleteMergeRequest { operation_id: input.operation_id, ticket_id, - expected_revision_id: input.expected_revision_id, - expected_merge_result_id: final_result.merge_result_id.clone(), - observed_target_commit: target_commit - .ok_or(merge_request::MergeRequestError::FinalMergeResultNotApplied)?, + expected_revision_id: mr.current_revision.revision_id, + target_commit: input.target_commit.clone(), + source_commit: input.source_commit, + result_commit: input.result_commit.clone(), + strategy: input.strategy, + resolution: input.resolution, implementation_assignment_id: assignment.assignment_id, completion_actor_runtime_id: source.runtime_id, completion_actor_worker_id: source.worker_id, now: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), - })?; - Ok(Json(outcome)) + }; + match store.complete(completion) { + Ok(outcome) => Ok(Json(outcome)), + Err(error) => { + if !target_was_already_updated { + let _ = repositories.update_merge_target( + &mr.repository_id, + selector, + &input.result_commit, + &input.target_commit, + ); + } + Err(error.into()) + } + } } async fn scoped_reopen_merge_request( @@ -4089,9 +4066,7 @@ async fn scoped_reopen_merge_request( ) -> ApiResult> { let workspace_id = parse_workspace_id(&workspace_id)?; require_workspace_access(&workspace_id, &api)?; - if headers.contains_key("authorization") { - return Err(Error::BrowserReopenConfirmationRequired.into()); - } + reject_non_browser_reopen_auth(&headers)?; let _actor = require_actor(&api, &headers).await?; if !input.explicit_confirmation { return Err(Error::BrowserReopenConfirmationRequired.into()); @@ -4103,6 +4078,13 @@ async fn scoped_reopen_merge_request( )?)) } +fn reject_non_browser_reopen_auth(headers: &HeaderMap) -> Result<()> { + if headers.contains_key("authorization") { + return Err(Error::BrowserReopenConfirmationRequired); + } + Ok(()) +} + async fn scoped_close_ticket_record( State(api): State, AxumPath((workspace_id, id)): AxumPath<(String, String)>, @@ -11719,6 +11701,17 @@ mod tests { ObjectiveTicketLinkRecord, SqliteWorkspaceStore, WorkspaceRecord, }; + #[test] + fn reopen_confirmation_rejects_api_token_actor_before_session_resolution() { + let mut headers = HeaderMap::new(); + headers.insert("authorization", "Bearer api-token".parse().unwrap()); + assert!(matches!( + reject_non_browser_reopen_auth(&headers), + Err(Error::BrowserReopenConfirmationRequired) + )); + assert!(reject_non_browser_reopen_auth(&HeaderMap::new()).is_ok()); + } + #[test] fn flow_or_generic_worker_state_change_is_not_ticket_completion_authority() { let operation = TicketBackendOperation::SetWorkflowState { @@ -12365,42 +12358,36 @@ mod tests { async fn merge_request_completion_endpoint_rejects_coder_and_accepts_orchestrator() { let workspace = tempfile::tempdir().unwrap(); init_clean_git_workspace(workspace.path()); - let git_output = |args: &[&str]| { + 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(), "git {:?} failed", args); + assert!(output.status.success()); String::from_utf8(output.stdout).unwrap().trim().to_string() }; - let base_commit = git_output(&["rev-parse", "HEAD"]); - std::fs::write(workspace.path().join("README.md"), "completion candidate\n").unwrap(); + 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(["add", "README.md"]) - .status() - .unwrap() - .success() - ); - assert!( - std::process::Command::new("git") - .arg("-C") - .arg(workspace.path()) - .args(["commit", "-m", "candidate"]) - .status() - .unwrap() - .success() - ); - let head_commit = git_output(&["rev-parse", "HEAD"]); - assert!( - std::process::Command::new("git") - .arg("-C") - .arg(workspace.path()) - .args(["reset", "--hard", &base_commit]) + .args(["reset", "--hard", &target_commit]) .status() .unwrap() .success() @@ -12441,12 +12428,12 @@ 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(), + target_ref_selector: target_ref.clone(), revision: merge_request::MergeRequestRevision { revision_id: "V1".into(), ordinal: 1, - base_commit: base_commit.clone(), - head_commit: head_commit.clone(), + base_commit: target_commit.clone(), + head_commit: source_commit.clone(), diff_digest: "sha256:diff".into(), changed_paths: vec!["src/lib.rs".into()], @@ -12472,7 +12459,6 @@ 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(), @@ -12485,7 +12471,6 @@ 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(), @@ -12493,31 +12478,6 @@ mod tests { now: "t3".into(), }) .unwrap(); - mr_store - .record_merge_result(merge_request::RecordMergeResult { - operation_id: "record-complete-result".into(), - merge_result_id: "result-complete".into(), - ticket_id: ticket.id.clone(), - expected_revision_id: "V1".into(), - target_commit: base_commit.clone(), - source_commit: head_commit.clone(), - result_commit: head_commit.clone(), - strategy: merge_request::MergeStrategy::FastForward, - resolution: merge_request::MergeResolution::None, - actor_runtime_id: coder.worker_ref.runtime_id.clone(), - actor_worker_id: coder.worker_ref.worker_id.clone(), - created_at: "t4".into(), - }) - .unwrap(); - assert!( - std::process::Command::new("git") - .arg("-C") - .arg(workspace.path()) - .args(["reset", "--hard", &head_commit]) - .status() - .unwrap() - .success() - ); let worker_headers = |worker: &RuntimeWorkerRef| { let mut headers = HeaderMap::new(); @@ -12534,6 +12494,11 @@ mod tests { let request = || CompleteMergeRequestRequest { operation_id: "complete-operation".into(), expected_revision_id: "V1".into(), + target_commit: target_commit.clone(), + source_commit: source_commit.clone(), + result_commit: source_commit.clone(), + strategy: merge_request::MergeStrategy::FastForward, + resolution: merge_request::MergeResolution::None, }; let coder_error = scoped_complete_merge_request( State(api.clone()), @@ -12574,6 +12539,35 @@ mod tests { .workflow_state, TicketWorkflowState::Done ); + assert_eq!( + api.repository_reader() + .observe_merge_target(TEST_REPOSITORY_ID, Some(&target_ref)) + .unwrap() + .commit, + source_commit + ); + let merged = mr_store.show_for_ticket(&ticket.id).unwrap().unwrap(); + assert_eq!( + merged.merged_target_commit.as_deref(), + Some(target_commit.as_str()) + ); + assert_eq!( + merged.merged_result_commit.as_deref(), + Some(source_commit.as_str()) + ); + assert_eq!( + merged.merge_strategy, + Some(merge_request::MergeStrategy::FastForward) + ); + 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!(replayed.replayed); let conn = rusqlite::Connection::open(&api.config.database_path).unwrap(); let actor: String = conn .query_row( @@ -13418,7 +13412,6 @@ mod tests { for args in [ vec!["add", "README.md", ".gitignore"], vec!["commit", "-m", "init"], - vec!["branch", "-M", "develop"], ] { let status = std::process::Command::new("git") .arg("-C") diff --git a/web/workspace/src/lib/workspace/console/worker-console.ui.test.ts b/web/workspace/src/lib/workspace/console/worker-console.ui.test.ts index 6b70bbb1..0a653e59 100644 --- a/web/workspace/src/lib/workspace/console/worker-console.ui.test.ts +++ b/web/workspace/src/lib/workspace/console/worker-console.ui.test.ts @@ -249,7 +249,7 @@ Deno.test("workspace Tickets surface provides Kanban and lifecycle controls", as ticketDetailPage.includes('mutate("state", "/state"') && ticketDetailPage.includes('mutate("queue", "/queue"') && !ticketDetailPage.includes("/merge-request/merge") && - ticketDetailPage.includes("final_merge_result") && + ticketDetailPage.includes("merged_result_commit") && !ticketDetailPage.includes('mutate("review", "/review"') && ticketDetailPage.includes('mutate("close", "/close"') && ticketDetailPage.includes("ticketWorkerLaunchHref") && 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 1e09683e..1a32d5c0 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte @@ -25,16 +25,13 @@ 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; - final_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_revision_id?: string | null; + merged_target_commit?: string | null; + merged_result_commit?: string | null; + merge_strategy?: "fast_forward" | "merge" | null; + merge_resolution?: "none" | "clean" | "conflicts_resolved" | null; + merged_by_runtime_id?: string | null; + merged_by_worker_id?: string | null; merged_at?: string | null; }; @@ -361,15 +358,18 @@ {#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.final_merge_result} + {#if mergeRequest.merged_result_commit}

- Final MergeResult {mergeRequest.final_merge_result.merge_result_id} · - {mergeRequest.final_merge_result.strategy} / {mergeRequest.final_merge_result.resolution} · - {mergeRequest.final_merge_result.target_status} / {mergeRequest.final_merge_result.review_status} + Final merge · {mergeRequest.merge_strategy} / {mergeRequest.merge_resolution} +

+

+ Target before {mergeRequest.merged_target_commit} · result + {mergeRequest.merged_result_commit} +

+

+ Revision {mergeRequest.merged_revision_id} · completed by + {mergeRequest.merged_by_runtime_id}/{mergeRequest.merged_by_worker_id}

-

Result {mergeRequest.final_merge_result.result_commit}

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

No final MergeResult has been selected for this revision.

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

{mergeRequest.current_revision.summary}

{/if} {#if mergeRequest.current_review}