Merge branch 'refs/heads/develop' into work/T-565-worker-launch-rest-dto

This commit is contained in:
2026-09-03 14:54:19 +09:00
23 changed files with 791 additions and 412 deletions
+3 -4
View File
@@ -3772,7 +3772,6 @@ fn embedded_worker_status_label(status: EmbeddedWorkerStatus) -> &'static str {
EmbeddedWorkerStatus::Running => "running",
EmbeddedWorkerStatus::Paused => "paused",
EmbeddedWorkerStatus::Stopped => "stopped",
EmbeddedWorkerStatus::Cancelled => "cancelled",
}
}
@@ -5678,7 +5677,7 @@ mod tests {
json!({
"workers": [
worker_json_with_status("remote:primary", &worker_ids[0], "stopped"),
worker_json_with_status("remote:primary", &worker_ids[1], "cancelled"),
worker_json_with_status("remote:primary", &worker_ids[1], "running"),
worker_json_with_status("remote:primary", &worker_ids[2], "paused"),
worker_json_with_status("remote:primary", &worker_ids[3], "idle")
]
@@ -5717,11 +5716,11 @@ mod tests {
let workers = registry.list_workers(10);
assert_eq!(workers.items.len(), 4);
assert!(!workers.items[0].capabilities.can_stop);
assert!(!workers.items[1].capabilities.can_stop);
assert!(workers.items[1].capabilities.can_stop);
assert!(workers.items[2].capabilities.can_stop);
assert!(workers.items[3].capabilities.can_stop);
assert_eq!(workers.items[0].state, "stopped");
assert_eq!(workers.items[1].state, "cancelled");
assert_eq!(workers.items[1].state, "running");
assert_eq!(workers.items[2].state, "paused");
assert_eq!(workers.items[3].state, "idle");
+7 -15
View File
@@ -13167,22 +13167,13 @@ fn compensate_failed_worker_spawn(
let cancellation = api
.runtime
.cancel_worker(&worker.worker, lifecycle_request.clone());
let cancellation_accepted = cancellation
let stop = api.runtime.stop_worker(&worker.worker, lifecycle_request);
let stop_accepted = stop
.as_ref()
.is_ok_and(|result| result.state == WorkerOperationState::Accepted);
let stop = (!cancellation_accepted)
.then(|| api.runtime.stop_worker(&worker.worker, lifecycle_request));
let stop_accepted = stop.as_ref().is_some_and(|result| {
result
.as_ref()
.is_ok_and(|result| result.state == WorkerOperationState::Accepted)
});
let termination_detail = (!cancellation_accepted && !stop_accepted).then(|| {
let termination_detail = (!stop_accepted).then(|| {
let cancellation = lifecycle_failure_detail("cancel", &cancellation);
let stop = stop
.as_ref()
.map(|result| lifecycle_failure_detail("stop", result))
.unwrap_or_else(|| "stop was not attempted".to_string());
let stop = lifecycle_failure_detail("stop", &stop);
format!("{cancellation}; {stop}")
});
@@ -20256,7 +20247,7 @@ mod tests {
.iter()
.any(|assignment| assignment.role == "coder")
);
assert_eq!(api.runtime.worker(&worker).unwrap().state, "cancelled");
assert_eq!(api.runtime.worker(&worker).unwrap().state, "idle");
let Json(replayed) = scoped_cancel_ticket_implementation(State(api), path(), request())
.await
@@ -26620,13 +26611,14 @@ mod tests {
"embedded_worker_runtime"
);
let repository_id = test_repository_id(&api);
let workdir_id = "external-workdir";
api.store
.upsert_workdir_registry(&WorkdirRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
workdir_id: workdir_id.to_string(),
runtime_id: "external-workdir-runtime".to_string(),
repository_id: "main".to_string(),
repository_id,
creation_selector: None,
creation_ref: None,
creation_tree: None,
+54 -8
View File
@@ -7123,6 +7123,10 @@ fn migrate_repository_identity_to_keys(conn: &Connection) -> Result<()> {
stmt.query_map([], |row| row.get::<_, String>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?
};
let mut repository_reference_tables = REPOSITORY_REFERENCE_TABLES
.iter()
.map(|table| (*table).to_string())
.collect::<Vec<_>>();
for table in &tables {
let quoted = table.replace('"', "\"\"");
let pragma = format!("PRAGMA table_info(\"{quoted}\")");
@@ -7132,23 +7136,45 @@ fn migrate_repository_identity_to_keys(conn: &Connection) -> Result<()> {
.collect::<std::result::Result<Vec<_>, _>>()?
.iter()
.any(|column| column == "repository_id");
if has_repository_id
&& table != "repositories"
&& table != "legacy_repositories"
&& !REPOSITORY_REFERENCE_TABLES.contains(&table.as_str())
if !has_repository_id
|| table == "repositories"
|| table == "legacy_repositories"
|| REPOSITORY_REFERENCE_TABLES.contains(&table.as_str())
{
continue;
}
let pragma = format!("PRAGMA foreign_key_list(\"{quoted}\")");
let mut stmt = conn.prepare(&pragma)?;
let references_repository_authority = stmt
.query_map([], |row| {
Ok((
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, String>(4)?,
))
})?
.collect::<std::result::Result<Vec<_>, _>>()?
.iter()
.any(|(parent, from, to)| {
parent == "repositories" && from == "repository_id" && to == "repository_id"
});
if references_repository_authority {
repository_reference_tables.push(table.clone());
} else {
return Err(Error::Store(format!(
"repository identity migration does not recognize repository_id authority in table {table}"
)));
}
}
for table in REPOSITORY_REFERENCE_TABLES {
for table in &repository_reference_tables {
if !tables.iter().any(|candidate| candidate == table) {
continue;
}
let quoted = table.replace('"', "\"\"");
let sql = format!(
r#"SELECT COUNT(*)
FROM "{table}" AS child
FROM "{quoted}" AS child
LEFT JOIN repositories AS repository
ON repository.workspace_id = child.workspace_id
AND repository.repository_id = child.repository_id
@@ -7183,12 +7209,13 @@ fn migrate_repository_identity_to_keys(conn: &Connection) -> Result<()> {
)?;
}
for table in REPOSITORY_REFERENCE_TABLES {
for table in &repository_reference_tables {
if !tables.iter().any(|candidate| candidate == table) {
continue;
}
let quoted = table.replace('"', "\"\"");
let sql = format!(
r#"UPDATE "{table}" AS child
r#"UPDATE "{quoted}" AS child
SET repository_id = (
SELECT mapping.new_repository_id
FROM repository_identity_v50 AS mapping
@@ -12597,6 +12624,17 @@ INSERT INTO worker_registry (
'1', '1', 'local_path', '/repo-a', 1, 'sha256:a', 'unverified'),
('workspace-b', 'main', 'Legacy B', 'git', 'git', '/repo-b', 'develop',
'1', '1', 'local_path', '/repo-b', 1, 'sha256:b', 'unverified');
CREATE TABLE legacy_v5_merge_requests (
workspace_id TEXT NOT NULL,
merge_request_id TEXT NOT NULL,
repository_id TEXT NOT NULL,
PRIMARY KEY (workspace_id, merge_request_id),
FOREIGN KEY (workspace_id, repository_id)
REFERENCES repositories(workspace_id, repository_id)
);
INSERT INTO legacy_v5_merge_requests (
workspace_id, merge_request_id, repository_id
) VALUES ('workspace-a', 'legacy-mr-a', 'main');
INSERT INTO typed_tickets (
workspace_id, ticket_id, slug, title, status, kind, priority, body,
workflow_state, workflow_state_explicit, repository_id
@@ -12651,6 +12689,14 @@ INSERT INTO worker_registry (
assert_eq!(repositories[0].2, "main");
assert_eq!(repositories[1].2, "main");
assert_ne!(repositories[0].1, repositories[1].1);
let legacy_repository_id: String = conn
.query_row(
"SELECT repository_id FROM legacy_v5_merge_requests WHERE workspace_id = 'workspace-a' AND merge_request_id = 'legacy-mr-a'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(legacy_repository_id, repositories[0].1);
for (_, repository_id, _) in &repositories {
assert_eq!(Uuid::parse_str(repository_id).unwrap().get_version_num(), 7);
}