server: authorize orchestrator MR completion

This commit is contained in:
2026-08-14 03:14:08 +09:00
parent d039359386
commit 38dad4e865
3 changed files with 451 additions and 48 deletions
+92 -18
View File
@@ -11,7 +11,7 @@ use std::path::{Path, PathBuf};
use std::time::Duration;
use thiserror::Error;
const SCHEMA_VERSION: i64 = 7;
const SCHEMA_VERSION: i64 = 8;
const REVIEWER_PROFILE: &str = "builtin:reviewer";
const MAX_SUMMARY_BYTES: usize = 16 * 1024;
const MAX_REVIEW_BODY_BYTES: usize = 64 * 1024;
@@ -231,9 +231,9 @@ pub struct CompleteMergeRequest {
pub operation_id: String,
pub ticket_id: String,
pub expected_revision_id: String,
pub assignment_id: String,
pub authenticated_runtime_id: String,
pub authenticated_worker_id: String,
pub implementation_assignment_id: String,
pub completion_actor_runtime_id: String,
pub completion_actor_worker_id: String,
pub now: String,
}
@@ -550,6 +550,18 @@ impl SqliteMergeRequestStore {
("operation_id", input.operation_id.as_str()),
("ticket_id", input.ticket_id.as_str()),
("revision_id", input.expected_revision_id.as_str()),
(
"implementation_assignment_id",
input.implementation_assignment_id.as_str(),
),
(
"completion_actor_runtime_id",
input.completion_actor_runtime_id.as_str(),
),
(
"completion_actor_worker_id",
input.completion_actor_worker_id.as_str(),
),
] {
nonempty(name, value)?;
}
@@ -566,17 +578,22 @@ impl SqliteMergeRequestStore {
}
} else {
conn.execute(
"INSERT INTO merge_request_completion_operations (workspace_id, operation_id, ticket_id, revision_id, assignment_id, fingerprint, status, created_at, updated_at) VALUES (?1,?2,?3,?4,?5,?6,'pending',?7,?7)",
params![self.workspace_id, input.operation_id, input.ticket_id, input.expected_revision_id, input.assignment_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, fingerprint, status, created_at, updated_at) VALUES (?1,?2,?3,?4,'workspace_orchestrator',?5,?6,?7,?8,'pending',?9,?9)",
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, 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,
)?;
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 });
}
validate_current_assignment(conn, &self.workspace_id, &input.ticket_id, &input.assignment_id, &input.authenticated_runtime_id, &input.authenticated_worker_id)?;
if mr.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",
@@ -671,6 +688,11 @@ pub fn migrate(conn: &Connection) -> Result<()> {
.map_err(db)?;
archive_incompatible_legacy_tables(conn, version)?;
conn.execute_batch(SCHEMA_V1).map_err(db)?;
if version < SCHEMA_VERSION
&& column_exists(conn, "merge_request_completion_operations", "assignment_id")?
{
migrate_completion_authority_v8(conn)?;
}
if version < 1 {
conn.execute(
"INSERT INTO merge_request_schema_migrations(version) VALUES (1)",
@@ -684,7 +706,7 @@ pub fn migrate(conn: &Connection) -> Result<()> {
// preserved and revalidated by the current typed store rather than rewritten.
if column_exists(conn, "merge_request_schema_migrations", "name")? {
conn.execute(
"INSERT OR IGNORE INTO merge_request_schema_migrations(version,name) VALUES (?1,'fresh_bounded_context_authority')",
"INSERT OR IGNORE INTO merge_request_schema_migrations(version,name) VALUES (?1,'separate_completion_authority')",
params![SCHEMA_VERSION],
).map_err(db)?;
} else {
@@ -698,6 +720,29 @@ pub fn migrate(conn: &Connection) -> Result<()> {
verify(conn)
}
fn migrate_completion_authority_v8(conn: &Connection) -> Result<()> {
conn.execute_batch(
"ALTER TABLE merge_request_completion_operations RENAME TO merge_request_completion_operations_v7;
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)
);
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,fingerprint,status,result_ticket_state,created_at,updated_at
) SELECT workspace_id,operation_id,ticket_id,revision_id,'legacy_assigned_coder',assignment_id,
NULL,NULL,fingerprint,status,result_ticket_state,created_at,updated_at
FROM merge_request_completion_operations_v7;
DROP TABLE merge_request_completion_operations_v7;",
)
.map_err(db)
}
pub fn verify(conn: &Connection) -> Result<()> {
let version: i64 = conn
.query_row(
@@ -817,7 +862,10 @@ pub fn verify(conn: &Connection) -> Result<()> {
"operation_id",
"ticket_id",
"revision_id",
"assignment_id",
"authority_kind",
"implementation_assignment_id",
"completion_actor_runtime_id",
"completion_actor_worker_id",
"fingerprint",
"status",
"result_ticket_state",
@@ -901,7 +949,9 @@ CREATE TABLE IF NOT EXISTS merge_request_review_findings (
);
CREATE TABLE IF NOT EXISTS 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')),
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)
@@ -1166,6 +1216,26 @@ fn load_review(
}))
}
fn validate_current_implementation_assignment(
conn: &Connection,
workspace_id: &str,
ticket_id: &str,
assignment_id: &str,
) -> Result<()> {
let current: Option<String> = conn
.query_row(
"SELECT assignment_id FROM ticket_current_worker_assignments WHERE workspace_id=?1 AND ticket_id=?2",
params![workspace_id, ticket_id],
|row| row.get(0),
)
.optional()
.map_err(db)?;
if current.as_deref() != Some(assignment_id) {
return Err(MergeRequestError::AssignmentMismatch);
}
Ok(())
}
fn validate_current_assignment(
conn: &Connection,
workspace_id: &str,
@@ -1187,16 +1257,20 @@ fn append_completion_event(
input: &CompleteMergeRequest,
) -> Result<()> {
let index:i64=conn.query_row("SELECT COALESCE(MAX(event_index),-1)+1 FROM typed_ticket_events WHERE workspace_id=?1 AND ticket_id=?2",params![workspace_id,input.ticket_id],|r|r.get(0)).map_err(db)?;
conn.execute("INSERT INTO typed_ticket_events (workspace_id,ticket_id,event_index,kind,author,at,from_state,to_state,heading,body) VALUES (?1,?2,?3,'state_changed',?4,?5,'inprogress','done','Merge Request completed',?6)",params![workspace_id,input.ticket_id,index,format!("worker:{}:{}",input.authenticated_runtime_id,input.authenticated_worker_id),input.now,format!("Approved immutable revision `{}` completed implementation.",input.expected_revision_id)]).map_err(db)?;
conn.execute("INSERT INTO typed_ticket_events (workspace_id,ticket_id,event_index,kind,author,at,from_state,to_state,heading,body) VALUES (?1,?2,?3,'state_changed',?4,?5,'inprogress','done','Merge Request completed',?6)",params![workspace_id,input.ticket_id,index,format!("worker:{}:{}",input.completion_actor_runtime_id,input.completion_actor_worker_id),input.now,format!("Approved immutable revision `{}` completed implementation.",input.expected_revision_id)]).map_err(db)?;
for (key, value) in [
("assignment_id", input.assignment_id.as_str()),
(
"implementation_assignment_id",
input.implementation_assignment_id.as_str(),
),
(
"merge_request_revision_id",
input.expected_revision_id.as_str(),
),
("operation_id", input.operation_id.as_str()),
("runtime_id", input.authenticated_runtime_id.as_str()),
("worker_id", input.authenticated_worker_id.as_str()),
("completion_authority", "workspace_orchestrator"),
("runtime_id", input.completion_actor_runtime_id.as_str()),
("worker_id", input.completion_actor_worker_id.as_str()),
] {
conn.execute("INSERT INTO typed_ticket_event_attributes (workspace_id,ticket_id,event_index,key,value) VALUES (?1,?2,?3,?4,?5)",params![workspace_id,input.ticket_id,index,key,value]).map_err(db)?;
}
@@ -1299,12 +1373,12 @@ fn token_hash(token: &str) -> String {
}
fn completion_fingerprint(input: &CompleteMergeRequest) -> String {
token_hash(&format!(
"{}\0{}\0{}\0{}\0{}",
"workspace_orchestrator\0{}\0{}\0{}\0{}\0{}",
input.ticket_id,
input.expected_revision_id,
input.assignment_id,
input.authenticated_runtime_id,
input.authenticated_worker_id
input.implementation_assignment_id,
input.completion_actor_runtime_id,
input.completion_actor_worker_id
))
}
fn db(error: rusqlite::Error) -> MergeRequestError {
+103 -12
View File
@@ -192,6 +192,51 @@ fn rejected_v6_schema_missing_diff_digest_is_archived_before_fresh_v7() {
}
}
#[test]
fn v7_completion_operations_are_preserved_as_legacy_assigned_coder_authority() {
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 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));\
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);
SqliteMergeRequestStore::open(&path, "ws-a").unwrap();
let conn = Connection::open(&path).unwrap();
let row: (String, String, Option<String>, Option<String>, String) = conn
.query_row(
"SELECT authority_kind,implementation_assignment_id,completion_actor_runtime_id,completion_actor_worker_id,fingerprint FROM merge_request_completion_operations WHERE operation_id='legacy-op'",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?, row.get(4)?)),
)
.unwrap();
assert_eq!(
row,
(
"legacy_assigned_coder".into(),
"A1".into(),
None,
None,
"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, 8);
}
#[test]
fn request_changes_new_revision_resets_and_exact_completion_replay_converges() {
let (_dir, store) = setup();
@@ -223,9 +268,9 @@ fn request_changes_new_revision_resets_and_exact_completion_replay_converges() {
operation_id: "OP1".into(),
ticket_id: "T1".into(),
expected_revision_id: "V2".into(),
assignment_id: "A1".into(),
authenticated_runtime_id: "R1".into(),
authenticated_worker_id: "W1".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();
@@ -273,6 +318,32 @@ fn request_changes_new_revision_resets_and_exact_completion_replay_converges() {
.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]
@@ -340,9 +411,9 @@ fn concurrent_exact_completion_replays_commit_one_ticket_side_effect() {
operation_id: "OP-concurrent".into(),
ticket_id: "T1".into(),
expected_revision_id: "V1".into(),
assignment_id: "A1".into(),
authenticated_runtime_id: "R1".into(),
authenticated_worker_id: "W1".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();
@@ -374,7 +445,7 @@ fn concurrent_exact_completion_replays_commit_one_ticket_side_effect() {
}
#[test]
fn operation_key_mismatch_and_assignment_takeover_are_fenced() {
fn operation_key_mismatch_and_actor_or_assignment_change_are_fenced() {
let (_dir, store) = setup();
open(&store);
attempt(&store, "AT", "V1", "token", "child");
@@ -383,19 +454,39 @@ fn operation_key_mismatch_and_assignment_takeover_are_fenced() {
operation_id: "OP".into(),
ticket_id: "T1".into(),
expected_revision_id: "V1".into(),
assignment_id: "A1".into(),
authenticated_runtime_id: "R1".into(),
authenticated_worker_id: "W1".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();
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();
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();
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),