merge-request: add target and merge result authority

This commit is contained in:
2026-08-16 21:03:30 +09:00
parent 4f042cae84
commit 027f60d262
8 changed files with 1439 additions and 73 deletions
+563 -36
View File
@@ -11,7 +11,7 @@ use std::path::{Path, PathBuf};
use std::time::Duration;
use thiserror::Error;
const SCHEMA_VERSION: i64 = 8;
const SCHEMA_VERSION: i64 = 9;
const REVIEWER_PROFILE: &str = "builtin:reviewer";
const MAX_SUMMARY_BYTES: usize = 16 * 1024;
const MAX_REVIEW_BODY_BYTES: usize = 64 * 1024;
@@ -51,6 +51,14 @@ pub enum MergeRequestError {
NotOpen(String),
#[error("completion operation id was reused with different input")]
OperationConflict,
#[error("merge result operation id was reused with different input")]
MergeResultOperationConflict,
#[error("merge result {0} was not found for the current Merge Request revision")]
MergeResultNotFound(String),
#[error("merge result evidence is invalid: {0}")]
InvalidMergeResult(String),
#[error("Merge Request target is unknown and must be resolved explicitly")]
UnknownTarget,
#[error("Ticket must be inprogress before Merge Request completion (current: {0})")]
TicketStateConflict(String),
#[error("only an authenticated user with explicit confirmation may merge")]
@@ -117,13 +125,86 @@ pub enum ReviewStatus {
ChangesRequested,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MergeRequestTargetStatus {
Known,
Unknown,
}
impl MergeRequestTargetStatus {
fn parse(value: &str) -> Self {
match value {
"known" => Self::Known,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MergeStrategy {
FastForward,
Merge,
}
impl MergeStrategy {
fn as_str(self) -> &'static str {
match self {
Self::FastForward => "fast_forward",
Self::Merge => "merge",
}
}
fn parse(value: &str) -> Self {
match value {
"fast_forward" => Self::FastForward,
_ => Self::Merge,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MergeResolution {
None,
Clean,
ConflictsResolved,
}
impl MergeResolution {
fn as_str(self) -> &'static str {
match self {
Self::None => "none",
Self::Clean => "clean",
Self::ConflictsResolved => "conflicts_resolved",
}
}
fn parse(value: &str) -> Self {
match value {
"none" => Self::None,
"conflicts_resolved" => Self::ConflictsResolved,
_ => Self::Clean,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MergeResultTargetStatus {
Current,
Applied,
Stale,
Unknown,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MergeRequestRevision {
pub revision_id: String,
pub ordinal: u64,
pub base_commit: String,
pub head_commit: String,
pub head_tree: String,
pub diff_digest: String,
pub changed_paths: Vec<String>,
pub summary: String,
@@ -144,6 +225,8 @@ pub struct ReviewFinding {
pub struct MergeRequestReview {
pub attempt_id: String,
pub revision_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub merge_result_id: Option<String>,
pub decision: ReviewDecision,
pub body: String,
pub findings: Vec<ReviewFinding>,
@@ -155,17 +238,46 @@ pub struct MergeRequestReview {
pub submitted_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MergeResult {
pub merge_result_id: String,
pub revision_id: String,
pub target_commit: String,
pub source_commit: String,
pub result_commit: String,
pub strategy: MergeStrategy,
pub resolution: MergeResolution,
pub created_by_runtime_id: String,
pub created_by_worker_id: String,
pub created_at: String,
pub operation_id: String,
pub validated_at: String,
pub target_status: MergeResultTargetStatus,
pub review_status: ReviewStatus,
pub current_review: Option<MergeRequestReview>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MergeRequest {
pub merge_request_id: String,
pub workspace_id: String,
pub ticket_id: String,
pub repository_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub target_ref_selector: Option<String>,
pub target_status: MergeRequestTargetStatus,
#[serde(skip_serializing_if = "Option::is_none")]
pub observed_target_commit: Option<String>,
pub state: MergeRequestState,
pub lifecycle_generation: u64,
pub current_revision: MergeRequestRevision,
pub review_status: ReviewStatus,
pub current_review: Option<MergeRequestReview>,
pub merge_results: Vec<MergeResult>,
#[serde(skip_serializing_if = "Option::is_none")]
pub current_merge_result: Option<MergeResult>,
#[serde(skip_serializing_if = "Option::is_none")]
pub applied_merge_result: Option<MergeResult>,
pub created_at: String,
pub updated_at: String,
pub merged_by_account_id: Option<String>,
@@ -177,6 +289,7 @@ pub struct OpenMergeRequest {
pub merge_request_id: String,
pub ticket_id: String,
pub repository_id: String,
pub target_ref_selector: String,
pub revision: MergeRequestRevision,
pub authenticated_runtime_id: String,
pub authenticated_worker_id: String,
@@ -206,6 +319,7 @@ pub struct RegisterReviewAttempt {
pub attempt_id: String,
pub ticket_id: String,
pub revision_id: String,
pub merge_result_id: Option<String>,
pub parent_assignment_id: String,
pub parent_runtime_id: String,
pub parent_worker_id: String,
@@ -219,6 +333,7 @@ pub struct RegisterReviewAttempt {
pub struct SubmitReview {
pub ticket_id: String,
pub revision_id: String,
pub merge_result_id: Option<String>,
pub capability_token: String,
pub decision: ReviewDecision,
pub body: String,
@@ -226,6 +341,28 @@ pub struct SubmitReview {
pub now: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RecordMergeResult {
pub merge_result_id: String,
pub ticket_id: String,
pub expected_revision_id: String,
pub target_commit: String,
pub source_commit: String,
pub result_commit: String,
pub strategy: MergeStrategy,
pub resolution: MergeResolution,
pub operation_id: String,
pub actor_runtime_id: String,
pub actor_worker_id: String,
pub created_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RecordMergeResultOutcome {
pub merge_result: MergeResult,
pub replayed: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompleteMergeRequest {
pub operation_id: String,
@@ -251,8 +388,12 @@ pub struct MergeRequestReadiness {
pub ticket_id: String,
pub merge_request_id: String,
pub revision_id: String,
pub target_ref_selector: Option<String>,
pub observed_target_commit: Option<String>,
pub ready: bool,
pub review_status: ReviewStatus,
pub merge_result_id: Option<String>,
pub merge_result_review_status: Option<ReviewStatus>,
pub blockers: Vec<String>,
}
@@ -327,33 +468,87 @@ impl SqliteMergeRequestStore {
}
pub fn show_for_ticket(&self, ticket_id: &str) -> Result<Option<MergeRequest>> {
self.show_for_ticket_with_target(ticket_id, None)
}
pub fn show_for_ticket_with_target(
&self,
ticket_id: &str,
observed_target_commit: Option<&str>,
) -> Result<Option<MergeRequest>> {
nonempty("ticket_id", ticket_id)?;
let conn = self.connect()?;
verify(&conn)?;
load_merge_request(&conn, &self.workspace_id, ticket_id)
let mut mr = load_merge_request(&conn, &self.workspace_id, ticket_id)?;
if let Some(mr) = mr.as_mut() {
apply_target_observation(mr, observed_target_commit);
}
Ok(mr)
}
pub fn readiness_for_ticket(&self, ticket_id: &str) -> Result<MergeRequestReadiness> {
self.readiness_for_ticket_with_target(ticket_id, None)
}
pub fn readiness_for_ticket_with_target(
&self,
ticket_id: &str,
observed_target_commit: Option<&str>,
) -> Result<MergeRequestReadiness> {
let mr = self
.show_for_ticket(ticket_id)?
.show_for_ticket_with_target(ticket_id, observed_target_commit)?
.ok_or_else(|| MergeRequestError::NotFound(ticket_id.to_string()))?;
let mut blockers = Vec::new();
if mr.state != MergeRequestState::Open {
blockers.push(format!("merge request is {}", mr.state.as_str()));
}
match mr.review_status {
ReviewStatus::Pending => blockers.push("current revision has no review result".into()),
ReviewStatus::Pending => {
blockers.push("current source revision has no review result".into())
}
ReviewStatus::ChangesRequested => {
blockers.push("current revision has request_changes".into())
blockers.push("current source revision has request_changes".into())
}
ReviewStatus::Approved => {}
}
if mr.target_status != MergeRequestTargetStatus::Known || observed_target_commit.is_none() {
blockers.push("merge request target is unknown or could not be resolved".into());
}
let current_result = mr
.current_merge_result
.as_ref()
.or(mr.applied_merge_result.as_ref());
if current_result.is_none() && observed_target_commit.is_some() {
if mr
.merge_results
.iter()
.any(|result| result.target_status == MergeResultTargetStatus::Stale)
{
blockers.push("target moved after the latest MergeResult was recorded".into());
} else {
blockers.push(
"current source revision has no validated MergeResult for the current target"
.into(),
);
}
}
if let Some(result) = current_result {
if matches!(result.strategy, MergeStrategy::Merge)
&& result.review_status != ReviewStatus::Approved
{
blockers.push("non-fast-forward MergeResult is not independently approved".into());
}
}
Ok(MergeRequestReadiness {
ticket_id: ticket_id.to_string(),
merge_request_id: mr.merge_request_id,
revision_id: mr.current_revision.revision_id,
target_ref_selector: mr.target_ref_selector,
observed_target_commit: mr.observed_target_commit,
ready: blockers.is_empty(),
review_status: mr.review_status,
merge_result_id: current_result.map(|result| result.merge_result_id.clone()),
merge_result_review_status: current_result.map(|result| result.review_status),
blockers,
})
}
@@ -364,6 +559,7 @@ impl SqliteMergeRequestStore {
("merge_request_id", input.merge_request_id.as_str()),
("ticket_id", input.ticket_id.as_str()),
("repository_id", input.repository_id.as_str()),
("target_ref_selector", input.target_ref_selector.as_str()),
("runtime_id", input.authenticated_runtime_id.as_str()),
("worker_id", input.authenticated_worker_id.as_str()),
] {
@@ -379,8 +575,8 @@ impl SqliteMergeRequestStore {
&input.authenticated_worker_id,
)?;
conn.execute(
"INSERT INTO merge_requests (workspace_id, merge_request_id, repository_id, state, lifecycle_generation, current_revision_id, created_at, updated_at) VALUES (?1, ?2, ?3, 'open', 1, ?4, ?5, ?5)",
params![self.workspace_id, input.merge_request_id, input.repository_id, input.revision.revision_id, input.now],
"INSERT INTO merge_requests (workspace_id, merge_request_id, repository_id, target_ref_selector, target_status, state, lifecycle_generation, current_revision_id, created_at, updated_at) VALUES (?1, ?2, ?3, ?4, 'known', 'open', 1, ?5, ?6, ?6)",
params![self.workspace_id, input.merge_request_id, input.repository_id, input.target_ref_selector, input.revision.revision_id, input.now],
).map_err(db)?;
conn.execute(
"INSERT INTO merge_request_ticket_relations (workspace_id,merge_request_id,ticket_id,relation_kind,created_at) VALUES (?1,?2,?3,'implements',?4)",
@@ -414,13 +610,13 @@ impl SqliteMergeRequestStore {
if input.revision.ordinal != current.current_revision.ordinal + 1 {
return Err(MergeRequestError::RevisionConflict(input.revision.revision_id.clone()));
}
let existing: Option<(String, String, String, String)> = conn.query_row(
"SELECT base_commit, head_commit, head_tree, diff_digest FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3",
let existing: Option<(String, String, String)> = conn.query_row(
"SELECT base_commit, head_commit, diff_digest FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3",
params![self.workspace_id, current.merge_request_id, input.revision.revision_id],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
).optional().map_err(db)?;
if let Some(existing) = existing {
if existing == (input.revision.base_commit.clone(), input.revision.head_commit.clone(), input.revision.head_tree.clone(), input.revision.diff_digest.clone()) {
if existing == (input.revision.base_commit.clone(), input.revision.head_commit.clone(), input.revision.diff_digest.clone()) {
return Ok(current);
}
return Err(MergeRequestError::RevisionConflict(input.revision.revision_id.clone()));
@@ -468,7 +664,19 @@ impl SqliteMergeRequestStore {
if mr.current_revision.revision_id != input.revision_id {
return Err(MergeRequestError::StaleRevision { expected: input.revision_id.clone(), current: mr.current_revision.revision_id });
}
validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id, &input.parent_runtime_id, &input.parent_worker_id)?;
if let Some(merge_result_id) = input.merge_result_id.as_deref() {
let exists: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM merge_request_merge_results WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3 AND merge_result_id=?4)",
params![self.workspace_id,mr.merge_request_id,input.revision_id,merge_result_id],
|row| row.get(0),
).map_err(db)?;
if !exists {
return Err(MergeRequestError::MergeResultNotFound(merge_result_id.into()));
}
validate_current_assignment_id(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id)?;
} else {
validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &input.parent_assignment_id, &input.parent_runtime_id, &input.parent_worker_id)?;
}
let effective_profile: Option<String> = conn.query_row(
"SELECT effective_profile FROM merge_request_reviewer_child_sessions WHERE workspace_id=?1 AND child_session_id=?2 AND parent_runtime_id=?3 AND parent_worker_id=?4",
params![self.workspace_id,input.child_session_id,input.parent_runtime_id,input.parent_worker_id],
@@ -478,8 +686,8 @@ impl SqliteMergeRequestStore {
return Err(MergeRequestError::InvalidReviewer);
}
conn.execute(
"INSERT INTO merge_request_review_attempts (workspace_id, attempt_id, merge_request_id, ticket_id, revision_id, lifecycle_generation, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, capability_token_sha256, status, created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,'open',?13)",
params![self.workspace_id, input.attempt_id, mr.merge_request_id, input.ticket_id, input.revision_id, mr.lifecycle_generation as i64, input.parent_assignment_id, input.parent_runtime_id, input.parent_worker_id, input.child_session_id, REVIEWER_PROFILE, token_hash(&input.capability_token), input.now],
"INSERT INTO merge_request_review_attempts (workspace_id, attempt_id, merge_request_id, ticket_id, revision_id, merge_result_id, lifecycle_generation, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, capability_token_sha256, status, created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,'open',?14)",
params![self.workspace_id, input.attempt_id, mr.merge_request_id, input.ticket_id, input.revision_id, input.merge_result_id, mr.lifecycle_generation as i64, input.parent_assignment_id, input.parent_runtime_id, input.parent_worker_id, input.child_session_id, REVIEWER_PROFILE, token_hash(&input.capability_token), input.now],
).map_err(|_| MergeRequestError::InvalidReviewAttempt)?;
Ok(())
})
@@ -505,12 +713,12 @@ impl SqliteMergeRequestStore {
validate_review_input(&input)?;
self.write(|conn| {
let token = token_hash(&input.capability_token);
let attempt: Option<(String,String,String,String,String,String,String,String,i64)> = conn.query_row(
"SELECT attempt_id, merge_request_id, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, status, lifecycle_generation FROM merge_request_review_attempts WHERE workspace_id=?1 AND ticket_id=?2 AND revision_id=?3 AND capability_token_sha256=?4",
params![self.workspace_id, input.ticket_id, input.revision_id, token],
|row| Ok((row.get(0)?,row.get(1)?,row.get(2)?,row.get(3)?,row.get(4)?,row.get(5)?,row.get(6)?,row.get(7)?,row.get(8)?)),
let attempt: Option<(String,String,Option<String>,String,String,String,String,String,String,i64)> = conn.query_row(
"SELECT attempt_id, merge_request_id, merge_result_id, parent_assignment_id, parent_runtime_id, parent_worker_id, child_session_id, child_effective_profile, status, lifecycle_generation FROM merge_request_review_attempts WHERE workspace_id=?1 AND ticket_id=?2 AND revision_id=?3 AND ((?4 IS NULL AND merge_result_id IS NULL) OR merge_result_id=?4) AND capability_token_sha256=?5",
params![self.workspace_id, input.ticket_id, input.revision_id, input.merge_result_id, token],
|row| Ok((row.get(0)?,row.get(1)?,row.get(2)?,row.get(3)?,row.get(4)?,row.get(5)?,row.get(6)?,row.get(7)?,row.get(8)?,row.get(9)?)),
).optional().map_err(db)?;
let Some((attempt_id, mr_id, assignment_id, runtime_id, worker_id, child_session_id, effective_profile, status, lifecycle_generation)) = attempt else {
let Some((attempt_id, mr_id, merge_result_id, assignment_id, runtime_id, worker_id, child_session_id, effective_profile, status, lifecycle_generation)) = attempt else {
return Err(MergeRequestError::InvalidReviewAttempt);
};
if status != "open" || effective_profile != REVIEWER_PROFILE || child_session_id == worker_id {
@@ -525,10 +733,14 @@ impl SqliteMergeRequestStore {
if mr.current_revision.revision_id != input.revision_id {
return Err(MergeRequestError::StaleRevision { expected: input.revision_id.clone(), current: mr.current_revision.revision_id });
}
validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &assignment_id, &runtime_id, &worker_id)?;
if merge_result_id.is_some() {
validate_current_assignment_id(conn, &self.workspace_id, &input.ticket_id, &assignment_id)?;
} else {
validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &assignment_id, &runtime_id, &worker_id)?;
}
conn.execute(
"INSERT INTO merge_request_reviews (workspace_id, attempt_id, merge_request_id, revision_id, decision, body, submitted_at) VALUES (?1,?2,?3,?4,?5,?6,?7)",
params![self.workspace_id, attempt_id, mr_id, input.revision_id, input.decision.as_str(), input.body, input.now],
"INSERT INTO merge_request_reviews (workspace_id, attempt_id, merge_request_id, revision_id, merge_result_id, decision, body, submitted_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8)",
params![self.workspace_id, attempt_id, mr_id, input.revision_id, merge_result_id, input.decision.as_str(), input.body, input.now],
).map_err(|_| MergeRequestError::InvalidReviewAttempt)?;
for (ordinal, finding) in input.findings.iter().enumerate() {
nonempty("finding.body", &finding.body)?;
@@ -545,6 +757,88 @@ impl SqliteMergeRequestStore {
})
}
pub fn record_merge_result(
&self,
input: RecordMergeResult,
) -> Result<RecordMergeResultOutcome> {
for (name, value) in [
("merge_result_id", input.merge_result_id.as_str()),
("ticket_id", input.ticket_id.as_str()),
("expected_revision_id", input.expected_revision_id.as_str()),
("target_commit", input.target_commit.as_str()),
("source_commit", input.source_commit.as_str()),
("result_commit", input.result_commit.as_str()),
("operation_id", input.operation_id.as_str()),
("actor_runtime_id", input.actor_runtime_id.as_str()),
("actor_worker_id", input.actor_worker_id.as_str()),
] {
nonempty(name, value)?;
}
if matches!(input.strategy, MergeStrategy::FastForward)
&& (input.result_commit != input.source_commit
|| !matches!(input.resolution, MergeResolution::None))
{
return Err(MergeRequestError::InvalidMergeResult(
"fast-forward result must equal the source commit and use resolution=none".into(),
));
}
if matches!(input.strategy, MergeStrategy::Merge)
&& matches!(input.resolution, MergeResolution::None)
{
return Err(MergeRequestError::InvalidMergeResult(
"merge strategy requires clean or conflicts_resolved resolution".into(),
));
}
let fingerprint = merge_result_fingerprint(&input);
self.write(|conn| {
if let Some((stored, merge_result_id, generation)) = conn
.query_row(
"SELECT r.operation_fingerprint,r.merge_result_id,mr.lifecycle_generation FROM merge_request_merge_results r JOIN merge_requests mr ON mr.workspace_id=r.workspace_id AND mr.merge_request_id=r.merge_request_id WHERE r.workspace_id=?1 AND r.operation_id=?2",
params![self.workspace_id,input.operation_id],
|row| Ok((row.get::<_,String>(0)?,row.get::<_,String>(1)?,row.get::<_,i64>(2)?)),
)
.optional()
.map_err(db)?
{
if stored != fingerprint {
return Err(MergeRequestError::MergeResultOperationConflict);
}
let merge_result = load_merge_result(
conn,
&self.workspace_id,
&merge_result_id,
generation as u64,
)?
.ok_or_else(|| MergeRequestError::MergeResultNotFound(merge_result_id.clone()))?;
return Ok(RecordMergeResultOutcome { merge_result, replayed: true });
}
let mr = load_merge_request(conn, &self.workspace_id, &input.ticket_id)?
.ok_or_else(|| MergeRequestError::NotFound(input.ticket_id.clone()))?;
ensure_open(&mr)?;
if mr.target_status != MergeRequestTargetStatus::Known {
return Err(MergeRequestError::UnknownTarget);
}
if mr.current_revision.revision_id != input.expected_revision_id {
return Err(MergeRequestError::StaleRevision {
expected: input.expected_revision_id.clone(),
current: mr.current_revision.revision_id,
});
}
if mr.current_revision.head_commit != input.source_commit {
return Err(MergeRequestError::InvalidMergeResult(
"source commit does not match the current source revision".into(),
));
}
conn.execute(
"INSERT INTO merge_request_merge_results (workspace_id,merge_result_id,merge_request_id,ticket_id,revision_id,target_commit,source_commit,result_commit,strategy,resolution,created_by_runtime_id,created_by_worker_id,created_at,operation_id,operation_fingerprint,validated_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?13)",
params![self.workspace_id,input.merge_result_id,mr.merge_request_id,input.ticket_id,input.expected_revision_id,input.target_commit,input.source_commit,input.result_commit,input.strategy.as_str(),input.resolution.as_str(),input.actor_runtime_id,input.actor_worker_id,input.created_at,input.operation_id,fingerprint],
).map_err(db)?;
let merge_result = load_merge_result(conn, &self.workspace_id, &input.merge_result_id, mr.lifecycle_generation)?
.ok_or_else(|| MergeRequestError::MergeResultNotFound(input.merge_result_id.clone()))?;
Ok(RecordMergeResultOutcome { merge_result, replayed: false })
})
}
pub fn complete(&self, input: CompleteMergeRequest) -> Result<CompletionOutcome> {
for (name, value) in [
("operation_id", input.operation_id.as_str()),
@@ -714,6 +1008,9 @@ fn migrate_transaction(conn: &Connection) -> Result<()> {
{
migrate_completion_authority_v8(conn)?;
}
if version < 9 && !column_exists(conn, "merge_requests", "target_ref_selector")? {
migrate_merge_result_authority_v9(conn)?;
}
if version < 1 {
conn.execute(
"INSERT INTO merge_request_schema_migrations(version) VALUES (1)",
@@ -764,6 +1061,54 @@ fn migrate_completion_authority_v8(conn: &Connection) -> Result<()> {
.map_err(db)
}
fn migrate_merge_result_authority_v9(conn: &Connection) -> Result<()> {
conn.pragma_update(None, "foreign_keys", "OFF")
.map_err(db)?;
let result = conn.execute_batch(
"ALTER TABLE merge_requests ADD COLUMN target_ref_selector TEXT;
ALTER TABLE merge_requests ADD COLUMN target_status TEXT NOT NULL DEFAULT 'unknown' CHECK(target_status IN ('known','unknown'));
CREATE TABLE merge_request_revisions_v9 (
workspace_id TEXT NOT NULL, merge_request_id TEXT NOT NULL, revision_id TEXT NOT NULL,
ordinal INTEGER NOT NULL, base_commit TEXT NOT NULL, head_commit TEXT NOT NULL,
diff_digest TEXT NOT NULL, summary TEXT NOT NULL, assignment_id TEXT NOT NULL, created_at TEXT NOT NULL,
PRIMARY KEY(workspace_id,merge_request_id,revision_id),
UNIQUE(workspace_id,merge_request_id,ordinal),
FOREIGN KEY(workspace_id,merge_request_id) REFERENCES merge_requests(workspace_id,merge_request_id) ON DELETE CASCADE
);
INSERT INTO merge_request_revisions_v9(
workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at
) SELECT workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at
FROM merge_request_revisions;
DROP TABLE merge_request_revisions;
ALTER TABLE merge_request_revisions_v9 RENAME TO merge_request_revisions;
CREATE TABLE merge_request_merge_results (
workspace_id TEXT NOT NULL, merge_result_id TEXT NOT NULL, merge_request_id TEXT NOT NULL,
ticket_id TEXT NOT NULL, revision_id TEXT NOT NULL, target_commit TEXT NOT NULL,
source_commit TEXT NOT NULL, result_commit TEXT NOT NULL,
strategy TEXT NOT NULL CHECK(strategy IN ('fast_forward','merge')),
resolution TEXT NOT NULL CHECK(resolution IN ('none','clean','conflicts_resolved')),
created_by_runtime_id TEXT NOT NULL, created_by_worker_id TEXT NOT NULL,
created_at TEXT NOT NULL, operation_id TEXT NOT NULL, operation_fingerprint TEXT NOT NULL,
validated_at TEXT NOT NULL,
PRIMARY KEY(workspace_id,merge_result_id),
UNIQUE(workspace_id,operation_id),
FOREIGN KEY(workspace_id,merge_request_id,revision_id)
REFERENCES merge_request_revisions(workspace_id,merge_request_id,revision_id),
FOREIGN KEY(workspace_id,ticket_id) REFERENCES typed_tickets(workspace_id,ticket_id)
);
CREATE INDEX merge_request_merge_results_current_idx
ON merge_request_merge_results(workspace_id,merge_request_id,revision_id,target_commit,created_at);
ALTER TABLE merge_request_review_attempts ADD COLUMN merge_result_id TEXT;
ALTER TABLE merge_request_reviews ADD COLUMN merge_result_id TEXT;",
);
let restore = conn.pragma_update(None, "foreign_keys", "ON").map_err(db);
result.map_err(db)?;
restore
}
pub fn verify(conn: &Connection) -> Result<()> {
let version: i64 = conn
.query_row(
@@ -786,6 +1131,7 @@ pub fn verify(conn: &Connection) -> Result<()> {
"merge_request_review_attempts",
"merge_request_reviews",
"merge_request_review_findings",
"merge_request_merge_results",
"merge_request_completion_operations",
] {
let present: Option<i64> = conn
@@ -809,6 +1155,8 @@ pub fn verify(conn: &Connection) -> Result<()> {
"workspace_id",
"merge_request_id",
"repository_id",
"target_ref_selector",
"target_status",
"state",
"lifecycle_generation",
"current_revision_id",
@@ -832,7 +1180,6 @@ pub fn verify(conn: &Connection) -> Result<()> {
"ordinal",
"base_commit",
"head_commit",
"head_tree",
"diff_digest",
"assignment_id",
] as &[_],
@@ -855,6 +1202,7 @@ pub fn verify(conn: &Connection) -> Result<()> {
"merge_request_id",
"ticket_id",
"revision_id",
"merge_result_id",
"lifecycle_generation",
"parent_assignment_id",
"parent_runtime_id",
@@ -872,10 +1220,29 @@ pub fn verify(conn: &Connection) -> Result<()> {
"attempt_id",
"merge_request_id",
"revision_id",
"merge_result_id",
"decision",
"body",
] as &[_],
),
(
"merge_request_merge_results",
&[
"workspace_id",
"merge_result_id",
"merge_request_id",
"ticket_id",
"revision_id",
"target_commit",
"source_commit",
"result_commit",
"strategy",
"resolution",
"operation_id",
"operation_fingerprint",
"validated_at",
] as &[_],
),
(
"merge_request_completion_operations",
&[
@@ -991,6 +1358,8 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result
"workspace_id",
"merge_request_id",
"repository_id",
"target_ref_selector",
"target_status",
"state",
"lifecycle_generation",
"current_revision_id",
@@ -1016,7 +1385,6 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result
"ordinal",
"base_commit",
"head_commit",
"head_tree",
"diff_digest",
"assignment_id",
],
@@ -1029,6 +1397,7 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result
"merge_request_reviews",
"merge_request_review_attempts",
"merge_request_reviewer_child_sessions",
"merge_request_merge_results",
"merge_request_completion_operations",
"merge_request_revision_paths",
"merge_request_ticket_relations",
@@ -1111,14 +1480,16 @@ fn load_merge_request(
workspace_id: &str,
ticket_id: &str,
) -> Result<Option<MergeRequest>> {
let row: Option<(String,String,String,String,i64,String,String,String,Option<String>,Option<String>)> = conn.query_row(
"SELECT mr.merge_request_id,rel.ticket_id,mr.repository_id,mr.state,mr.lifecycle_generation,mr.current_revision_id,mr.created_at,mr.updated_at,mr.merged_by_account_id,mr.merged_at FROM merge_requests mr JOIN merge_request_ticket_relations rel ON rel.workspace_id=mr.workspace_id AND rel.merge_request_id=mr.merge_request_id WHERE mr.workspace_id=?1 AND rel.ticket_id=?2 AND rel.relation_kind='implements' ORDER BY mr.updated_at DESC,mr.merge_request_id DESC LIMIT 1",
params![workspace_id,ticket_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?)),
let row: Option<(String,String,String,Option<String>,String,String,i64,String,String,String,Option<String>,Option<String>)> = conn.query_row(
"SELECT mr.merge_request_id,rel.ticket_id,mr.repository_id,mr.target_ref_selector,mr.target_status,mr.state,mr.lifecycle_generation,mr.current_revision_id,mr.created_at,mr.updated_at,mr.merged_by_account_id,mr.merged_at FROM merge_requests mr JOIN merge_request_ticket_relations rel ON rel.workspace_id=mr.workspace_id AND rel.merge_request_id=mr.merge_request_id WHERE mr.workspace_id=?1 AND rel.ticket_id=?2 AND rel.relation_kind='implements' ORDER BY mr.updated_at DESC,mr.merge_request_id DESC LIMIT 1",
params![workspace_id,ticket_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?,r.get(10)?,r.get(11)?)),
).optional().map_err(db)?;
let Some((
mr_id,
ticket_id,
repository_id,
target_ref_selector,
target_status,
state,
generation,
revision_id,
@@ -1137,16 +1508,24 @@ fn load_merge_request(
Some(ReviewDecision::RequestChanges) => ReviewStatus::ChangesRequested,
None => ReviewStatus::Pending,
};
let merge_results =
load_merge_results(conn, workspace_id, &mr_id, &revision_id, generation as u64)?;
Ok(Some(MergeRequest {
merge_request_id: mr_id,
workspace_id: workspace_id.into(),
ticket_id,
repository_id,
target_ref_selector,
target_status: MergeRequestTargetStatus::parse(&target_status),
observed_target_commit: None,
state: MergeRequestState::parse(&state),
lifecycle_generation: generation as u64,
current_revision: revision,
review_status,
current_review,
merge_results,
current_merge_result: None,
applied_merge_result: None,
created_at,
updated_at,
merged_by_account_id,
@@ -1161,8 +1540,8 @@ fn load_revision(
revision_id: &str,
) -> Result<MergeRequestRevision> {
let mut revision: MergeRequestRevision = conn.query_row(
"SELECT revision_id,ordinal,base_commit,head_commit,head_tree,diff_digest,summary,assignment_id,created_at FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3",
params![workspace_id,mr_id,revision_id], |r| Ok(MergeRequestRevision { revision_id:r.get(0)?, ordinal:r.get::<_,i64>(1)? as u64, base_commit:r.get(2)?, head_commit:r.get(3)?, head_tree:r.get(4)?, diff_digest:r.get(5)?, changed_paths:Vec::new(), summary:r.get(6)?, assignment_id:r.get(7)?, created_at:r.get(8)? }),
"SELECT revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at FROM merge_request_revisions WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3",
params![workspace_id,mr_id,revision_id], |r| Ok(MergeRequestRevision { revision_id:r.get(0)?, ordinal:r.get::<_,i64>(1)? as u64, base_commit:r.get(2)?, head_commit:r.get(3)?, diff_digest:r.get(4)?, changed_paths:Vec::new(), summary:r.get(5)?, assignment_id:r.get(6)?, created_at:r.get(7)? }),
).map_err(db)?;
let mut statement = conn.prepare("SELECT path FROM merge_request_revision_paths WHERE workspace_id=?1 AND merge_request_id=?2 AND revision_id=?3 ORDER BY ordinal").map_err(db)?;
revision.changed_paths = statement
@@ -1179,7 +1558,7 @@ fn insert_revision(
mr_id: &str,
revision: &MergeRequestRevision,
) -> Result<()> {
conn.execute("INSERT INTO merge_request_revisions (workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,head_tree,diff_digest,summary,assignment_id,created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11)", params![workspace_id,mr_id,revision.revision_id,revision.ordinal as i64,revision.base_commit,revision.head_commit,revision.head_tree,revision.diff_digest,revision.summary,revision.assignment_id,revision.created_at]).map_err(db)?;
conn.execute("INSERT INTO merge_request_revisions (workspace_id,merge_request_id,revision_id,ordinal,base_commit,head_commit,diff_digest,summary,assignment_id,created_at) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10)", params![workspace_id,mr_id,revision.revision_id,revision.ordinal as i64,revision.base_commit,revision.head_commit,revision.diff_digest,revision.summary,revision.assignment_id,revision.created_at]).map_err(db)?;
for (ordinal, path) in revision.changed_paths.iter().enumerate() {
conn.execute("INSERT INTO merge_request_revision_paths (workspace_id,merge_request_id,revision_id,ordinal,path) VALUES (?1,?2,?3,?4,?5)", params![workspace_id,mr_id,revision.revision_id,ordinal as i64,path]).map_err(db)?;
}
@@ -1193,7 +1572,7 @@ fn load_latest_review(
revision_id: &str,
generation: i64,
) -> Result<Option<MergeRequestReview>> {
let attempt: Option<String> = conn.query_row("SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_request_id=?2 AND r.revision_id=?3 AND a.lifecycle_generation=?4 ORDER BY r.submitted_at DESC, r.attempt_id DESC LIMIT 1", params![workspace_id,mr_id,revision_id,generation], |r| r.get(0)).optional().map_err(db)?;
let attempt: Option<String> = conn.query_row("SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_request_id=?2 AND r.revision_id=?3 AND r.merge_result_id IS NULL AND a.lifecycle_generation=?4 ORDER BY r.submitted_at DESC, r.attempt_id DESC LIMIT 1", params![workspace_id,mr_id,revision_id,generation], |r| r.get(0)).optional().map_err(db)?;
match attempt {
Some(id) => load_review(conn, workspace_id, &id),
None => Ok(None),
@@ -1205,12 +1584,13 @@ fn load_review(
workspace_id: &str,
attempt_id: &str,
) -> Result<Option<MergeRequestReview>> {
let row: Option<(String,String,String,String,String,String,String,String,String)> = conn.query_row(
"SELECT r.revision_id,r.decision,r.body,a.parent_assignment_id,a.parent_runtime_id,a.parent_worker_id,a.child_session_id,a.child_effective_profile,r.submitted_at FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.attempt_id=?2",
params![workspace_id,attempt_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?)),
let row: Option<(String,Option<String>,String,String,String,String,String,String,String,String)> = conn.query_row(
"SELECT r.revision_id,r.merge_result_id,r.decision,r.body,a.parent_assignment_id,a.parent_runtime_id,a.parent_worker_id,a.child_session_id,a.child_effective_profile,r.submitted_at FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.attempt_id=?2",
params![workspace_id,attempt_id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?,r.get(6)?,r.get(7)?,r.get(8)?,r.get(9)?)),
).optional().map_err(db)?;
let Some((
revision_id,
merge_result_id,
decision,
body,
assignment,
@@ -1240,6 +1620,7 @@ fn load_review(
Ok(Some(MergeRequestReview {
attempt_id: attempt_id.into(),
revision_id,
merge_result_id,
decision: ReviewDecision::parse(&decision),
body,
findings,
@@ -1272,6 +1653,15 @@ fn validate_current_implementation_assignment(
Ok(())
}
fn validate_current_assignment_id(
conn: &Connection,
workspace_id: &str,
ticket_id: &str,
assignment_id: &str,
) -> Result<()> {
validate_current_implementation_assignment(conn, workspace_id, ticket_id, assignment_id)
}
fn validate_current_assignment(
conn: &Connection,
workspace_id: &str,
@@ -1318,7 +1708,6 @@ fn validate_revision(revision: &MergeRequestRevision) -> Result<()> {
("revision_id", revision.revision_id.as_str()),
("base_commit", revision.base_commit.as_str()),
("head_commit", revision.head_commit.as_str()),
("head_tree", revision.head_tree.as_str()),
("diff_digest", revision.diff_digest.as_str()),
("assignment_id", revision.assignment_id.as_str()),
] {
@@ -1354,6 +1743,144 @@ fn validate_revision(revision: &MergeRequestRevision) -> Result<()> {
Ok(())
}
fn merge_result_fingerprint(input: &RecordMergeResult) -> String {
token_hash(&format!(
"{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}\0{}",
input.ticket_id,
input.expected_revision_id,
input.target_commit,
input.source_commit,
input.result_commit,
input.strategy.as_str(),
input.resolution.as_str(),
input.actor_runtime_id,
input.actor_worker_id,
))
}
fn apply_target_observation(mr: &mut MergeRequest, observed_target_commit: Option<&str>) {
mr.observed_target_commit = observed_target_commit.map(str::to_owned);
for result in &mut mr.merge_results {
result.target_status = match observed_target_commit {
Some(commit) if result.target_commit == commit => MergeResultTargetStatus::Current,
Some(commit) if result.result_commit == commit => MergeResultTargetStatus::Applied,
Some(_) => MergeResultTargetStatus::Stale,
None => MergeResultTargetStatus::Unknown,
};
}
mr.current_merge_result = observed_target_commit.and_then(|commit| {
mr.merge_results
.iter()
.rev()
.find(|result| result.target_commit == commit)
.cloned()
});
mr.applied_merge_result = observed_target_commit.and_then(|commit| {
mr.merge_results
.iter()
.rev()
.find(|result| result.result_commit == commit)
.cloned()
});
}
fn load_merge_result(
conn: &Connection,
workspace_id: &str,
merge_result_id: &str,
generation: u64,
) -> Result<Option<MergeResult>> {
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<Vec<MergeResult>> {
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::<std::result::Result<Vec<_>, _>>()
.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<Option<MergeRequestReview>> {
let attempt: Option<String> = conn.query_row(
"SELECT r.attempt_id FROM merge_request_reviews r JOIN merge_request_review_attempts a ON a.workspace_id=r.workspace_id AND a.attempt_id=r.attempt_id WHERE r.workspace_id=?1 AND r.merge_result_id=?2 AND a.lifecycle_generation=?3 ORDER BY r.submitted_at DESC,r.attempt_id DESC LIMIT 1",
params![workspace_id,merge_result_id,generation as i64], |row| row.get(0),
).optional().map_err(db)?;
attempt
.map(|attempt| load_review(conn, workspace_id, &attempt))
.transpose()
.map(|review| review.flatten())
}
fn validate_review_input(input: &SubmitReview) -> Result<()> {
if input.body.len() > MAX_REVIEW_BODY_BYTES {
return Err(MergeRequestError::TooLarge {