diff --git a/crates/merge-request/src/lib.rs b/crates/merge-request/src/lib.rs index 22844355..31c847d3 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 = 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 = 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 { diff --git a/crates/merge-request/tests/store.rs b/crates/merge-request/tests/store.rs index 8763799c..26e6743d 100644 --- a/crates/merge-request/tests/store.rs +++ b/crates/merge-request/tests/store.rs @@ -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, Option, 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), diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 8a0aac22..9e0f3042 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -3844,28 +3844,21 @@ async fn scoped_complete_merge_request( let workspace_id = parse_workspace_id(&workspace_id)?; require_workspace_access(&workspace_id, &api)?; let source = authenticate_worker_mutation_source(&api, &workspace_id, &headers)?; + require_online_workspace_orchestrator_source(&api, &source)?; let assignment = api .store .get_current_ticket_worker_assignment(&workspace_id, &ticket_id)? .ok_or_else(|| { Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into()) })?; - if assignment.worker.runtime_id != source.runtime_id - || assignment.worker.worker_id != source.worker_id - { - return Err(Error::TicketAssignmentConflict( - "authenticated Worker is not the current Ticket assignee".into(), - ) - .into()); - } let outcome = merge_request_store(&api, &workspace_id)?.complete( merge_request::CompleteMergeRequest { operation_id: input.operation_id, ticket_id, expected_revision_id: input.expected_revision_id, - assignment_id: assignment.assignment_id, - authenticated_runtime_id: source.runtime_id, - authenticated_worker_id: source.worker_id, + 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), }, )?; @@ -4826,17 +4819,42 @@ fn bounded_orchestrator_attention_text(input: &str, max_chars: usize) -> String output } +fn require_online_workspace_orchestrator_source( + api: &WorkspaceApi, + source: &WorkerMutationSource, +) -> Result<()> { + let orchestrator = find_online_workspace_orchestrator(api).ok_or_else(|| { + Error::TicketAssignmentConflict( + "Workspace has no current online Workspace Orchestrator".into(), + ) + })?; + if orchestrator.worker != *source { + return Err(Error::TicketAssignmentConflict( + "Merge Request completion requires the current online Workspace Orchestrator".into(), + )); + } + Ok(()) +} + +fn find_online_workspace_orchestrator(api: &WorkspaceApi) -> Option { + api.runtime + .list_workers(1000) + .items + .into_iter() + .find(|worker| { + worker.singleton_key.as_deref() + == Some(crate::hosts::WORKSPACE_ORCHESTRATOR_SINGLETON_KEY) + && worker.workspace.workspace_id.as_deref() + == Some(api.config.workspace_id.as_str()) + && matches!(worker.state.as_str(), "idle" | "running" | "paused") + }) +} + fn find_workspace_orchestrator(api: &WorkspaceApi) -> Option { let is_orchestrator = |worker: &WorkerSummary| { worker.singleton_key.as_deref() == Some(crate::hosts::WORKSPACE_ORCHESTRATOR_SINGLETON_KEY) }; - if let Some(worker) = api - .runtime - .list_workers(1000) - .items - .into_iter() - .find(is_orchestrator) - { + if let Some(worker) = find_online_workspace_orchestrator(api) { return Some(worker); } for runtime in api.runtime.list_runtimes(1000).items { @@ -12012,6 +12030,226 @@ mod tests { assert!(matches!(error, Error::WorkerSourceIdentity(_))); } + #[tokio::test] + async fn merge_request_completion_authority_requires_current_online_orchestrator() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let workspace_id = api.config.workspace_id.clone(); + let Json(generic) = create_workspace_worker( + State(api.clone()), + HeaderMap::new(), + Json(CreateWorkspaceWorkerRequest { + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + display_name: "Generic Worker".to_string(), + profile: Some("builtin:coder".to_string()), + ticket_assignment: None, + initial_submit: Vec::new(), + working_directory: None, + }), + ) + .await + .unwrap(); + + assert!(matches!( + require_online_workspace_orchestrator_source(&api, &generic.worker_ref), + Err(Error::TicketAssignmentConflict(_)) + )); + + let Json(started) = scoped_start_workspace_orchestrator( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + ) + .await + .unwrap(); + let orchestrator = started.worker.unwrap().worker; + require_online_workspace_orchestrator_source(&api, &orchestrator).unwrap(); + assert!(matches!( + require_online_workspace_orchestrator_source(&api, &generic.worker_ref), + Err(Error::TicketAssignmentConflict(_)) + )); + + api.runtime + .stop_worker( + &orchestrator, + WorkerLifecycleRequest { + reason: Some("completion authority regression test".into()), + ticket_assignment: None, + }, + ) + .unwrap(); + assert!(find_workspace_orchestrator(&api).is_some()); + assert!(find_online_workspace_orchestrator(&api).is_none()); + assert!(matches!( + require_online_workspace_orchestrator_source(&api, &orchestrator), + Err(Error::TicketAssignmentConflict(_)) + )); + } + + #[tokio::test] + async fn merge_request_completion_endpoint_rejects_coder_and_accepts_orchestrator() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let workspace_id = api.config.workspace_id.clone(); + let backend = browser_ticket_backend(&api).unwrap(); + let mut input = ticket::NewTicket::new("Orchestrator completion authority"); + input.workflow_state = Some(TicketWorkflowState::InProgress); + let ticket = backend.create(input).unwrap(); + let Json(coder) = create_workspace_worker( + State(api.clone()), + HeaderMap::new(), + Json(CreateWorkspaceWorkerRequest { + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + display_name: "Assigned Coder".to_string(), + profile: Some("builtin:coder".to_string()), + ticket_assignment: Some(CreateWorkspaceWorkerTicketAssignmentRequest { + ticket_id: ticket.id.clone(), + operation_id: "completion-coder-assignment".to_string(), + }), + initial_submit: vec![Segment::Flow { + selector: "builtin:coder-review".to_string(), + }], + working_directory: None, + }), + ) + .await + .unwrap(); + let assignment = api + .store + .get_current_ticket_worker_assignment(&workspace_id, &ticket.id) + .unwrap() + .unwrap(); + let mr_store = merge_request_store(&api, &workspace_id).unwrap(); + mr_store + .open_merge_request(merge_request::OpenMergeRequest { + merge_request_id: "MR-server-completion".into(), + ticket_id: ticket.id.clone(), + repository_id: TEST_REPOSITORY_ID.into(), + revision: merge_request::MergeRequestRevision { + revision_id: "V1".into(), + ordinal: 1, + base_commit: "base".into(), + head_commit: "head".into(), + head_tree: "tree".into(), + diff_digest: "sha256:diff".into(), + changed_paths: vec!["src/lib.rs".into()], + summary: "approved revision".into(), + assignment_id: assignment.assignment_id.clone(), + created_at: "t1".into(), + }, + authenticated_runtime_id: coder.worker_ref.runtime_id.clone(), + authenticated_worker_id: coder.worker_ref.worker_id.clone(), + now: "t1".into(), + }) + .unwrap(); + mr_store + .register_reviewer_child_session(merge_request::RegisterReviewerChildSession { + parent_runtime_id: coder.worker_ref.runtime_id.clone(), + parent_worker_id: coder.worker_ref.worker_id.clone(), + child_session_id: "reviewer-child".into(), + now: "t2".into(), + }) + .unwrap(); + mr_store + .register_review_attempt(merge_request::RegisterReviewAttempt { + attempt_id: "attempt".into(), + ticket_id: ticket.id.clone(), + revision_id: "V1".into(), + 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(), + child_session_id: "reviewer-child".into(), + capability_token: "review-token".into(), + now: "t2".into(), + }) + .unwrap(); + mr_store + .submit_review(merge_request::SubmitReview { + ticket_id: ticket.id.clone(), + revision_id: "V1".into(), + capability_token: "review-token".into(), + decision: merge_request::ReviewDecision::Approve, + body: "approved".into(), + findings: Vec::new(), + now: "t3".into(), + }) + .unwrap(); + + let worker_headers = |worker: &RuntimeWorkerRef| { + let mut headers = HeaderMap::new(); + headers.insert( + "x-yoi-runtime-id", + axum::http::HeaderValue::from_str(&worker.runtime_id).unwrap(), + ); + headers.insert( + "x-yoi-worker-id", + axum::http::HeaderValue::from_str(&worker.worker_id).unwrap(), + ); + headers + }; + let request = || CompleteMergeRequestRequest { + operation_id: "complete-operation".into(), + expected_revision_id: "V1".into(), + }; + let coder_error = scoped_complete_merge_request( + State(api.clone()), + worker_headers(&coder.worker_ref), + AxumPath((workspace_id.clone(), ticket.id.clone())), + Json(request()), + ) + .await + .unwrap_err(); + assert!(matches!( + coder_error.error, + Error::TicketAssignmentConflict(_) + )); + + let Json(started) = scoped_start_workspace_orchestrator( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + ) + .await + .unwrap(); + let orchestrator = started.worker.unwrap().worker; + let Json(completed) = scoped_complete_merge_request( + State(api.clone()), + worker_headers(&orchestrator), + AxumPath((workspace_id, ticket.id.clone())), + Json(request()), + ) + .await + .unwrap(); + assert!(!completed.replayed); + assert_eq!( + backend + .show(ticket.id.clone().into()) + .unwrap() + .meta + .workflow_state, + TicketWorkflowState::Done + ); + let conn = rusqlite::Connection::open(&api.config.database_path).unwrap(); + let actor: String = conn + .query_row( + "SELECT author FROM typed_ticket_events WHERE workspace_id=?1 AND ticket_id=?2 AND kind='state_changed'", + rusqlite::params![api.config.workspace_id, ticket.id], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + actor, + format!( + "worker:{}:{}", + orchestrator.runtime_id, orchestrator.worker_id + ) + ); + } + #[tokio::test] async fn production_profile_backend_launches_and_restores_workspace_orchestrator() { let workspace = tempfile::tempdir().unwrap();