merge: worker workdir registry

# Conflicts:
#	.yoi/tickets/00001KX6BPY7M/item.md
#	.yoi/tickets/00001KX6BPY7M/thread.md
This commit is contained in:
Keisuke Hirata 2026-07-11 02:35:48 +09:00
commit 03652f8210
No known key found for this signature in database
6 changed files with 1292 additions and 56 deletions

View File

@ -117,6 +117,10 @@ pub struct WorkingDirectoryRequest {
pub materializer: MaterializerKind, pub materializer: MaterializerKind,
#[serde(default)] #[serde(default)]
pub dirty_state_policy: DirtyStatePolicy, pub dirty_state_policy: DirtyStatePolicy,
/// Backend-assigned stable Workdir id. Runtimes use this when present so the
/// Backend can create canonical registry rows before materialization.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend_workdir_id: Option<String>,
} }
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
@ -158,6 +162,10 @@ pub struct WorkingDirectorySummary {
#[serde(default, skip_serializing_if = "Option::is_none")] #[serde(default, skip_serializing_if = "Option::is_none")]
pub cleanup_policy: Option<String>, pub cleanup_policy: Option<String>,
pub status: WorkingDirectoryStatusKind, pub status: WorkingDirectoryStatusKind,
/// Backend projection metadata. Runtimes leave this absent; Workspace Browser
/// APIs fill it with `backend_managed` or `runtime_unmanaged`.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub management_kind: Option<String>,
} }
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]

View File

@ -1036,6 +1036,7 @@ mod tests {
}, },
materializer: MaterializerKind::LocalGitWorktree, materializer: MaterializerKind::LocalGitWorktree,
dirty_state_policy: DirtyStatePolicy::CleanPointOnly, dirty_state_policy: DirtyStatePolicy::CleanPointOnly,
backend_workdir_id: None,
} }
} }

View File

@ -51,6 +51,7 @@ impl WorkingDirectory {
cleanup_target: Some(self.cleanup_target.clone()), cleanup_target: Some(self.cleanup_target.clone()),
cleanup_policy: Some(self.cleanup_policy.clone()), cleanup_policy: Some(self.cleanup_policy.clone()),
status: self.status.clone(), status: self.status.clone(),
management_kind: None,
} }
} }
} }
@ -395,7 +396,10 @@ impl WorkingDirectoryMaterializer for LocalGitWorktreeMaterializer {
&self, &self,
request: &WorkingDirectoryRequest, request: &WorkingDirectoryRequest,
) -> Result<WorkingDirectoryBinding, WorkingDirectoryDiagnostic> { ) -> Result<WorkingDirectoryBinding, WorkingDirectoryDiagnostic> {
let working_directory_id = next_working_directory_id(&request.repository.id); let working_directory_id = request
.backend_workdir_id
.clone()
.unwrap_or_else(|| next_working_directory_id(&request.repository.id));
self.materialize_with_working_directory_id( self.materialize_with_working_directory_id(
working_directory_id, working_directory_id,
request, request,
@ -719,6 +723,7 @@ mod tests {
}, },
materializer: MaterializerKind::LocalGitWorktree, materializer: MaterializerKind::LocalGitWorktree,
dirty_state_policy: DirtyStatePolicy::CleanPointOnly, dirty_state_policy: DirtyStatePolicy::CleanPointOnly,
backend_workdir_id: None,
} }
} }

View File

@ -24,10 +24,11 @@ use crate::config::{RemoteRuntimeConfigFile, WorkspaceBackendConfigFile, resolve
use crate::hosts::{ use crate::hosts::{
ConfigBundleCheckResult, ConfigBundleSyncResult, DiagnosticSeverity, EmbeddedWorkerRuntime, ConfigBundleCheckResult, ConfigBundleSyncResult, DiagnosticSeverity, EmbeddedWorkerRuntime,
HostSummary, RemoteRuntimeConfig, RemoteWorkerRuntime, RuntimeDiagnostic, RuntimeRegistry, HostSummary, RemoteRuntimeConfig, RemoteWorkerRuntime, RuntimeDiagnostic, RuntimeRegistry,
RuntimeRegistryUnregisterResult, RuntimeSummary, WorkerInputRequest, WorkerInputResult, RuntimeRegistryUnregisterResult, RuntimeSummary, WorkerCapabilitySummary,
WorkerLifecycleRequest, WorkerLifecycleResult, WorkerOperationState, WorkerImplementationSummary, WorkerInputRequest, WorkerInputResult, WorkerLifecycleRequest,
WorkerSpawnAcceptanceRequirement, WorkerSpawnIntent, WorkerSpawnRequest, WorkerSpawnResult, WorkerLifecycleResult, WorkerOperationState, WorkerSpawnAcceptanceRequirement,
WorkerSpawnWorkingDirectoryRequest, WorkerSummary, WorkerTranscriptProjection, WorkerSpawnIntent, WorkerSpawnRequest, WorkerSpawnResult, WorkerSpawnWorkingDirectoryRequest,
WorkerSummary, WorkerTranscriptProjection, WorkerWorkspaceSummary,
}; };
use crate::identity::WorkspaceIdentity; use crate::identity::WorkspaceIdentity;
use crate::observation::{ use crate::observation::{
@ -48,12 +49,16 @@ use crate::repositories::{
RepositoryRegistryReader, RepositorySummary, RepositoryRegistryReader, RepositorySummary,
}; };
use crate::resource_broker::BackendResourceBroker; use crate::resource_broker::BackendResourceBroker;
use crate::store::{ControlPlaneStore, WorkspaceRecord}; use crate::store::{
ControlPlaneStore, WorkdirRegistryRecord, WorkerRegistryRecord, WorkerWorkdirLinkRecord,
WorkspaceRecord,
};
use crate::{Error, Result}; use crate::{Error, Result};
use worker_runtime::catalog::{ use worker_runtime::catalog::{
ConfigBundleRef, DirtyStatePolicy, MaterializerKind, ProfileSelector, ConfigBundleRef, DirtyStatePolicy, MaterializerKind, ProfileSelector,
RepositorySelector as RuntimeRepositorySelector, WorkingDirectoryClaim, RepositorySelector as RuntimeRepositorySelector, WorkingDirectoryClaim,
WorkingDirectoryRepository, WorkingDirectoryRequest, WorkingDirectorySummary, WorkingDirectoryRepository, WorkingDirectoryRequest, WorkingDirectoryStatusKind,
WorkingDirectorySummary,
}; };
use worker_runtime::config_bundle::ConfigBundle; use worker_runtime::config_bundle::ConfigBundle;
use worker_runtime::http_server::{ use worker_runtime::http_server::{
@ -1182,18 +1187,11 @@ async fn scoped_list_runtime_working_directories(
AxumPath(path): AxumPath<ScopedRuntimePath>, AxumPath(path): AxumPath<ScopedRuntimePath>,
) -> ApiResult<Json<BrowserWorkingDirectoryListResponse>> { ) -> ApiResult<Json<BrowserWorkingDirectoryListResponse>> {
validate_workspace_scope(&api, &path.workspace_id)?; validate_workspace_scope(&api, &path.workspace_id)?;
let list = api let (items, diagnostics) = runtime_working_directory_summaries(&api, &path.runtime_id)?;
.runtime
.list_working_directories(&path.runtime_id)
.map_err(|err| err.into_error())?;
Ok(Json(BrowserWorkingDirectoryListResponse { Ok(Json(BrowserWorkingDirectoryListResponse {
workspace_id: api.config.workspace_id.clone(), workspace_id: api.config.workspace_id.clone(),
items: list items,
.items diagnostics,
.into_iter()
.map(|status| status.summary)
.collect(),
diagnostics: list.diagnostics,
})) }))
} }
@ -1274,12 +1272,36 @@ fn create_working_directory_for_runtime(
request: BrowserWorkingDirectoryCreateRequest, request: BrowserWorkingDirectoryCreateRequest,
) -> ApiResult<Json<BrowserWorkingDirectoryDetailResponse>> { ) -> ApiResult<Json<BrowserWorkingDirectoryDetailResponse>> {
let runtime_id = request.runtime_id.clone(); let runtime_id = request.runtime_id.clone();
let working_directory_request = working_directory_request_for_browser(&api, request)?; let mut working_directory_request = working_directory_request_for_browser(&api, request)?;
let workdir_id = next_backend_workdir_id(&working_directory_request.repository.id);
working_directory_request.backend_workdir_id = Some(workdir_id.clone());
let pending = WorkdirRegistryRecord {
workspace_id: api.config.workspace_id.clone(),
workdir_id: workdir_id.clone(),
runtime_id: runtime_id.clone(),
repository_id: working_directory_request.repository.id.clone(),
selector: working_directory_request
.repository
.selector
.as_ref()
.map(|selector| selector.as_ref().to_string()),
resolved_commit: None,
materialization_status: "pending".to_string(),
cleanliness: "clean".to_string(),
management_kind: "backend_managed".to_string(),
created_at: now_registry_timestamp(),
updated_at: now_registry_timestamp(),
};
api.store.upsert_workdir_registry(&pending)?;
let result = api let result = api
.runtime .runtime
.create_working_directory(&runtime_id, working_directory_request) .create_working_directory(&runtime_id, working_directory_request)
.map_err(|err| err.into_error())?; .map_err(|err| err.into_error())?;
let Some(working_directory) = result.working_directory else { let Some(working_directory) = result.working_directory else {
let mut failed = pending;
failed.materialization_status = "failed".to_string();
failed.updated_at = now_registry_timestamp();
api.store.upsert_workdir_registry(&failed)?;
return Err(ApiError::with_diagnostics( return Err(ApiError::with_diagnostics(
Error::RuntimeOperationFailed { Error::RuntimeOperationFailed {
runtime_id, runtime_id,
@ -1289,6 +1311,13 @@ fn create_working_directory_for_runtime(
result.diagnostics, result.diagnostics,
)); ));
}; };
let record = workdir_record_from_summary(
&api,
&runtime_id,
&working_directory.summary,
"backend_managed",
);
api.store.upsert_workdir_registry(&record)?;
Ok(Json(BrowserWorkingDirectoryDetailResponse { Ok(Json(BrowserWorkingDirectoryDetailResponse {
workspace_id: api.config.workspace_id.clone(), workspace_id: api.config.workspace_id.clone(),
item: working_directory.summary, item: working_directory.summary,
@ -1305,21 +1334,43 @@ fn working_directory_detail_for_runtime(
.runtime .runtime
.working_directory(runtime_id, working_directory_id) .working_directory(runtime_id, working_directory_id)
.map_err(|err| err.into_error())?; .map_err(|err| err.into_error())?;
let Some(working_directory) = result.working_directory else { if let Some(working_directory) = result.working_directory {
return Err(ApiError::with_diagnostics( let management_kind = api
Error::RuntimeOperationFailed { .store
runtime_id: runtime_id.to_string(), .get_workdir_registry(&api.config.workspace_id, working_directory_id)?
code: "workspace_working_directory_lookup_failed".to_string(), .map(|record| record.management_kind)
message: "Runtime did not return working directory".to_string(), .unwrap_or_else(|| "runtime_unmanaged".to_string());
}, let record = workdir_record_from_summary(
result.diagnostics, &api,
)); runtime_id,
}; &working_directory.summary,
Ok(Json(BrowserWorkingDirectoryDetailResponse { management_kind.as_str(),
workspace_id: api.config.workspace_id.clone(), );
item: working_directory.summary, api.store.upsert_workdir_registry(&record)?;
diagnostics: result.diagnostics, return Ok(Json(BrowserWorkingDirectoryDetailResponse {
})) workspace_id: api.config.workspace_id.clone(),
item: working_directory.summary,
diagnostics: result.diagnostics,
}));
}
if let Some(record) = api
.store
.get_workdir_registry(&api.config.workspace_id, working_directory_id)?
{
return Ok(Json(BrowserWorkingDirectoryDetailResponse {
workspace_id: api.config.workspace_id.clone(),
item: workdir_summary_from_record(&record),
diagnostics: result.diagnostics,
}));
}
Err(ApiError::with_diagnostics(
Error::RuntimeOperationFailed {
runtime_id: runtime_id.to_string(),
code: "workspace_working_directory_lookup_failed".to_string(),
message: "Runtime did not return working directory".to_string(),
},
result.diagnostics,
))
} }
fn cleanup_working_directory_for_runtime( fn cleanup_working_directory_for_runtime(
@ -1341,6 +1392,18 @@ fn cleanup_working_directory_for_runtime(
result.diagnostics, result.diagnostics,
)); ));
}; };
let management_kind = api
.store
.get_workdir_registry(&api.config.workspace_id, working_directory_id)?
.map(|record| record.management_kind)
.unwrap_or_else(|| "runtime_unmanaged".to_string());
let record = workdir_record_from_summary(
&api,
runtime_id,
&working_directory.summary,
management_kind.as_str(),
);
api.store.upsert_workdir_registry(&record)?;
Ok(Json(BrowserWorkingDirectoryDetailResponse { Ok(Json(BrowserWorkingDirectoryDetailResponse {
workspace_id: api.config.workspace_id.clone(), workspace_id: api.config.workspace_id.clone(),
item: working_directory.summary, item: working_directory.summary,
@ -1966,6 +2029,7 @@ fn working_directory_request_from_repository(
}, },
materializer: MaterializerKind::LocalGitWorktree, materializer: MaterializerKind::LocalGitWorktree,
dirty_state_policy: DirtyStatePolicy::CleanPointOnly, dirty_state_policy: DirtyStatePolicy::CleanPointOnly,
backend_workdir_id: None,
} }
} }
@ -2026,6 +2090,10 @@ async fn create_workspace_worker(
content: initial_text, content: initial_text,
}) })
}; };
let selected_working_directory_id = request
.working_directory
.as_ref()
.map(|selection| selection.working_directory_id.clone());
let resolved_working_directory = let resolved_working_directory =
request request
.working_directory .working_directory
@ -2039,7 +2107,7 @@ async fn create_workspace_worker(
.spawn_worker( .spawn_worker(
&request.runtime_id, &request.runtime_id,
WorkerSpawnRequest { WorkerSpawnRequest {
requested_worker_name: Some(display_name), requested_worker_name: Some(display_name.clone()),
intent: WorkerSpawnIntent::WorkspaceCoding, intent: WorkerSpawnIntent::WorkspaceCoding,
acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted {
expected_segments: if initial_input.is_some() { 1 } else { 0 }, expected_segments: if initial_input.is_some() { 1 } else { 0 },
@ -2064,6 +2132,37 @@ async fn create_workspace_worker(
code: "workspace_worker_create_missing_summary".to_string(), code: "workspace_worker_create_missing_summary".to_string(),
message: "Runtime completed worker creation without returning a Worker summary".to_string(), message: "Runtime completed worker creation without returning a Worker summary".to_string(),
})?; })?;
let worker_record = sync_worker_observation(&api, &worker)?;
if let Some(workdir_id) = selected_working_directory_id.as_deref() {
if api
.store
.get_workdir_registry(&api.config.workspace_id, workdir_id)?
.is_none()
{
if let Ok(result) = api
.runtime
.working_directory(worker.runtime_id.as_str(), workdir_id)
.map_err(|err| err.into_error())
{
if let Some(status) = result.working_directory {
let record = workdir_record_from_summary(
&api,
worker.runtime_id.as_str(),
&status.summary,
"runtime_unmanaged",
);
api.store.upsert_workdir_registry(&record)?;
}
}
}
if api
.store
.get_workdir_registry(&api.config.workspace_id, workdir_id)?
.is_some()
{
link_worker_to_workdir(&api, &worker_record, workdir_id)?;
}
}
let runtime_id = worker.runtime_id.clone(); let runtime_id = worker.runtime_id.clone();
let worker_id = worker.worker_id.clone(); let worker_id = worker.worker_id.clone();
let workspace_id = api.workspace_id().to_string(); let workspace_id = api.workspace_id().to_string();
@ -2172,10 +2271,47 @@ async fn create_runtime_worker(
configured_working_directory_request(&api.config, working_directory) configured_working_directory_request(&api.config, working_directory)
}) })
.transpose()?; .transpose()?;
let prepared_workdir_id = if let Some(working_directory_request) =
request.resolved_working_directory_request.as_mut()
{
Some(upsert_pending_backend_workdir(
&api,
&runtime_id,
working_directory_request,
)?)
} else {
request
.resolved_working_directory
.as_ref()
.map(|claim| claim.working_directory_id.clone())
};
let result = api let result = api
.runtime .runtime
.spawn_worker(&runtime_id, request) .spawn_worker(&runtime_id, request)
.map_err(|err| err.into_error())?; .map_err(|err| err.into_error())?;
if let Some(worker) = result.worker.as_ref() {
let record = sync_worker_observation(&api, worker)?;
if worker.working_directory.is_none() {
if let Some(workdir_id) = prepared_workdir_id.as_deref() {
if api
.store
.get_workdir_registry(&api.config.workspace_id, workdir_id)?
.is_some()
{
link_worker_to_workdir(&api, &record, workdir_id)?;
}
}
}
} else if let Some(workdir_id) = prepared_workdir_id.as_deref() {
if let Some(mut record) = api
.store
.get_workdir_registry(&api.config.workspace_id, workdir_id)?
{
record.materialization_status = "failed".to_string();
record.updated_at = now_registry_timestamp();
api.store.upsert_workdir_registry(&record)?;
}
}
Ok(Json(result)) Ok(Json(result))
} }
@ -2230,6 +2366,16 @@ async fn stop_runtime_worker(
.runtime .runtime
.stop_worker(&runtime_id, &worker_id, request) .stop_worker(&runtime_id, &worker_id, request)
.map_err(|err| err.into_error())?; .map_err(|err| err.into_error())?;
let backend_id = backend_worker_id(&runtime_id, &worker_id);
if let Some(mut record) = api
.store
.get_worker_registry(&api.config.workspace_id, backend_id.as_str())?
{
record.lifecycle_state = "stopped".to_string();
record.updated_at = now_registry_timestamp();
api.store.upsert_worker_registry(&record)?;
sync_linked_workdir_after_worker_stop(&api, &runtime_id, &record)?;
}
Ok(Json(result)) Ok(Json(result))
} }
@ -2423,11 +2569,37 @@ async fn list_host_workers(
fn workers_response(api: WorkspaceApi) -> ApiResult<RuntimeListResponse<WorkerSummary>> { fn workers_response(api: WorkspaceApi) -> ApiResult<RuntimeListResponse<WorkerSummary>> {
let limit = api.config.max_records.min(200); let limit = api.config.max_records.min(200);
let runtime_workers = api.runtime.list_workers(limit); let runtime_workers = api.runtime.list_workers(limit);
let mut observed = std::collections::BTreeMap::new();
for worker in &runtime_workers.items {
let _ = sync_worker_observation(&api, worker);
observed.insert(
backend_worker_id(worker.runtime_id.as_str(), worker.worker_id.as_str()),
worker.clone(),
);
}
let worker_records = api
.store
.list_worker_registry(&api.config.workspace_id, limit)?;
let workdir_records = api
.store
.list_workdir_registry(&api.config.workspace_id, 500)?;
let mut items = Vec::new();
for record in worker_records {
let links = api
.store
.list_worker_workdir_links(&api.config.workspace_id, record.worker_id.as_str())?;
items.push(merge_worker_registry_projection(
observed.get(record.worker_id.as_str()),
&record,
links,
&workdir_records,
));
}
Ok(RuntimeListResponse { Ok(RuntimeListResponse {
workspace_id: api.config.workspace_id, workspace_id: api.config.workspace_id,
limit, limit,
items: runtime_workers.items, items,
source: "worker_runtime_registry".to_string(), source: "backend_worker_registry".to_string(),
diagnostics: runtime_workers.diagnostics, diagnostics: runtime_workers.diagnostics,
}) })
} }
@ -3126,25 +3298,377 @@ fn working_directory_repository_options(
} }
fn working_directory_summaries(api: &WorkspaceApi) -> ApiResult<Vec<WorkingDirectorySummary>> { fn working_directory_summaries(api: &WorkspaceApi) -> ApiResult<Vec<WorkingDirectorySummary>> {
let list = api let _ = sync_all_runtime_workdir_observations(api);
.runtime let records = api
.list_working_directories(EMBEDDED_WORKER_RUNTIME_ID) .store
.map_err(|err| err.into_error())?; .list_managed_workdir_registry(&api.config.workspace_id, 200)?;
if !list.diagnostics.is_empty() { Ok(records
return Err(ApiError::with_diagnostics( .iter()
Error::RuntimeOperationFailed { .map(workdir_summary_from_record)
runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), .collect::<Vec<_>>())
code: "workspace_working_directory_list_failed".to_string(), }
message: "Runtime did not list working directories".to_string(),
}, fn runtime_working_directory_summaries(
list.diagnostics, api: &WorkspaceApi,
)); runtime_id: &str,
) -> ApiResult<(Vec<WorkingDirectorySummary>, Vec<RuntimeDiagnostic>)> {
let diagnostics = sync_runtime_workdir_observations(api, runtime_id)?;
let records = api
.store
.list_workdir_registry(&api.config.workspace_id, 200)?;
let items = records
.iter()
.filter(|record| record.runtime_id == runtime_id)
.map(workdir_summary_from_record)
.collect::<Vec<_>>();
Ok((items, diagnostics))
}
fn backend_worker_id(runtime_id: &str, runtime_worker_id: &str) -> String {
format!("{runtime_id}/{runtime_worker_id}")
}
fn now_registry_timestamp() -> String {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis().to_string())
.unwrap_or_else(|_| "0".to_string())
}
fn registry_safe_id_component(value: &str) -> String {
let sanitized: String = value
.chars()
.map(|ch| {
if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
ch
} else {
'-'
}
})
.collect();
sanitized.trim_matches('-').chars().take(48).collect()
}
fn next_backend_workdir_id(repository_id: &str) -> String {
let repository = registry_safe_id_component(repository_id);
format!(
"backend-{}-{}",
now_registry_timestamp(),
if repository.is_empty() {
"workdir"
} else {
&repository
}
)
}
fn record_worker_summary(
api: &WorkspaceApi,
worker: &WorkerSummary,
display_name: &str,
profile: Option<String>,
) -> ApiResult<WorkerRegistryRecord> {
let timestamp = now_registry_timestamp();
let worker_id = backend_worker_id(worker.runtime_id.as_str(), worker.worker_id.as_str());
let record = WorkerRegistryRecord {
workspace_id: api.config.workspace_id.clone(),
worker_id: worker_id.clone(),
runtime_id: worker.runtime_id.as_str().to_string(),
runtime_worker_id: worker.worker_id.as_str().to_string(),
display_name: display_name.to_string(),
profile,
lifecycle_state: worker.status.clone(),
retention_state: "normal".to_string(),
transcript_ref: Some(format!(
"runtime://{}/workers/{}/transcript",
worker.runtime_id.as_str(),
worker.worker_id.as_str()
)),
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: timestamp.clone(),
updated_at: timestamp,
};
api.store.upsert_worker_registry(&record)?;
Ok(api
.store
.get_worker_registry(&api.config.workspace_id, worker_id.as_str())?
.unwrap_or(record))
}
fn worker_summary_from_registry(record: &WorkerRegistryRecord) -> WorkerSummary {
WorkerSummary {
worker_id: record.runtime_worker_id.clone(),
runtime_id: record.runtime_id.clone(),
host_id: "backend-registry".to_string(),
role: None,
label: record.display_name.clone(),
status: record.lifecycle_state.clone(),
state: record.lifecycle_state.clone(),
last_seen_at: Some(record.updated_at.clone()),
capabilities: WorkerCapabilitySummary {
can_accept_input: false,
can_stop: false,
can_spawn_followup: false,
},
workspace: WorkerWorkspaceSummary {
visibility: "backend_registry".to_string(),
identity: record.workspace_id.clone(),
},
profile: record.profile.clone(),
implementation: WorkerImplementationSummary {
kind: "backend_worker_registry".to_string(),
display_hint: "Archived Worker".to_string(),
},
working_directory: None,
diagnostics: vec![RuntimeDiagnostic {
code: "backend_worker_registry_only".to_string(),
severity: DiagnosticSeverity::Info,
message:
"Worker is preserved in the Backend registry without a live Runtime observation"
.to_string(),
}],
} }
Ok(list }
.items
fn merge_worker_registry_projection(
live: Option<&WorkerSummary>,
record: &WorkerRegistryRecord,
links: Vec<WorkerWorkdirLinkRecord>,
workdirs: &[WorkdirRegistryRecord],
) -> WorkerSummary {
let mut summary = live
.cloned()
.unwrap_or_else(|| worker_summary_from_registry(record));
summary.label = record.display_name.clone();
summary.status = record.lifecycle_state.clone();
summary.state = record.lifecycle_state.clone();
summary.profile = record.profile.clone();
summary.working_directory = links.iter().find_map(|link| {
workdirs
.iter()
.find(|workdir| workdir.workdir_id == link.workdir_id)
.map(|workdir| workdir_summary_from_record(workdir))
});
summary
}
fn sync_worker_observation(
api: &WorkspaceApi,
worker: &WorkerSummary,
) -> ApiResult<WorkerRegistryRecord> {
let record = record_worker_summary(api, worker, worker.label.as_str(), worker.profile.clone())?;
if let Some(working_directory) = worker.working_directory.as_ref() {
let management_kind = api
.store
.get_workdir_registry(
&api.config.workspace_id,
&working_directory.working_directory_id,
)?
.map(|existing| existing.management_kind)
.unwrap_or_else(|| "runtime_unmanaged".to_string());
let workdir_record = workdir_record_from_summary(
api,
worker.runtime_id.as_str(),
working_directory,
management_kind.as_str(),
);
api.store.upsert_workdir_registry(&workdir_record)?;
link_worker_to_workdir(api, &record, &working_directory.working_directory_id)?;
}
Ok(record)
}
fn upsert_pending_backend_workdir(
api: &WorkspaceApi,
runtime_id: &str,
request: &mut WorkingDirectoryRequest,
) -> ApiResult<String> {
let workdir_id = request
.backend_workdir_id
.clone()
.unwrap_or_else(|| next_backend_workdir_id(&request.repository.id));
request.backend_workdir_id = Some(workdir_id.clone());
let timestamp = now_registry_timestamp();
api.store.upsert_workdir_registry(&WorkdirRegistryRecord {
workspace_id: api.config.workspace_id.clone(),
workdir_id: workdir_id.clone(),
runtime_id: runtime_id.to_string(),
repository_id: request.repository.id.clone(),
selector: request
.repository
.selector
.as_ref()
.map(|selector| selector.as_ref().to_string()),
resolved_commit: None,
materialization_status: "pending".to_string(),
cleanliness: "clean".to_string(),
management_kind: "backend_managed".to_string(),
created_at: timestamp.clone(),
updated_at: timestamp,
})?;
Ok(workdir_id)
}
fn sync_runtime_workdir_observations(
api: &WorkspaceApi,
runtime_id: &str,
) -> ApiResult<Vec<RuntimeDiagnostic>> {
let response = api
.runtime
.list_working_directories(runtime_id)
.map_err(|err| err.into_error())?;
let mut observed = std::collections::BTreeSet::new();
for status in &response.items {
observed.insert(status.summary.working_directory_id.clone());
let management_kind = api
.store
.get_workdir_registry(
&api.config.workspace_id,
&status.summary.working_directory_id,
)?
.map(|existing| existing.management_kind)
.unwrap_or_else(|| "runtime_unmanaged".to_string());
let record =
workdir_record_from_summary(api, runtime_id, &status.summary, management_kind.as_str());
api.store.upsert_workdir_registry(&record)?;
}
for mut record in api
.store
.list_workdir_registry(&api.config.workspace_id, 500)?
.into_iter() .into_iter()
.map(|status| status.summary) .filter(|record| record.runtime_id == runtime_id && !observed.contains(&record.workdir_id))
.collect()) {
if record.materialization_status == "present" {
record.materialization_status = "missing".to_string();
record.updated_at = now_registry_timestamp();
api.store.upsert_workdir_registry(&record)?;
}
}
Ok(response.diagnostics)
}
fn sync_all_runtime_workdir_observations(api: &WorkspaceApi) -> Vec<RuntimeDiagnostic> {
let mut diagnostics = Vec::new();
let runtimes = api.runtime.list_runtimes(api.config.max_records.min(200));
for runtime in runtimes.items {
if runtime.capabilities.supports_worktrees {
match sync_runtime_workdir_observations(api, runtime.runtime_id.as_str()) {
Ok(mut runtime_diagnostics) => diagnostics.append(&mut runtime_diagnostics),
Err(err) => diagnostics.extend(err.diagnostics),
}
}
}
diagnostics
}
fn sync_linked_workdir_after_worker_stop(
api: &WorkspaceApi,
runtime_id: &str,
worker_record: &WorkerRegistryRecord,
) -> ApiResult<()> {
let links = api
.store
.list_worker_workdir_links(&api.config.workspace_id, worker_record.worker_id.as_str())?;
for link in links {
let result = api
.runtime
.working_directory(runtime_id, link.workdir_id.as_str())
.map_err(|err| err.into_error())?;
if let Some(status) = result.working_directory {
let management_kind = api
.store
.get_workdir_registry(&api.config.workspace_id, link.workdir_id.as_str())?
.map(|record| record.management_kind)
.unwrap_or_else(|| "runtime_unmanaged".to_string());
let record = workdir_record_from_summary(
api,
runtime_id,
&status.summary,
management_kind.as_str(),
);
api.store.upsert_workdir_registry(&record)?;
} else if let Some(mut record) = api
.store
.get_workdir_registry(&api.config.workspace_id, link.workdir_id.as_str())?
{
record.materialization_status = "missing".to_string();
record.updated_at = now_registry_timestamp();
api.store.upsert_workdir_registry(&record)?;
}
}
Ok(())
}
fn workdir_record_from_summary(
api: &WorkspaceApi,
runtime_id: &str,
summary: &WorkingDirectorySummary,
management_kind: &str,
) -> WorkdirRegistryRecord {
let timestamp = now_registry_timestamp();
WorkdirRegistryRecord {
workspace_id: api.config.workspace_id.clone(),
workdir_id: summary.working_directory_id.clone(),
runtime_id: runtime_id.to_string(),
repository_id: summary.repository_id.clone(),
selector: summary.requested_selector.clone(),
resolved_commit: summary.resolved_commit.clone(),
materialization_status: match summary.status {
WorkingDirectoryStatusKind::Active => "present",
WorkingDirectoryStatusKind::Removed => "removed",
WorkingDirectoryStatusKind::CleanupPending => "pending",
}
.to_string(),
cleanliness: "clean".to_string(),
management_kind: management_kind.to_string(),
created_at: timestamp.clone(),
updated_at: timestamp,
}
}
fn workdir_summary_from_record(record: &WorkdirRegistryRecord) -> WorkingDirectorySummary {
let status = match record.materialization_status.as_str() {
"present" => WorkingDirectoryStatusKind::Active,
"pending" => WorkingDirectoryStatusKind::CleanupPending,
_ => WorkingDirectoryStatusKind::Removed,
};
WorkingDirectorySummary {
working_directory_id: record.workdir_id.clone(),
repository_id: record.repository_id.clone(),
requested_selector: record.selector.clone(),
materializer_kind: MaterializerKind::LocalGitWorktree,
dirty_state_policy: DirtyStatePolicy::CleanPointOnly,
resolved_commit: record.resolved_commit.clone(),
resolved_tree: None,
cleanup_target: Some(worker_runtime::catalog::WorkingDirectoryCleanupTarget {
kind: "local_git_worktree".to_string(),
working_directory_id: record.workdir_id.clone(),
repository_id: record.repository_id.clone(),
}),
cleanup_policy: Some("manual_or_worker_stop".to_string()),
status,
management_kind: Some(record.management_kind.clone()),
}
}
fn link_worker_to_workdir(
api: &WorkspaceApi,
worker_record: &WorkerRegistryRecord,
workdir_id: &str,
) -> ApiResult<()> {
let timestamp = now_registry_timestamp();
api.store
.upsert_worker_workdir_link(&WorkerWorkdirLinkRecord {
workspace_id: api.config.workspace_id.clone(),
worker_id: worker_record.worker_id.clone(),
workdir_id: workdir_id.to_string(),
role: "primary_cwd".to_string(),
linked_at: timestamp,
unlinked_at: None,
})?;
Ok(())
} }
fn validate_working_directory_claim_for_browser( fn validate_working_directory_claim_for_browser(
@ -3207,6 +3731,7 @@ fn working_directory_request_for_browser(
}, },
materializer: MaterializerKind::LocalGitWorktree, materializer: MaterializerKind::LocalGitWorktree,
dirty_state_policy: DirtyStatePolicy::CleanPointOnly, dirty_state_policy: DirtyStatePolicy::CleanPointOnly,
backend_workdir_id: None,
}) })
} }
@ -3711,6 +4236,148 @@ mod tests {
const TEST_REPOSITORY_ID: &str = "main"; const TEST_REPOSITORY_ID: &str = "main";
const TEST_CREATED_AT: &str = "2026-06-23T06:43:28Z"; const TEST_CREATED_AT: &str = "2026-06-23T06:43:28Z";
#[test]
fn backend_worker_projection_preserves_archive_rows_links_and_redacts_paths() {
let worker = WorkerRegistryRecord {
workspace_id: "workspace-1".to_string(),
worker_id: "embedded/worker-1".to_string(),
runtime_id: "embedded".to_string(),
runtime_worker_id: "worker-1".to_string(),
display_name: "Archived Worker".to_string(),
profile: Some("builtin:coder".to_string()),
lifecycle_state: "stopped".to_string(),
retention_state: "pinned".to_string(),
transcript_ref: Some("runtime://embedded/workers/worker-1/transcript".to_string()),
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: "1".to_string(),
updated_at: "2".to_string(),
};
let workdir = WorkdirRegistryRecord {
workspace_id: "workspace-1".to_string(),
workdir_id: "backend-1-repo".to_string(),
runtime_id: "embedded".to_string(),
repository_id: "repo".to_string(),
selector: Some("develop".to_string()),
resolved_commit: Some("abcdef".to_string()),
materialization_status: "missing".to_string(),
cleanliness: "clean".to_string(),
management_kind: "backend_managed".to_string(),
created_at: "1".to_string(),
updated_at: "3".to_string(),
};
let link = WorkerWorkdirLinkRecord {
workspace_id: "workspace-1".to_string(),
worker_id: worker.worker_id.clone(),
workdir_id: workdir.workdir_id.clone(),
role: "primary_cwd".to_string(),
linked_at: "4".to_string(),
unlinked_at: None,
};
let projected = merge_worker_registry_projection(None, &worker, vec![link], &[workdir]);
assert_eq!(projected.status, "stopped");
assert_eq!(
projected.working_directory.as_ref().unwrap().status,
WorkingDirectoryStatusKind::Removed
);
assert_eq!(
projected
.working_directory
.as_ref()
.unwrap()
.management_kind
.as_deref(),
Some("backend_managed")
);
let serialized = serde_json::to_string(&projected).unwrap();
assert!(!serialized.contains("/tmp/"));
assert!(!serialized.contains("materialized_path"));
}
#[tokio::test]
async fn workspace_managed_workdir_summaries_exclude_runtime_unmanaged_rows() {
let dir = tempfile::tempdir().unwrap();
let api = test_api(dir.path()).await;
api.store
.upsert_workdir_registry(&WorkdirRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
workdir_id: "managed".to_string(),
runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(),
repository_id: "repo".to_string(),
selector: None,
resolved_commit: None,
materialization_status: "present".to_string(),
cleanliness: "clean".to_string(),
management_kind: "backend_managed".to_string(),
created_at: "1".to_string(),
updated_at: "1".to_string(),
})
.unwrap();
api.store
.upsert_workdir_registry(&WorkdirRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
workdir_id: "runtime-direct".to_string(),
runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(),
repository_id: "repo".to_string(),
selector: None,
resolved_commit: None,
materialization_status: "present".to_string(),
cleanliness: "unknown".to_string(),
management_kind: "runtime_unmanaged".to_string(),
created_at: "1".to_string(),
updated_at: "2".to_string(),
})
.unwrap();
let managed = working_directory_summaries(&api)
.unwrap_or_else(|err| panic!("working_directory_summaries failed: {}", err.error));
assert_eq!(managed.len(), 1);
assert_eq!(managed[0].working_directory_id, "managed");
assert_eq!(
managed[0].management_kind.as_deref(),
Some("backend_managed")
);
let (runtime_projection, _) =
runtime_working_directory_summaries(&api, EMBEDDED_WORKER_RUNTIME_ID).unwrap_or_else(
|err| panic!("runtime_working_directory_summaries failed: {}", err.error),
);
assert!(runtime_projection.iter().any(|summary| {
summary.working_directory_id == "runtime-direct"
&& summary.management_kind.as_deref() == Some("runtime_unmanaged")
}));
}
#[test]
fn unmanaged_runtime_workdir_projection_is_typed_and_diagnostic_safe() {
let workdir = WorkdirRegistryRecord {
workspace_id: "workspace-1".to_string(),
workdir_id: "runtime-direct".to_string(),
runtime_id: "embedded".to_string(),
repository_id: "repo".to_string(),
selector: None,
resolved_commit: None,
materialization_status: "present".to_string(),
cleanliness: "unknown".to_string(),
management_kind: "runtime_unmanaged".to_string(),
created_at: "1".to_string(),
updated_at: "2".to_string(),
};
let projected = workdir_summary_from_record(&workdir);
assert_eq!(projected.status, WorkingDirectoryStatusKind::Active);
assert_eq!(
projected.management_kind.as_deref(),
Some("runtime_unmanaged")
);
let serialized = serde_json::to_string(&projected).unwrap();
assert!(!serialized.contains("/tmp/"));
assert!(!serialized.contains("materialized_path"));
}
#[test] #[test]
fn worker_profile_candidates_are_backend_published_and_mapped() { fn worker_profile_candidates_are_backend_published_and_mapped() {
let candidates = worker_profile_candidates(); let candidates = worker_profile_candidates();

View File

@ -27,6 +27,11 @@ const MIGRATIONS: &[Migration] = &[
name: "align legacy workspace bootstrap with schema v0", name: "align legacy workspace bootstrap with schema v0",
apply: align_legacy_bootstrap_schema, apply: align_legacy_bootstrap_schema,
}, },
Migration {
version: 3,
name: "backend worker workdir registry schema",
apply: create_worker_workdir_registry_tables,
},
]; ];
struct Migration { struct Migration {
@ -44,11 +49,99 @@ pub struct WorkspaceRecord {
pub updated_at: String, pub updated_at: String,
} }
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkerRegistryRecord {
pub workspace_id: String,
/// Backend-owned archival Worker id. In v0 it is derived from runtime_id + runtime_worker_id.
pub worker_id: String,
pub runtime_id: String,
pub runtime_worker_id: String,
pub display_name: String,
pub profile: Option<String>,
pub lifecycle_state: String,
/// Retention state is explicit so `pinned` can be represented before prune exists.
pub retention_state: String,
pub transcript_ref: Option<String>,
pub session_ref: Option<String>,
pub summary_ref: Option<String>,
pub diagnostics_ref: Option<String>,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkdirRegistryRecord {
pub workspace_id: String,
pub workdir_id: String,
pub runtime_id: String,
pub repository_id: String,
pub selector: Option<String>,
pub resolved_commit: Option<String>,
pub materialization_status: String,
pub cleanliness: String,
/// `backend_managed` rows are authored by this Backend; `runtime_unmanaged` is for diagnostics only.
pub management_kind: String,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkerWorkdirLinkRecord {
pub workspace_id: String,
pub worker_id: String,
pub workdir_id: String,
pub role: String,
pub linked_at: String,
pub unlinked_at: Option<String>,
}
#[async_trait] #[async_trait]
pub trait ControlPlaneStore: Send + Sync { pub trait ControlPlaneStore: Send + Sync {
async fn schema_version(&self) -> Result<i64>; async fn schema_version(&self) -> Result<i64>;
async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()>; async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()>;
async fn get_workspace(&self, workspace_id: &str) -> Result<Option<WorkspaceRecord>>; async fn get_workspace(&self, workspace_id: &str) -> Result<Option<WorkspaceRecord>>;
fn upsert_worker_registry(&self, record: &WorkerRegistryRecord) -> Result<()>;
fn get_worker_registry(
&self,
workspace_id: &str,
worker_id: &str,
) -> Result<Option<WorkerRegistryRecord>>;
fn get_worker_registry_by_runtime(
&self,
workspace_id: &str,
runtime_id: &str,
runtime_worker_id: &str,
) -> Result<Option<WorkerRegistryRecord>>;
fn list_worker_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkerRegistryRecord>>;
fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()>;
fn get_workdir_registry(
&self,
workspace_id: &str,
workdir_id: &str,
) -> Result<Option<WorkdirRegistryRecord>>;
fn list_workdir_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkdirRegistryRecord>>;
fn list_managed_workdir_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkdirRegistryRecord>>;
fn upsert_worker_workdir_link(&self, record: &WorkerWorkdirLinkRecord) -> Result<()>;
fn list_worker_workdir_links(
&self,
workspace_id: &str,
worker_id: &str,
) -> Result<Vec<WorkerWorkdirLinkRecord>>;
} }
#[derive(Clone)] #[derive(Clone)]
@ -131,6 +224,355 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
.map_err(Error::from) .map_err(Error::from)
}) })
} }
fn upsert_worker_registry(&self, record: &WorkerRegistryRecord) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"INSERT INTO worker_registry (
workspace_id, worker_id, runtime_id, runtime_worker_id, display_name, profile,
lifecycle_state, retention_state, transcript_ref, session_ref, summary_ref,
diagnostics_ref, created_at, updated_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
ON CONFLICT(workspace_id, worker_id) DO UPDATE SET
runtime_id = excluded.runtime_id,
runtime_worker_id = excluded.runtime_worker_id,
display_name = excluded.display_name,
profile = excluded.profile,
lifecycle_state = excluded.lifecycle_state,
retention_state = CASE
WHEN worker_registry.retention_state = 'pinned' AND excluded.retention_state = 'normal'
THEN worker_registry.retention_state
ELSE excluded.retention_state
END,
transcript_ref = excluded.transcript_ref,
session_ref = excluded.session_ref,
summary_ref = excluded.summary_ref,
diagnostics_ref = excluded.diagnostics_ref,
updated_at = excluded.updated_at"#,
params![
record.workspace_id,
record.worker_id,
record.runtime_id,
record.runtime_worker_id,
record.display_name,
record.profile,
record.lifecycle_state,
record.retention_state,
record.transcript_ref,
record.session_ref,
record.summary_ref,
record.diagnostics_ref,
record.created_at,
record.updated_at,
],
)?;
Ok(())
})
}
fn get_worker_registry(
&self,
workspace_id: &str,
worker_id: &str,
) -> Result<Option<WorkerRegistryRecord>> {
self.with_conn(|conn| {
conn.query_row(
worker_registry_select_sql("WHERE workspace_id = ?1 AND worker_id = ?2").as_str(),
params![workspace_id, worker_id],
read_worker_registry_record,
)
.optional()
.map_err(Error::from)
})
}
fn get_worker_registry_by_runtime(
&self,
workspace_id: &str,
runtime_id: &str,
runtime_worker_id: &str,
) -> Result<Option<WorkerRegistryRecord>> {
self.with_conn(|conn| {
conn.query_row(
worker_registry_select_sql(
"WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3",
)
.as_str(),
params![workspace_id, runtime_id, runtime_worker_id],
read_worker_registry_record,
)
.optional()
.map_err(Error::from)
})
}
fn list_worker_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkerRegistryRecord>> {
self.with_conn(|conn| {
let sql = worker_registry_select_sql(
"WHERE workspace_id = ?1 ORDER BY updated_at DESC LIMIT ?2",
);
let mut stmt = conn.prepare(sql.as_str())?;
let rows = stmt.query_map(
params![workspace_id, limit as i64],
read_worker_registry_record,
)?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"INSERT INTO workdir_registry (
workspace_id, workdir_id, runtime_id, repository_id, selector, resolved_commit,
materialization_status, cleanliness, management_kind, created_at, updated_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
ON CONFLICT(workspace_id, workdir_id) DO UPDATE SET
runtime_id = excluded.runtime_id,
repository_id = excluded.repository_id,
selector = excluded.selector,
resolved_commit = excluded.resolved_commit,
materialization_status = excluded.materialization_status,
cleanliness = excluded.cleanliness,
management_kind = excluded.management_kind,
updated_at = excluded.updated_at"#,
params![
record.workspace_id,
record.workdir_id,
record.runtime_id,
record.repository_id,
record.selector,
record.resolved_commit,
record.materialization_status,
record.cleanliness,
record.management_kind,
record.created_at,
record.updated_at,
],
)?;
Ok(())
})
}
fn get_workdir_registry(
&self,
workspace_id: &str,
workdir_id: &str,
) -> Result<Option<WorkdirRegistryRecord>> {
self.with_conn(|conn| {
conn.query_row(
workdir_registry_select_sql("WHERE workspace_id = ?1 AND workdir_id = ?2").as_str(),
params![workspace_id, workdir_id],
read_workdir_registry_record,
)
.optional()
.map_err(Error::from)
})
}
fn list_workdir_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkdirRegistryRecord>> {
self.with_conn(|conn| {
let sql = workdir_registry_select_sql(
"WHERE workspace_id = ?1 ORDER BY updated_at DESC LIMIT ?2",
);
let mut stmt = conn.prepare(sql.as_str())?;
let rows = stmt.query_map(
params![workspace_id, limit as i64],
read_workdir_registry_record,
)?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
fn list_managed_workdir_registry(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<WorkdirRegistryRecord>> {
self.with_conn(|conn| {
let sql = workdir_registry_select_sql(
"WHERE workspace_id = ?1 AND management_kind = 'backend_managed' ORDER BY updated_at DESC LIMIT ?2",
);
let mut stmt = conn.prepare(sql.as_str())?;
let rows = stmt.query_map(params![workspace_id, limit as i64], read_workdir_registry_record)?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
fn upsert_worker_workdir_link(&self, record: &WorkerWorkdirLinkRecord) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"INSERT INTO worker_workdir_links (
workspace_id, worker_id, workdir_id, role, linked_at, unlinked_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT(workspace_id, worker_id, workdir_id, role) DO UPDATE SET
linked_at = excluded.linked_at,
unlinked_at = excluded.unlinked_at"#,
params![
record.workspace_id,
record.worker_id,
record.workdir_id,
record.role,
record.linked_at,
record.unlinked_at,
],
)?;
Ok(())
})
}
fn list_worker_workdir_links(
&self,
workspace_id: &str,
worker_id: &str,
) -> Result<Vec<WorkerWorkdirLinkRecord>> {
self.with_conn(|conn| {
let mut stmt = conn.prepare(
r#"SELECT workspace_id, worker_id, workdir_id, role, linked_at, unlinked_at
FROM worker_workdir_links
WHERE workspace_id = ?1 AND worker_id = ?2 AND unlinked_at IS NULL
ORDER BY linked_at DESC"#,
)?;
let rows = stmt.query_map(params![workspace_id, worker_id], |row| {
Ok(WorkerWorkdirLinkRecord {
workspace_id: row.get(0)?,
worker_id: row.get(1)?,
workdir_id: row.get(2)?,
role: row.get(3)?,
linked_at: row.get(4)?,
unlinked_at: row.get(5)?,
})
})?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
}
fn worker_registry_select_sql(where_clause: &str) -> String {
format!(
"SELECT workspace_id, worker_id, runtime_id, runtime_worker_id, display_name, profile, \
lifecycle_state, retention_state, transcript_ref, session_ref, summary_ref, diagnostics_ref, \
created_at, updated_at FROM worker_registry {where_clause}"
)
}
fn read_worker_registry_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<WorkerRegistryRecord> {
Ok(WorkerRegistryRecord {
workspace_id: row.get(0)?,
worker_id: row.get(1)?,
runtime_id: row.get(2)?,
runtime_worker_id: row.get(3)?,
display_name: row.get(4)?,
profile: row.get(5)?,
lifecycle_state: row.get(6)?,
retention_state: row.get(7)?,
transcript_ref: row.get(8)?,
session_ref: row.get(9)?,
summary_ref: row.get(10)?,
diagnostics_ref: row.get(11)?,
created_at: row.get(12)?,
updated_at: row.get(13)?,
})
}
fn workdir_registry_select_sql(where_clause: &str) -> String {
format!(
"SELECT workspace_id, workdir_id, runtime_id, repository_id, selector, resolved_commit, \
materialization_status, cleanliness, management_kind, created_at, updated_at \
FROM workdir_registry {where_clause}"
)
}
fn read_workdir_registry_record(
row: &rusqlite::Row<'_>,
) -> rusqlite::Result<WorkdirRegistryRecord> {
Ok(WorkdirRegistryRecord {
workspace_id: row.get(0)?,
workdir_id: row.get(1)?,
runtime_id: row.get(2)?,
repository_id: row.get(3)?,
selector: row.get(4)?,
resolved_commit: row.get(5)?,
materialization_status: row.get(6)?,
cleanliness: row.get(7)?,
management_kind: row.get(8)?,
created_at: row.get(9)?,
updated_at: row.get(10)?,
})
}
fn create_worker_workdir_registry_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
CREATE TABLE IF NOT EXISTS worker_registry (
workspace_id TEXT NOT NULL,
worker_id TEXT NOT NULL,
runtime_id TEXT NOT NULL,
runtime_worker_id TEXT NOT NULL,
display_name TEXT NOT NULL,
profile TEXT,
lifecycle_state TEXT NOT NULL,
retention_state TEXT NOT NULL CHECK (retention_state IN ('normal', 'pinned')),
transcript_ref TEXT,
session_ref TEXT,
summary_ref TEXT,
diagnostics_ref TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (workspace_id, worker_id),
UNIQUE (workspace_id, runtime_id, runtime_worker_id),
FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE
);
CREATE TABLE IF NOT EXISTS workdir_registry (
workspace_id TEXT NOT NULL,
workdir_id TEXT NOT NULL,
runtime_id TEXT NOT NULL,
repository_id TEXT NOT NULL,
selector TEXT,
resolved_commit TEXT,
materialization_status TEXT NOT NULL CHECK (materialization_status IN ('pending', 'present', 'missing', 'removed', 'failed')),
cleanliness TEXT NOT NULL CHECK (cleanliness IN ('clean', 'dirty', 'unknown')),
management_kind TEXT NOT NULL CHECK (management_kind IN ('backend_managed', 'runtime_unmanaged')),
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (workspace_id, workdir_id),
FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE
);
CREATE TABLE IF NOT EXISTS worker_workdir_links (
workspace_id TEXT NOT NULL,
worker_id TEXT NOT NULL,
workdir_id TEXT NOT NULL,
role TEXT NOT NULL,
linked_at TEXT NOT NULL,
unlinked_at TEXT,
PRIMARY KEY (workspace_id, worker_id, workdir_id, role),
FOREIGN KEY (workspace_id, worker_id) REFERENCES worker_registry(workspace_id, worker_id) ON DELETE CASCADE,
FOREIGN KEY (workspace_id, workdir_id) REFERENCES workdir_registry(workspace_id, workdir_id) ON DELETE CASCADE
);
CREATE INDEX IF NOT EXISTS idx_worker_registry_workspace_updated
ON worker_registry(workspace_id, updated_at DESC);
CREATE INDEX IF NOT EXISTS idx_workdir_registry_workspace_updated
ON workdir_registry(workspace_id, updated_at DESC);
CREATE INDEX IF NOT EXISTS idx_worker_workdir_links_worker
ON worker_workdir_links(workspace_id, worker_id, linked_at DESC);
"#,
)?;
Ok(())
} }
fn configure_sqlite(conn: &Connection) -> Result<()> { fn configure_sqlite(conn: &Connection) -> Result<()> {
@ -488,7 +930,7 @@ mod tests {
let db = dir.path().join("control-plane.sqlite"); let db = dir.path().join("control-plane.sqlite");
let store = SqliteWorkspaceStore::open(&db).unwrap(); let store = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 2); assert_eq!(store.schema_version().await.unwrap(), 3);
let record = WorkspaceRecord { let record = WorkspaceRecord {
workspace_id: "local-dev".to_string(), workspace_id: "local-dev".to_string(),
@ -500,7 +942,7 @@ mod tests {
store.upsert_workspace(&record).await.unwrap(); store.upsert_workspace(&record).await.unwrap();
let reopened = SqliteWorkspaceStore::open(&db).unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(reopened.schema_version().await.unwrap(), 2); assert_eq!(reopened.schema_version().await.unwrap(), 3);
assert_eq!( assert_eq!(
reopened.get_workspace("local-dev").await.unwrap(), reopened.get_workspace("local-dev").await.unwrap(),
Some(record) Some(record)
@ -527,6 +969,9 @@ mod tests {
"ticket_worker_links", "ticket_worker_links",
"artifacts", "artifacts",
"audit_events", "audit_events",
"worker_registry",
"workdir_registry",
"worker_workdir_links",
] { ] {
assert!( assert!(
tables.contains(expected), tables.contains(expected),
@ -688,7 +1133,7 @@ mod tests {
.unwrap(); .unwrap();
let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 2); assert_eq!(store.schema_version().await.unwrap(), 3);
store store
.with_conn(|conn| { .with_conn(|conn| {
@ -773,6 +1218,115 @@ mod tests {
); );
} }
#[tokio::test]
async fn worker_workdir_registry_round_trips_and_preserves_pinned_retention() {
let temp = tempfile::tempdir().unwrap();
let db = temp.path().join("workspace.db");
let store = SqliteWorkspaceStore::open(&db).unwrap();
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
display_name: "Local Dev".to_string(),
state: "active".to_string(),
created_at: "1".to_string(),
updated_at: "1".to_string(),
};
store.upsert_workspace(&workspace).await.unwrap();
let worker = WorkerRegistryRecord {
workspace_id: "local-dev".to_string(),
worker_id: "embedded/browser-1".to_string(),
runtime_id: "embedded".to_string(),
runtime_worker_id: "browser-1".to_string(),
display_name: "Browser 1".to_string(),
profile: Some("builtin:companion".to_string()),
lifecycle_state: "idle".to_string(),
retention_state: "pinned".to_string(),
transcript_ref: Some("runtime://embedded/workers/browser-1/transcript".to_string()),
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: "2".to_string(),
updated_at: "2".to_string(),
};
store.upsert_worker_registry(&worker).unwrap();
let mut runtime_sync_worker = worker.clone();
runtime_sync_worker.lifecycle_state = "running".to_string();
runtime_sync_worker.retention_state = "normal".to_string();
runtime_sync_worker.updated_at = "5".to_string();
store.upsert_worker_registry(&runtime_sync_worker).unwrap();
let mut expected_worker = worker.clone();
expected_worker.lifecycle_state = "running".to_string();
expected_worker.updated_at = "5".to_string();
let workdir = WorkdirRegistryRecord {
workspace_id: "local-dev".to_string(),
workdir_id: "backend-2-repo".to_string(),
runtime_id: "embedded".to_string(),
repository_id: "repo".to_string(),
selector: Some("develop".to_string()),
resolved_commit: Some("abcdef".to_string()),
materialization_status: "removed".to_string(),
cleanliness: "clean".to_string(),
management_kind: "backend_managed".to_string(),
created_at: "2".to_string(),
updated_at: "3".to_string(),
};
store.upsert_workdir_registry(&workdir).unwrap();
let unmanaged_workdir = WorkdirRegistryRecord {
workspace_id: "local-dev".to_string(),
workdir_id: "runtime-direct".to_string(),
runtime_id: "embedded".to_string(),
repository_id: "repo".to_string(),
selector: Some("feature".to_string()),
resolved_commit: Some("123456".to_string()),
materialization_status: "present".to_string(),
cleanliness: "unknown".to_string(),
management_kind: "runtime_unmanaged".to_string(),
created_at: "3".to_string(),
updated_at: "4".to_string(),
};
store.upsert_workdir_registry(&unmanaged_workdir).unwrap();
let link = WorkerWorkdirLinkRecord {
workspace_id: "local-dev".to_string(),
worker_id: worker.worker_id.clone(),
workdir_id: workdir.workdir_id.clone(),
role: "primary_cwd".to_string(),
linked_at: "4".to_string(),
unlinked_at: None,
};
store.upsert_worker_workdir_link(&link).unwrap();
assert_eq!(
store
.get_worker_registry_by_runtime("local-dev", "embedded", "browser-1")
.unwrap(),
Some(expected_worker.clone())
);
assert_eq!(
store
.get_workdir_registry("local-dev", "backend-2-repo")
.unwrap(),
Some(workdir.clone())
);
assert_eq!(
store.list_workdir_registry("local-dev", 10).unwrap(),
vec![unmanaged_workdir.clone(), workdir.clone()]
);
assert_eq!(
store
.list_managed_workdir_registry("local-dev", 10)
.unwrap(),
vec![workdir]
);
assert_eq!(
store
.list_worker_workdir_links("local-dev", "embedded/browser-1")
.unwrap(),
vec![link]
);
}
fn table_names(conn: &Connection) -> BTreeSet<String> { fn table_names(conn: &Connection) -> BTreeSet<String> {
let mut stmt = conn let mut stmt = conn
.prepare( .prepare(

View File

@ -124,6 +124,7 @@ export type WorkingDirectorySummary = {
resolved_commit: string; resolved_commit: string;
resolved_tree?: string | null; resolved_tree?: string | null;
status: string; status: string;
management_kind?: "backend_managed" | "runtime_unmanaged" | string | null;
cleanup_policy: string; cleanup_policy: string;
cleanup_target: { cleanup_target: {
kind: string; kind: string;