feat: merge durable Workdir updates from develop

# Conflicts:
#	docs/README.md
#	docs/design/durable-operations.md
#	web/workspace/deno.json
This commit is contained in:
2026-09-03 02:47:50 +09:00
43 changed files with 6467 additions and 524 deletions
+1
View File
@@ -30,6 +30,7 @@ pub mod server;
pub mod skills;
pub mod store;
pub mod workdir_create_operations;
mod workdir_removal;
pub mod worker_source;
pub mod workspace_catalog;
mod workspace_subscription;
@@ -12,11 +12,15 @@ use worker_runtime::config_bundle::{
ConfigBundle, ConfigBundleMetadata, ConfigBundleProvenance, ConfigProfileDescriptor,
};
use worker_runtime::profile_archive::{ProfileSourceArchive, ProfileSourceArchiveInput};
use workspace_api::{
Diagnostic, DiagnosticSeverity, ProfileSettingsResponse, UpdateWorkspaceMetadataRequest,
WorkspaceMetadataSettingsResponse, WorkspaceProfileSourceProvenance,
WorkspaceProfileSourceSummary, WorkspaceProfileSummary,
};
use crate::config_source::{
WorkspaceConfigSchemaProvider, WorkspaceConfigState, evaluate_workspace_config_state,
};
use crate::hosts::{DiagnosticSeverity, RuntimeDiagnostic};
use crate::{Error, Result};
const PROFILE_SCHEMA_SOURCE: &str = r#"{
@@ -427,81 +431,6 @@ fn build_virtual_profile_archive(
.map_err(|error| profile_validation_error("profile_source_archive_invalid", &error.to_string()))
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspaceMetadataSettingsResponse {
pub workspace_id: String,
pub display_name: String,
pub created_at: String,
pub revision: String,
pub source: String,
pub diagnostics: Vec<RuntimeDiagnostic>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct UpdateWorkspaceMetadataRequest {
pub display_name: String,
pub revision: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspaceMetadataMutationResponse {
pub workspace: WorkspaceMetadataSettingsResponse,
pub diagnostics: Vec<RuntimeDiagnostic>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ProfileSettingsResponse {
pub workspace_id: String,
pub registry_revision: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub config_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tree_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub projection_digest: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub default_profile: Option<String>,
pub profiles: Vec<WorkspaceProfileSummary>,
pub sources: Vec<WorkspaceProfileSourceSummary>,
pub diagnostics: Vec<RuntimeDiagnostic>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspaceProfileSummary {
pub profile_id: String,
pub selector: String,
pub label: String,
pub source_kind: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub profile_source_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
pub editable: bool,
pub is_default: bool,
pub diagnostics: Vec<RuntimeDiagnostic>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspaceProfileSourceSummary {
pub profile_source_id: String,
pub display_path: String,
pub kind: String,
pub content_type: String,
pub content_digest: String,
pub provenance: WorkspaceProfileSourceProvenance,
pub editable: bool,
pub revision: String,
pub size_bytes: u64,
pub diagnostics: Vec<RuntimeDiagnostic>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WorkspaceProfileSourceProvenance {
ProjectProfileSourceTree,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct WorkspaceIdentityFile {
@@ -804,8 +733,8 @@ fn diagnostic(
code: impl Into<String>,
severity: DiagnosticSeverity,
message: impl Into<String>,
) -> RuntimeDiagnostic {
RuntimeDiagnostic {
) -> Diagnostic {
Diagnostic {
code: code.into(),
severity,
message: message.into(),
File diff suppressed because it is too large Load Diff
+154 -22
View File
@@ -267,6 +267,11 @@ const MIGRATIONS: &[Migration] = &[
name: "require one account owner for every Workspace",
apply: require_workspace_account_owner,
},
Migration {
version: 49,
name: "create durable Workdir removal operations",
apply: crate::workdir_removal::create_workdir_removal_operations,
},
];
struct Migration {
@@ -4884,6 +4889,19 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
"Workdir {workdir_id} is not registered in Workspace {workspace_id}"
)));
}
let removal_pending: bool = tx.query_row(
r#"SELECT EXISTS(
SELECT 1 FROM workdir_removal_operations
WHERE workspace_id = ?1 AND workdir_id = ?2 AND state = 'pending'
)"#,
params![workspace_id, workdir_id],
|row| row.get(0),
)?;
if removal_pending {
return Err(Error::WorkdirAttachmentConflict(format!(
"Workdir {workdir_id} has a pending durable removal operation"
)));
}
let occupied: bool = tx.query_row(
r#"SELECT EXISTS(
SELECT 1 FROM worker_workdir_links
@@ -10085,6 +10103,120 @@ mod tests {
.unwrap();
}
#[test]
fn schema_v49_upgrades_persisted_v48_workdir_fixture() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("server-v48.db");
{
let conn = Connection::open(&path).unwrap();
configure_sqlite(&conn).unwrap();
apply_migrations_through(&conn, 48).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert!(!table_exists(&conn, "workdir_removal_operations").unwrap());
conn.execute_batch(
r#"
INSERT INTO accounts (
account_id, kind, handle, display_name, created_at, updated_at
) VALUES ('owner-account', 'user', 'owner', 'Owner', '1', '1');
INSERT INTO workspaces (
workspace_id, owner_account_id, display_name, state, created_at, updated_at
) VALUES ('workspace-a', 'owner-account', 'Workspace A', 'active', '1', '1');
INSERT INTO repositories (
workspace_id, repository_id, name, kind, provider, uri,
source_kind, source_uri, default_ref, source_revision,
source_fingerprint, observed_status, observed_at, created_at, updated_at
) VALUES (
'workspace-a', 'repository-a', 'Repository A', 'git', 'local', '/repo-a',
'local_path', '/repo-a', 'develop', 1,
'sha256:source-a', 'unverified', NULL, '1', '1'
);
INSERT INTO workdir_registry (
workspace_id, workdir_id, runtime_id, repository_id,
creation_selector, creation_ref, creation_tree,
current_selector, current_ref, current_tree,
observed_at_epoch_seconds, materialization_status, cleanliness,
created_at, updated_at
) VALUES (
'workspace-a', 'workdir-a', 'runtime-a', 'repository-a',
'refs/heads/develop', 'abc', 'tree-a',
'refs/heads/work', 'def', 'tree-b',
1, 'present', 'clean', '1', '1'
);
"#,
)
.unwrap();
}
let store = SqliteWorkspaceStore::open(&path).unwrap();
store
.with_conn(|conn| {
assert_eq!(current_schema_version(conn)?, 49);
assert!(table_exists(conn, "workdir_removal_operations")?);
let columns = table_columns(conn, "workdir_removal_operations")?;
for required in [
"workspace_id",
"operation_id",
"request_fingerprint",
"workdir_id",
"runtime_id",
"repository_id",
"materialization_fingerprint",
"source_actor",
"reason",
"state",
"attempt_count",
"retryable",
"disposition",
"failure_category",
"attempt_owner_pid",
"attempt_owner_start_marker",
"created_at",
"updated_at",
"completed_at",
] {
assert!(
columns.iter().any(|column| column == required),
"missing {required}"
);
}
let preserved: (String, String, String) = conn.query_row(
"SELECT workspace_id, repository_id, materialization_status FROM workdir_registry WHERE workdir_id='workdir-a'",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
)?;
assert_eq!(
preserved,
(
"workspace-a".to_string(),
"repository-a".to_string(),
"present".to_string(),
)
);
let foreign_key_failures: i64 = conn.query_row(
"SELECT count(*) FROM pragma_foreign_key_check",
[],
|row| row.get(0),
)?;
assert_eq!(foreign_key_failures, 0);
Ok(())
})
.unwrap();
let workdir = store
.get_workdir_registry("workspace-a", "workdir-a")
.unwrap()
.unwrap();
let intent = crate::workdir_removal::workdir_removal_intent(
&workdir,
"migration-test",
"remove migrated Workdir",
)
.unwrap();
let operation = store.reserve_workdir_removal_operation(&intent).unwrap();
assert_eq!(operation.workspace_id, "workspace-a");
assert_eq!(operation.working_directory_id, "workdir-a");
}
#[test]
fn schema_v44_migrates_repository_sources_without_promoting_legacy_auth_refs() {
let conn = Connection::open_in_memory().unwrap();
@@ -10118,7 +10250,7 @@ mod tests {
assign_explicit_test_workspace_owner(&conn);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let remote = conn
.query_row(
"SELECT source_kind, source_uri, source_revision, source_fingerprint, observed_status \
@@ -10197,7 +10329,7 @@ mod tests {
let before = std::fs::read(&path).unwrap();
let plan = SqliteWorkspaceStore::migration_plan(&path).unwrap();
assert_eq!(plan.current_schema_version, 36);
assert_eq!(plan.target_schema_version, 48);
assert_eq!(plan.target_schema_version, 49);
assert!(plan.migration_required);
assert_eq!(plan.worker_count, 1);
assert_eq!(plan.mappings[0].legacy_worker_id, 7);
@@ -10211,7 +10343,7 @@ mod tests {
store
.with_conn(|conn| {
assert!(table_exists(conn, "worker_diagnostics_archives")?);
assert_eq!(current_schema_version(conn)?, 48);
assert_eq!(current_schema_version(conn)?, 49);
Ok(())
})
.unwrap();
@@ -10351,7 +10483,7 @@ mod tests {
),
]
);
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let foreign_key_error: Option<String> = conn
.query_row("PRAGMA foreign_key_check", [], |row| row.get(0))
.optional()
@@ -10481,7 +10613,7 @@ INSERT INTO worker_orphan_diagnostics (
assign_explicit_test_workspace_owner(&conn);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
assert!(!table_exists(&conn, "worker_control_delegation_operations").unwrap());
let controller_worker_id: String = conn
.query_row(
@@ -10599,7 +10731,7 @@ INSERT INTO worker_orphan_diagnostics (
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap());
}
@@ -10618,7 +10750,7 @@ INSERT INTO worker_orphan_diagnostics (
assign_explicit_test_workspace_owner(&conn);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let settings = conn
.query_row(
"SELECT settings_revision, language FROM workspace_memory_settings \
@@ -10659,7 +10791,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
assert!(table_exists(&conn, "flow_sources").unwrap());
assert!(table_exists(&conn, "flow_source_revisions").unwrap());
assert!(!table_exists(&conn, "flow_instances").unwrap());
@@ -10727,7 +10859,7 @@ INSERT INTO worker_workdir_attachment_reservations (
assign_explicit_test_workspace_owner(&conn);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let repositories_sql: String = conn
.query_row(
"SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'",
@@ -10910,7 +11042,7 @@ INSERT INTO workdir_registry (
let db = dir.path().join("control-plane.sqlite");
let store = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 48);
assert_eq!(store.schema_version().await.unwrap(), 49);
assert!(
!store
.with_conn(|conn| table_exists(conn, "worker_workspace_credentials"))
@@ -10927,7 +11059,7 @@ INSERT INTO workdir_registry (
store.upsert_workspace(&record).await.unwrap();
let reopened = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(reopened.schema_version().await.unwrap(), 48);
assert_eq!(reopened.schema_version().await.unwrap(), 49);
assert_eq!(
reopened.get_workspace("local-dev").await.unwrap(),
Some(record)
@@ -11693,7 +11825,7 @@ INSERT INTO worker_registry (
let migrated = SqliteWorkspaceStore::open(&db_path).unwrap();
migrated
.with_conn(|conn| {
assert_eq!(current_schema_version(conn)?, 48);
assert_eq!(current_schema_version(conn)?, 49);
assert_eq!(
conn.query_row("PRAGMA foreign_keys", [], |row| row.get::<_, i64>(0))?,
1,
@@ -12044,7 +12176,7 @@ INSERT INTO worker_registry (
assert_eq!(current_schema_version(&conn).unwrap(), 44);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
assert!(table_exists(&conn, "workdir_create_operations").unwrap());
let columns = table_columns(&conn, "workdir_create_operations").unwrap();
for required in [
@@ -12071,7 +12203,7 @@ INSERT INTO worker_registry (
assert_eq!(current_schema_version(&conn).unwrap(), 45);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
for table in [
"repository_ssh_credentials",
"repository_ssh_credential_revisions",
@@ -12098,7 +12230,7 @@ INSERT INTO worker_registry (
assert_eq!(current_schema_version(&conn).unwrap(), 46);
apply_migrations(&conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let columns = table_columns(&conn, "workdir_create_operations").unwrap();
for required in [
"source_kind",
@@ -12132,13 +12264,13 @@ INSERT INTO worker_registry (
configure_sqlite(&conn).unwrap();
apply_migrations(&conn).unwrap();
conn.execute(
"INSERT INTO __yoi_schema_migrations (version, name) VALUES (49, 'future')",
"INSERT INTO __yoi_schema_migrations (version, name) VALUES (50, 'future')",
[],
)
.unwrap();
let error = apply_migrations(&conn).unwrap_err().to_string();
assert!(error.contains("schema version 49 is newer"), "{error}");
assert!(error.contains("schema version 50 is newer"), "{error}");
assert!(error.contains("refusing to serve"), "{error}");
}
@@ -12362,7 +12494,7 @@ VALUES ('workspace-b', 'ticket-b', 'related', 'ticket-a', NULL, 'tester', '2026-
assign_explicit_test_workspace_owner(&conn);
apply_migrations(&mut conn).unwrap();
assert_eq!(current_schema_version(&conn).unwrap(), 48);
assert_eq!(current_schema_version(&conn).unwrap(), 49);
let workspace_id: Option<String> = conn
.query_row(
"SELECT workspace_id FROM trusted_runtime_records WHERE runtime_id = 'runtime-a'",
@@ -12991,7 +13123,7 @@ WHERE workspace_id = 'workspace-a'
assign_explicit_test_workspace_owner(&conn);
let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 48);
assert_eq!(store.schema_version().await.unwrap(), 49);
store
.with_conn(|conn| {
@@ -13180,7 +13312,7 @@ CREATE TABLE ticket_assignment_operations (
#[tokio::test]
async fn repository_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 48);
assert_eq!(store.schema_version().await.unwrap(), 49);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: "owner-account".to_string(),
@@ -13258,7 +13390,7 @@ CREATE TABLE ticket_assignment_operations (
#[tokio::test]
async fn memory_authority_records_round_trip_and_close_staging() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 48);
assert_eq!(store.schema_version().await.unwrap(), 49);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: "owner-account".to_string(),
@@ -13671,7 +13803,7 @@ CREATE TABLE ticket_assignment_operations (
#[tokio::test]
async fn account_and_login_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 48);
assert_eq!(store.schema_version().await.unwrap(), 49);
let now = "2026-07-22T00:00:00Z".to_string();
let account = AccountRecord {
account_id: "acct-user-alice".to_string(),
@@ -1,4 +1,4 @@
use rusqlite::{OptionalExtension, params};
use rusqlite::{OptionalExtension, TransactionBehavior, params};
use sha2::{Digest, Sha256};
use crate::store::WorkdirCreateOperationRecord;
@@ -103,6 +103,65 @@ impl SqliteWorkspaceStore {
})
}
pub fn begin_failed_workdir_create_retry(
&self,
workspace_id: &str,
operation_id: &str,
request_fingerprint: &str,
updated_at: &str,
) -> Result<WorkdirCreateOperationRecord> {
self.with_conn_mut(|conn| {
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let operation = read_workdir_create_operation(&tx, workspace_id, operation_id)?
.ok_or_else(|| {
Error::RegistryInconsistency(format!(
"Workdir create operation `{operation_id}` disappeared before retry"
))
})?;
if operation.request_fingerprint != request_fingerprint {
return Err(Error::InvalidInput(format!(
"Workdir create operation `{operation_id}` was reused with different input"
)));
}
if operation.state != "failed" {
return Err(Error::WorkdirAttachmentConflict(format!(
"Workdir create operation `{operation_id}` is not a failed retry"
)));
}
let removal_pending: bool = tx.query_row(
"SELECT EXISTS(SELECT 1 FROM workdir_removal_operations WHERE workspace_id=?1 AND workdir_id=?2 AND state='pending')",
params![workspace_id, operation.working_directory_id],
|row| row.get(0),
)?;
if removal_pending {
return Err(Error::WorkdirAttachmentConflict(format!(
"Workdir {} has a pending durable removal operation",
operation.working_directory_id
)));
}
let changed = tx.execute(
r#"UPDATE workdir_create_operations
SET state='pending', failure=NULL, updated_at=?1
WHERE workspace_id=?2 AND operation_id=?3
AND request_fingerprint=?4 AND state='failed'"#,
params![updated_at, workspace_id, operation_id, request_fingerprint],
)?;
if changed != 1 {
return Err(Error::WorkdirAttachmentConflict(format!(
"Workdir create operation `{operation_id}` retry was claimed concurrently"
)));
}
let updated = read_workdir_create_operation(&tx, workspace_id, operation_id)?
.ok_or_else(|| {
Error::RegistryInconsistency(format!(
"Workdir create operation `{operation_id}` disappeared after retry claim"
))
})?;
tx.commit()?;
Ok(updated)
})
}
pub fn bind_workdir_create_repository_access(
&self,
workspace_id: &str,
@@ -406,11 +465,32 @@ mod tests {
.unwrap();
assert_eq!(replayed, bound);
assert_eq!(replayed.source_uri.as_deref(), Some("/tmp/repo"));
let failed = store
.finish_workdir_create_operation(
"workspace",
"call-1",
&record.request_fingerprint,
false,
Some("provider failed"),
"2026-08-24T00:00:03Z",
)
.unwrap();
assert_eq!(failed.state, "failed");
let retry = store
.begin_failed_workdir_create_retry(
"workspace",
"call-1",
&record.request_fingerprint,
"2026-08-24T00:00:04Z",
)
.unwrap();
assert_eq!(retry.state, "pending");
assert_eq!(retry.failure, None);
assert_eq!(
store
.load_workdir_create_operation("workspace", "call-1")
.unwrap(),
Some(bound.clone())
Some(retry.clone())
);
let mut changed_input = record.clone();
changed_input.request_fingerprint =
File diff suppressed because it is too large Load Diff