6711 lines
256 KiB
Rust
6711 lines
256 KiB
Rust
use crate::catalog::{
|
|
ConfigBundleRef, CreateWorkerRequest, ProfileSelector, ProfileSourceArchiveSource,
|
|
RepositoryRefObservation, RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck,
|
|
WorkerRestoreIntent, WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest,
|
|
WorkingDirectoryRequest, WorkingDirectoryStatus as CatalogWorkingDirectoryStatus,
|
|
WorkspaceApiRef,
|
|
};
|
|
use crate::config_bundle::{
|
|
ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary, validate_config_bundle,
|
|
validate_config_bundle_ref,
|
|
};
|
|
use crate::diagnostics::{DiagnosticSeverity, RuntimeDiagnostic};
|
|
use crate::error::RuntimeError;
|
|
use crate::execution::WorkerExecutionRestoreRequest;
|
|
use crate::execution::{
|
|
WorkerExecutionBackend, WorkerExecutionBackendRef, WorkerExecutionHandle,
|
|
WorkerExecutionOperation, WorkerExecutionResult, WorkerExecutionSpawnRequest,
|
|
WorkerExecutionSpawnResult, WorkspaceConfigFetchRequest, WorkspaceConfigFetchResult,
|
|
};
|
|
#[cfg(feature = "fs-store")]
|
|
use crate::fs_store::{
|
|
FsRuntimeStore, FsRuntimeStoreOptions, PersistedRuntimeState, PersistedWorkerExecution,
|
|
PersistedWorkerExecutionBinding, PersistedWorkerRecord,
|
|
};
|
|
use crate::identity::{WorkerId, WorkerRef};
|
|
use crate::interaction::{WorkerInput, WorkerInputKind, WorkerInteractionAck};
|
|
use crate::management::{
|
|
RuntimeBackendKind, RuntimeOptions, RuntimeStatus, RuntimeSummary, WorkerDeleteResult,
|
|
};
|
|
#[cfg(feature = "ws-server")]
|
|
use crate::observation::{WorkerObservationCursor, WorkerObservationEvent};
|
|
use crate::resource::{
|
|
BackendResourceClient, BackendResourceError, BackendResourceFetchRequest, BackendResourceKind,
|
|
REPOSITORY_SSH_ACCESS_CONTENT_TYPE, RepositorySshAccessSecret,
|
|
};
|
|
#[cfg(feature = "fs-store")]
|
|
use crate::retention::{
|
|
FsWorkerRetentionProvider, WorkerRetentionExecutionRequest, WorkerRetentionExecutionResult,
|
|
WorkerRetentionInventory, WorkerRetentionInventorySnapshot, WorkerRetentionProvider,
|
|
};
|
|
use protocol::subscription::{
|
|
EventSubscriptionSelector, SubscriptionEventPayload, SubscriptionSnapshot,
|
|
SubscriptionValidationError, SubscriptionWorkdirId, SubscriptionWorker, SubscriptionWorkerId,
|
|
SubscriptionWorkerState,
|
|
};
|
|
use protocol::{Event, Method};
|
|
use std::collections::BTreeMap;
|
|
#[cfg(feature = "ws-server")]
|
|
use std::collections::VecDeque;
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
|
use std::sync::{Arc, Mutex, MutexGuard, Weak};
|
|
#[cfg(feature = "ws-server")]
|
|
use tokio::sync::broadcast;
|
|
use tokio::sync::mpsc;
|
|
use uuid::Uuid;
|
|
|
|
/// Workspace-scoped Runtime authorization context supplied by a trusted backend.
|
|
#[derive(Clone, Debug, Eq, PartialEq)]
|
|
pub struct RuntimeWorkspaceScope {
|
|
pub workspace_id: String,
|
|
pub server_id: String,
|
|
}
|
|
|
|
impl RuntimeWorkspaceScope {
|
|
pub fn new(workspace_id: impl Into<String>, server_id: impl Into<String>) -> Self {
|
|
Self {
|
|
workspace_id: workspace_id.into(),
|
|
server_id: server_id.into(),
|
|
}
|
|
}
|
|
}
|
|
|
|
const SUBSCRIPTION_QUEUE_CAPACITY: usize = 256;
|
|
|
|
#[derive(Clone, Debug)]
|
|
pub struct RuntimeSubscriptionUpdate {
|
|
pub subject_revision: u64,
|
|
pub payload: SubscriptionEventPayload,
|
|
}
|
|
|
|
#[derive(Debug, thiserror::Error)]
|
|
pub enum RuntimeSubscriptionRecvError {
|
|
#[error("Runtime event subscription lagged and requires a fresh snapshot")]
|
|
Lagged,
|
|
#[error("Runtime event subscription closed")]
|
|
Closed,
|
|
}
|
|
|
|
/// A gap-free snapshot/live subscription owned by one Runtime connection.
|
|
/// Dropping the subscription removes its bounded producer queue.
|
|
pub struct RuntimeEventSelectorSubscription {
|
|
subscription_id: u64,
|
|
selector: EventSubscriptionSelector,
|
|
snapshot_revision: u64,
|
|
snapshot: SubscriptionSnapshot,
|
|
receiver: mpsc::Receiver<RuntimeSubscriptionUpdate>,
|
|
lagged: Arc<AtomicBool>,
|
|
runtime: Weak<Mutex<RuntimeState>>,
|
|
}
|
|
|
|
impl RuntimeEventSelectorSubscription {
|
|
pub fn selector(&self) -> &EventSubscriptionSelector {
|
|
&self.selector
|
|
}
|
|
|
|
pub fn snapshot_revision(&self) -> u64 {
|
|
self.snapshot_revision
|
|
}
|
|
|
|
pub fn snapshot(&self) -> &SubscriptionSnapshot {
|
|
&self.snapshot
|
|
}
|
|
|
|
pub async fn recv(
|
|
&mut self,
|
|
) -> Result<RuntimeSubscriptionUpdate, RuntimeSubscriptionRecvError> {
|
|
if self.lagged.load(Ordering::Acquire) {
|
|
self.receiver.close();
|
|
return Err(RuntimeSubscriptionRecvError::Lagged);
|
|
}
|
|
match self.receiver.recv().await {
|
|
Some(_) if self.lagged.load(Ordering::Acquire) => {
|
|
self.receiver.close();
|
|
Err(RuntimeSubscriptionRecvError::Lagged)
|
|
}
|
|
Some(update) => Ok(update),
|
|
None if self.lagged.load(Ordering::Acquire) => {
|
|
Err(RuntimeSubscriptionRecvError::Lagged)
|
|
}
|
|
None => Err(RuntimeSubscriptionRecvError::Closed),
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Drop for RuntimeEventSelectorSubscription {
|
|
fn drop(&mut self) {
|
|
let Some(runtime) = self.runtime.upgrade() else {
|
|
return;
|
|
};
|
|
if let Ok(mut state) = runtime.lock() {
|
|
state.subscriptions.remove(&self.subscription_id);
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Concrete embedded Runtime domain entity.
|
|
///
|
|
/// The default implementation is memory-backed and tools/provider-less by
|
|
/// design. An optional `fs-store` feature adds filesystem persistence while
|
|
/// preserving the same typed authority boundary. It can later be adapted by
|
|
/// backend registries or web servers without making sockets, sessions, or paths
|
|
/// public authority.
|
|
#[derive(Clone, Debug)]
|
|
pub struct Runtime {
|
|
inner: Arc<Mutex<RuntimeState>>,
|
|
worker_operations: Arc<Mutex<BTreeMap<WorkerId, Arc<Mutex<()>>>>>,
|
|
}
|
|
|
|
impl Runtime {
|
|
/// Create a memory-backed Runtime with generated identity.
|
|
pub fn new_memory() -> Self {
|
|
Self::with_options(RuntimeOptions::default())
|
|
}
|
|
|
|
/// Create a memory-backed Runtime with explicit options.
|
|
pub fn with_options(options: RuntimeOptions) -> Self {
|
|
let state = RuntimeState::new(options.display_name);
|
|
Self {
|
|
inner: Arc::new(Mutex::new(state)),
|
|
worker_operations: Arc::new(Mutex::new(BTreeMap::new())),
|
|
}
|
|
}
|
|
|
|
/// Create a memory-backed Runtime with an attached execution backend.
|
|
pub fn with_execution_backend(
|
|
options: RuntimeOptions,
|
|
backend: Arc<dyn WorkerExecutionBackend>,
|
|
) -> Result<Self, RuntimeError> {
|
|
let runtime = Self::with_options(options);
|
|
runtime.install_execution_backend(backend)?;
|
|
Ok(runtime)
|
|
}
|
|
|
|
pub fn install_backend_resource_client(
|
|
&self,
|
|
client: Arc<dyn BackendResourceClient>,
|
|
) -> Result<(), RuntimeError> {
|
|
self.lock()?.backend_resource_client = Some(BackendResourceClientRef(client));
|
|
Ok(())
|
|
}
|
|
|
|
pub fn install_workspace_backend_resource_client(
|
|
&self,
|
|
workspace_id: impl Into<String>,
|
|
client: Arc<dyn BackendResourceClient>,
|
|
) -> Result<(), RuntimeError> {
|
|
let workspace_id = workspace_id.into();
|
|
if workspace_id.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Backend resource client Workspace id is empty".to_string(),
|
|
));
|
|
}
|
|
self.lock()?
|
|
.workspace_backend_resource_clients
|
|
.insert(workspace_id, BackendResourceClientRef(client));
|
|
Ok(())
|
|
}
|
|
|
|
/// Create or restore a filesystem-backed Runtime.
|
|
///
|
|
/// The store is scoped by `options.root`; if the directory already exists,
|
|
/// persisted state is loaded and validated. If it does not exist, a fresh
|
|
/// Runtime is initialized and durable files are created before return.
|
|
#[cfg(feature = "fs-store")]
|
|
pub fn with_fs_store(options: FsRuntimeStoreOptions) -> Result<Self, RuntimeError> {
|
|
Self::with_fs_store_inner(options, None)
|
|
}
|
|
|
|
/// Create or restore a filesystem-backed Runtime with an execution backend.
|
|
#[cfg(feature = "fs-store")]
|
|
pub fn with_fs_store_and_execution_backend(
|
|
options: FsRuntimeStoreOptions,
|
|
backend: Arc<dyn WorkerExecutionBackend>,
|
|
) -> Result<Self, RuntimeError> {
|
|
Self::with_fs_store_inner(options, Some(WorkerExecutionBackendRef::new(backend)?))
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn with_fs_store_inner(
|
|
options: FsRuntimeStoreOptions,
|
|
execution_backend: Option<WorkerExecutionBackendRef>,
|
|
) -> Result<Self, RuntimeError> {
|
|
let opened = FsRuntimeStore::open_or_create(options.root, &options.runtime_id)?;
|
|
let mut state = if let Some(persisted) = opened.state {
|
|
RuntimeState::from_persisted(persisted, opened.store)?
|
|
} else {
|
|
let state = RuntimeState::new_fs_backed(options.display_name, opened.store);
|
|
state.persist_runtime_snapshot()?;
|
|
state
|
|
};
|
|
state.execution_backend = execution_backend;
|
|
let runtime = Self {
|
|
inner: Arc::new(Mutex::new(state)),
|
|
worker_operations: Arc::new(Mutex::new(BTreeMap::new())),
|
|
};
|
|
runtime.restore_persisted_worker_executions()?;
|
|
Ok(runtime)
|
|
}
|
|
|
|
/// Management-plane summary.
|
|
pub fn summary(&self) -> Result<RuntimeSummary, RuntimeError> {
|
|
let state = self.lock()?;
|
|
let mut active_worker_count = 0;
|
|
let mut stopped_worker_count = 0;
|
|
for worker in state.workers.values() {
|
|
match worker.status {
|
|
WorkerStatus::Idle | WorkerStatus::Running | WorkerStatus::Paused => {
|
|
active_worker_count += 1;
|
|
}
|
|
WorkerStatus::Stopped => stopped_worker_count += 1,
|
|
}
|
|
}
|
|
|
|
Ok(RuntimeSummary {
|
|
display_name: state.display_name.clone(),
|
|
backend: state.backend,
|
|
status: state.status,
|
|
worker_count: state.workers.len(),
|
|
active_worker_count,
|
|
stopped_worker_count,
|
|
diagnostic_count: state.diagnostics.len(),
|
|
os: std::env::consts::OS.to_string(),
|
|
arch: std::env::consts::ARCH.to_string(),
|
|
worker_creation_available: state.execution_backend.is_some(),
|
|
})
|
|
}
|
|
|
|
/// Current Runtime lifecycle state.
|
|
pub fn status(&self) -> Result<RuntimeStatus, RuntimeError> {
|
|
Ok(self.lock()?.status)
|
|
}
|
|
|
|
/// Store a backend-synced Profile/config bundle for later Worker creation.
|
|
pub fn store_config_bundle(
|
|
&self,
|
|
bundle: ConfigBundle,
|
|
) -> Result<ConfigBundleAvailability, RuntimeError> {
|
|
validate_config_bundle(&bundle)?;
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
let reference = ConfigBundleRef {
|
|
id: bundle.metadata.id.clone(),
|
|
digest: bundle.metadata.digest.clone(),
|
|
};
|
|
let summary = bundle.summary();
|
|
if let Some(existing) = state.config_bundles.get(&bundle.metadata.id) {
|
|
if existing.metadata.digest != bundle.metadata.digest {
|
|
return Err(RuntimeError::ConfigBundleDigestMismatch {
|
|
bundle_id: bundle.metadata.id.clone(),
|
|
expected_digest: existing.metadata.digest.clone(),
|
|
actual_digest: bundle.metadata.digest.clone(),
|
|
});
|
|
}
|
|
return Ok(ConfigBundleAvailability {
|
|
reference,
|
|
summary: existing.summary(),
|
|
});
|
|
}
|
|
state
|
|
.config_bundles
|
|
.insert(bundle.metadata.id.clone(), bundle);
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(ConfigBundleAvailability { reference, summary })
|
|
}
|
|
|
|
/// List synced config bundles known to this Runtime.
|
|
pub fn list_config_bundles(&self) -> Result<Vec<ConfigBundleSummary>, RuntimeError> {
|
|
Ok(self
|
|
.lock()?
|
|
.config_bundles
|
|
.values()
|
|
.map(ConfigBundle::summary)
|
|
.collect())
|
|
}
|
|
|
|
/// Validate that a config bundle reference is present and digest-matched.
|
|
pub fn check_config_bundle(
|
|
&self,
|
|
reference: &ConfigBundleRef,
|
|
) -> Result<ConfigBundleAvailability, RuntimeError> {
|
|
let state = self.lock()?;
|
|
state.check_config_bundle_ref(reference)
|
|
}
|
|
|
|
/// Notify the execution backend of the Workspace's current immutable
|
|
/// Prompt projection. The Runtime keeps this cache outside persisted Worker
|
|
/// restore authority.
|
|
pub fn observe_workspace_prompt_projection(
|
|
&self,
|
|
projection: worker::WorkspacePromptProjection,
|
|
) -> Result<(), RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message:
|
|
"Workspace Prompt projection notification requires an execution backend"
|
|
.to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.observe_workspace_prompt_projection(projection)
|
|
.map_err(|message| RuntimeError::ExecutionBackendUnavailable { message })
|
|
}
|
|
|
|
/// Stop the Runtime. v0 keeps data readable after stop, but rejects new
|
|
/// create/send/worker lifecycle mutations.
|
|
pub fn stop_runtime(&self) -> Result<(), RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
if state.status == RuntimeStatus::Stopped {
|
|
return Ok(());
|
|
}
|
|
state.status = RuntimeStatus::Stopped;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_workers()?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Create a Runtime-owned working directory through the attached execution backend.
|
|
pub fn create_working_directory(
|
|
&self,
|
|
request: WorkingDirectoryRequest,
|
|
) -> Result<CatalogWorkingDirectoryStatus, RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "working directory creation requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.create_working_directory(&request)
|
|
.map_err(RuntimeError::from)
|
|
}
|
|
|
|
pub async fn create_working_directory_from_resource(
|
|
&self,
|
|
mut request: WorkingDirectoryRequest,
|
|
) -> Result<CatalogWorkingDirectoryStatus, RuntimeError> {
|
|
if let Some(ssh) = request
|
|
.materialization
|
|
.as_mut()
|
|
.and_then(|materialization| materialization.ssh.as_mut())
|
|
{
|
|
self.resolve_repository_access_resource(ssh).await?;
|
|
}
|
|
self.create_working_directory(request)
|
|
}
|
|
|
|
pub async fn observe_repository_ref_from_resource(
|
|
&self,
|
|
mut request: RepositoryRefObservationRequest,
|
|
) -> Result<RepositoryRefObservation, RuntimeError> {
|
|
if let Some(ssh) = request
|
|
.materialization
|
|
.as_mut()
|
|
.and_then(|materialization| materialization.ssh.as_mut())
|
|
{
|
|
self.resolve_repository_access_resource(ssh).await?;
|
|
}
|
|
self.observe_repository_ref(request)
|
|
}
|
|
|
|
pub fn observe_repository_ref(
|
|
&self,
|
|
request: RepositoryRefObservationRequest,
|
|
) -> Result<RepositoryRefObservation, RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "Repository ref observation requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.observe_repository_ref(&request)
|
|
.map_err(RuntimeError::from)
|
|
}
|
|
|
|
pub fn authorize_working_directory_repository_access(
|
|
&self,
|
|
request: WorkingDirectoryRepositoryAccessRequest,
|
|
) -> Result<(), RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "working directory Repository access requires an execution backend"
|
|
.to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.authorize_working_directory_repository_access(&request)
|
|
.map_err(RuntimeError::from)
|
|
}
|
|
|
|
async fn resolve_repository_access_resource(
|
|
&self,
|
|
ssh: &mut crate::catalog::RepositorySshMaterializationAccess,
|
|
) -> Result<(), RuntimeError> {
|
|
if ssh.credential_candidates.is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Repository SSH access requires at least one credential candidate".to_string(),
|
|
));
|
|
}
|
|
if ssh
|
|
.credential_candidates
|
|
.iter()
|
|
.all(|candidate| !candidate.private_key.expose().is_empty())
|
|
&& !ssh.known_hosts_entry.expose().is_empty()
|
|
{
|
|
return Ok(());
|
|
}
|
|
let (client, runtime_id) = {
|
|
let state = self.lock()?;
|
|
let client = state
|
|
.workspace_backend_resource_clients
|
|
.get(&ssh.secret_resource.workspace_id)
|
|
.cloned()
|
|
.or_else(|| state.backend_resource_client.clone())
|
|
.ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(format!(
|
|
"Backend Repository access resource client is unavailable for Workspace `{}`",
|
|
ssh.secret_resource.workspace_id
|
|
))
|
|
})?;
|
|
let runtime_id = state.runtime_identity.clone().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest("Runtime identity is unavailable".to_string())
|
|
})?;
|
|
(client, runtime_id)
|
|
};
|
|
tracing::info!(
|
|
target: "yoi::repository_access",
|
|
event = "repository_access_resource_fetch_started",
|
|
workspace_id = %ssh.secret_resource.workspace_id,
|
|
resource_id = %ssh.secret_resource.resource_id,
|
|
runtime_id = %runtime_id,
|
|
credential_candidate_count = ssh.credential_candidates.len(),
|
|
"fetching Repository SSH access resource from Workspace Backend"
|
|
);
|
|
let mut response = client
|
|
.0
|
|
.fetch_resource(BackendResourceFetchRequest {
|
|
handle: ssh.secret_resource.clone(),
|
|
runtime_id: runtime_id.clone(),
|
|
worker_id: None,
|
|
audit_correlation_id: ssh.secret_resource.audit_correlation_id.clone(),
|
|
})
|
|
.await
|
|
.map_err(repository_resource_error)?;
|
|
if response.kind != BackendResourceKind::RepositorySshAccess
|
|
|| response.content_type != REPOSITORY_SSH_ACCESS_CONTENT_TYPE
|
|
|| response.resource_id != ssh.secret_resource.resource_id
|
|
|| response.digest != ssh.secret_resource.digest
|
|
|| response.bytes.len() as u64 > ssh.secret_resource.max_bytes
|
|
{
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Backend Repository SSH access resource response was invalid".to_string(),
|
|
));
|
|
}
|
|
let secret = serde_json::from_slice::<RepositorySshAccessSecret>(&response.bytes);
|
|
response.bytes.fill(0);
|
|
let mut secret = secret.map_err(|_| {
|
|
RuntimeError::InvalidRequest(
|
|
"Backend Repository SSH access resource payload was invalid".to_string(),
|
|
)
|
|
})?;
|
|
if secret.credential_candidates.len() != ssh.credential_candidates.len()
|
|
|| secret
|
|
.credential_candidates
|
|
.iter()
|
|
.zip(&ssh.credential_candidates)
|
|
.any(|(secret, metadata)| {
|
|
secret.credential_id != metadata.credential_id
|
|
|| secret.credential_revision != metadata.credential_revision
|
|
})
|
|
{
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Backend Repository SSH access resource credential metadata was invalid"
|
|
.to_string(),
|
|
));
|
|
}
|
|
for (candidate, secret) in ssh
|
|
.credential_candidates
|
|
.iter_mut()
|
|
.zip(&mut secret.credential_candidates)
|
|
{
|
|
candidate.private_key =
|
|
crate::catalog::SensitiveString::new(std::mem::take(&mut secret.private_key));
|
|
}
|
|
ssh.known_hosts_entry =
|
|
crate::catalog::SensitiveString::new(std::mem::take(&mut secret.known_hosts_entry));
|
|
tracing::info!(
|
|
target: "yoi::repository_access",
|
|
event = "repository_access_resource_fetch_succeeded",
|
|
workspace_id = %ssh.secret_resource.workspace_id,
|
|
resource_id = %ssh.secret_resource.resource_id,
|
|
runtime_id = %runtime_id,
|
|
credential_candidate_count = ssh.credential_candidates.len(),
|
|
"fetched Repository SSH access resource from Workspace Backend"
|
|
);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn authorize_working_directory_repository_access_from_resource(
|
|
&self,
|
|
mut request: WorkingDirectoryRepositoryAccessRequest,
|
|
) -> Result<(), RuntimeError> {
|
|
let materialization_runtime_id = request.materialization.runtime_id.clone();
|
|
let ssh = request.materialization.ssh.as_mut().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest("Repository SSH access metadata is missing".to_string())
|
|
})?;
|
|
if let Err(error) = self.resolve_repository_access_resource(ssh).await {
|
|
tracing::warn!(
|
|
target: "yoi::repository_access",
|
|
event = "repository_access_resource_fetch_failed",
|
|
workspace_id = %ssh.secret_resource.workspace_id,
|
|
resource_id = %ssh.secret_resource.resource_id,
|
|
runtime_id = %materialization_runtime_id,
|
|
error = %error,
|
|
"failed to fetch Repository SSH access resource from Workspace Backend"
|
|
);
|
|
return Err(error);
|
|
}
|
|
self.authorize_working_directory_repository_access(request)
|
|
}
|
|
|
|
/// List Runtime-owned working directories through the attached execution backend.
|
|
pub fn list_working_directories(
|
|
&self,
|
|
) -> Result<Vec<CatalogWorkingDirectoryStatus>, RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "working directory listing requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
let statuses = backend.list_working_directories();
|
|
self.annotate_working_directory_statuses(statuses)
|
|
}
|
|
|
|
/// Get a Runtime-owned working directory status.
|
|
pub fn working_directory(
|
|
&self,
|
|
working_directory_id: &str,
|
|
) -> Result<CatalogWorkingDirectoryStatus, RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "working directory lookup requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
let status = backend
|
|
.working_directory(working_directory_id)
|
|
.map_err(RuntimeError::from)?;
|
|
self.annotate_working_directory_status(status)
|
|
}
|
|
|
|
/// Open a fresh Workdir operation session in this Runtime's authorized Workspace.
|
|
///
|
|
/// A same-Runtime owner Worker can be supplied as an additional persisted-binding check.
|
|
/// Cross-Runtime callers rely on the Runtime's Workspace capability scope and the existence
|
|
/// of the Runtime-owned materialization; Backend attachment remains occupancy authority.
|
|
pub fn open_workdir_session_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
working_directory_id: &str,
|
|
owner_worker_ref: Option<&WorkerRef>,
|
|
) -> Result<workdir::WorkdirSessionHandle, RuntimeError> {
|
|
let backend = {
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_workspace_owner(scope, false)?;
|
|
|
|
if let Some(worker_ref) = owner_worker_ref {
|
|
let owns_workdir = state
|
|
.workers
|
|
.get(&worker_ref.worker_id)
|
|
.is_some_and(|worker| {
|
|
worker.belongs_to_workspace(&scope.workspace_id)
|
|
&& worker.working_directory.as_ref().is_some_and(|status| {
|
|
status.summary.working_directory_id == working_directory_id
|
|
})
|
|
});
|
|
if !owns_workdir {
|
|
return Err(RuntimeError::WorkingDirectory(
|
|
crate::working_directory::WorkingDirectoryDiagnostic::rejected(
|
|
"working_directory_not_found",
|
|
"working directory was not assigned to the authorized owner Worker",
|
|
),
|
|
));
|
|
}
|
|
}
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "opening a Workdir session requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.working_directory(working_directory_id)
|
|
.map_err(RuntimeError::WorkingDirectory)?;
|
|
backend
|
|
.open_workdir_session(working_directory_id)
|
|
.map_err(RuntimeError::WorkingDirectory)
|
|
}
|
|
|
|
/// Cleanup a Runtime-owned working directory.
|
|
pub fn cleanup_working_directory(
|
|
&self,
|
|
working_directory_id: &str,
|
|
) -> Result<CatalogWorkingDirectoryStatus, RuntimeError> {
|
|
let backend = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
if let Some(worker_id) = state.primary_worker_id_for_workdir(working_directory_id) {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"working directory {working_directory_id} is assigned to worker {worker_id}"
|
|
)));
|
|
}
|
|
state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "working directory cleanup requires an execution backend".to_string(),
|
|
}
|
|
})?
|
|
};
|
|
backend
|
|
.cleanup_working_directory(working_directory_id)
|
|
.map_err(RuntimeError::from)
|
|
}
|
|
|
|
fn annotate_working_directory_statuses(
|
|
&self,
|
|
statuses: Vec<CatalogWorkingDirectoryStatus>,
|
|
) -> Result<Vec<CatalogWorkingDirectoryStatus>, RuntimeError> {
|
|
statuses
|
|
.into_iter()
|
|
.map(|status| self.annotate_working_directory_status(status))
|
|
.collect()
|
|
}
|
|
|
|
fn annotate_working_directory_status(
|
|
&self,
|
|
mut status: CatalogWorkingDirectoryStatus,
|
|
) -> Result<CatalogWorkingDirectoryStatus, RuntimeError> {
|
|
let state = self.lock()?;
|
|
status.summary.primary_worker_id = state
|
|
.primary_worker_id_for_workdir(status.summary.working_directory_id.as_str())
|
|
.map(|worker_id| worker_id.to_string());
|
|
Ok(status)
|
|
}
|
|
|
|
/// Create a Worker through the canonical profile-source + execution backend path.
|
|
pub fn create_worker(
|
|
&self,
|
|
request: CreateWorkerRequest,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
self.create_worker_with_workspace(request, None)
|
|
}
|
|
|
|
/// Create a Worker scoped to a workspace authorization context.
|
|
pub fn create_worker_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
request: CreateWorkerRequest,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
self.create_worker_with_workspace(request, Some(scope))
|
|
}
|
|
|
|
fn existing_worker_for_create(
|
|
&self,
|
|
request: &CreateWorkerRequest,
|
|
scope: Option<&RuntimeWorkspaceScope>,
|
|
) -> Result<Option<WorkerDetail>, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
validate_create_worker_request(request)?;
|
|
validate_create_workspace_scope(request, scope.map(|scope| scope.workspace_id.as_str()))?;
|
|
if let Some(scope) = scope {
|
|
state.ensure_workspace_owner(scope, false)?;
|
|
}
|
|
let workspace_id = scope.map(|scope| scope.workspace_id.as_str());
|
|
if let Some(existing) = state.workers.get(&request.worker_id) {
|
|
if existing.workspace_id.as_deref() != workspace_id {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} already belongs to another Workspace scope",
|
|
request.worker_id
|
|
)));
|
|
}
|
|
if existing.request.create_fingerprint != request.create_fingerprint {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} was already created with a different fingerprint",
|
|
request.worker_id
|
|
)));
|
|
}
|
|
return Ok(Some(existing.detail()));
|
|
}
|
|
Ok(None)
|
|
}
|
|
|
|
fn refresh_workspace_config(&self, request: &CreateWorkerRequest) -> Result<(), RuntimeError> {
|
|
let ProfileSourceArchiveSource::WorkspaceConfig {
|
|
archive: expected_archive,
|
|
} = &request.profile_source
|
|
else {
|
|
return Ok(());
|
|
};
|
|
let expected = request.config_bundle.as_ref().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Workspace Config profile source requires a config_bundle reference".to_string(),
|
|
)
|
|
})?;
|
|
let workspace_api = request.workspace_api.clone().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Workspace Config profile source requires a Workspace API reference".to_string(),
|
|
)
|
|
})?;
|
|
let profile_key = match &request.profile {
|
|
ProfileSelector::Builtin(value) | ProfileSelector::Named(value) => value.clone(),
|
|
};
|
|
let cache_key = format!(
|
|
"{}\u{0}{}\u{0}{}",
|
|
workspace_api.base_url, workspace_api.workspace_id, profile_key
|
|
);
|
|
let fetch_gate = {
|
|
let mut state = self.lock()?;
|
|
state
|
|
.workspace_config_fetch_gates
|
|
.entry(cache_key.clone())
|
|
.or_insert_with(|| Arc::new(Mutex::new(())))
|
|
.clone()
|
|
};
|
|
let _fetch_guard = fetch_gate.lock().map_err(|_| RuntimeError::StatePoisoned)?;
|
|
let (backend, cached) = {
|
|
let state = self.lock()?;
|
|
let cached_reference = state
|
|
.workspace_config_latest
|
|
.get(&cache_key)
|
|
.unwrap_or(expected);
|
|
let cached = state
|
|
.config_bundles
|
|
.get(&cached_reference.id)
|
|
.filter(|bundle| bundle.metadata.digest == cached_reference.digest)
|
|
.map(|bundle| ConfigBundleRef {
|
|
id: bundle.metadata.id.clone(),
|
|
digest: bundle.metadata.digest.clone(),
|
|
});
|
|
(state.execution_backend.clone(), cached)
|
|
};
|
|
let backend = backend.ok_or_else(|| RuntimeError::ExecutionBackendUnavailable {
|
|
message: "Workspace Config refresh requires an execution backend".to_string(),
|
|
})?;
|
|
let fetched = backend
|
|
.fetch_workspace_config(WorkspaceConfigFetchRequest {
|
|
workspace_api,
|
|
profile: request.profile.clone(),
|
|
expected: expected.clone(),
|
|
cached: cached.clone(),
|
|
})
|
|
.map_err(|message| {
|
|
RuntimeError::InvalidRequest(format!(
|
|
"failed to refresh latest Workspace Config: {message}"
|
|
))
|
|
})?;
|
|
let resolved = match fetched {
|
|
WorkspaceConfigFetchResult::NotModified => cached.ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"latest Workspace Config returned not-modified without a cached bundle"
|
|
.to_string(),
|
|
)
|
|
})?,
|
|
WorkspaceConfigFetchResult::Modified(bundle) => {
|
|
self.store_config_bundle(bundle)?.reference
|
|
}
|
|
};
|
|
if &resolved != expected {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"latest Workspace Config changed while Worker creation was being prepared: expected '{}@{}', got '{}@{}'",
|
|
expected.id, expected.digest, resolved.id, resolved.digest
|
|
)));
|
|
}
|
|
let mut state = self.lock()?;
|
|
let bundle = state.config_bundles.get(&resolved.id).ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"latest Workspace Config was not available after refresh".to_string(),
|
|
)
|
|
})?;
|
|
let archive = bundle.profile_source_archive.as_ref().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"latest Workspace Config is missing its profile source archive".to_string(),
|
|
)
|
|
})?;
|
|
if &archive.reference != expected_archive {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"latest Workspace Config profile source archive reference mismatch".to_string(),
|
|
));
|
|
}
|
|
state.workspace_config_latest.insert(cache_key, resolved);
|
|
Ok(())
|
|
}
|
|
|
|
fn create_worker_with_workspace(
|
|
&self,
|
|
request: CreateWorkerRequest,
|
|
scope: Option<&RuntimeWorkspaceScope>,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
let worker_id = request.worker_id;
|
|
let workspace_id = scope.map(|scope| scope.workspace_id.as_str());
|
|
let started_at = std::time::Instant::now();
|
|
let result = self.create_worker_with_workspace_inner(request, scope);
|
|
if let Err(error) = &result {
|
|
write_runtime_worker_create_failure(
|
|
worker_id,
|
|
workspace_id,
|
|
error,
|
|
started_at.elapsed(),
|
|
);
|
|
}
|
|
result
|
|
}
|
|
|
|
fn create_worker_with_workspace_inner(
|
|
&self,
|
|
request: CreateWorkerRequest,
|
|
scope: Option<&RuntimeWorkspaceScope>,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
let operation_lock = self.worker_operation_lock(request.worker_id)?;
|
|
let _operation_guard = operation_lock
|
|
.lock()
|
|
.map_err(|_| RuntimeError::StatePoisoned)?;
|
|
if let Some(existing) = self.existing_worker_for_create(&request, scope)? {
|
|
return Ok(existing);
|
|
}
|
|
self.refresh_workspace_config(&request)?;
|
|
let (backend, worker_ref, spawn_request) = {
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
validate_create_worker_request(&request)?;
|
|
validate_create_workspace_scope(
|
|
&request,
|
|
scope.map(|scope| scope.workspace_id.as_str()),
|
|
)?;
|
|
if let Some(scope) = scope {
|
|
state.ensure_workspace_owner(scope, true)?;
|
|
};
|
|
let workspace_id = scope.map(|scope| scope.workspace_id.as_str());
|
|
if let Some(existing) = state.workers.get(&request.worker_id) {
|
|
if existing.workspace_id.as_deref() != workspace_id {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} already belongs to another Workspace scope",
|
|
request.worker_id
|
|
)));
|
|
}
|
|
if existing.request.create_fingerprint != request.create_fingerprint {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} was already created with a different fingerprint",
|
|
request.worker_id
|
|
)));
|
|
}
|
|
return Ok(existing.detail());
|
|
}
|
|
state.validate_worker_config_boundary(&request)?;
|
|
if let Some(working_directory_id) = requested_primary_workdir_id(&request) {
|
|
if let Some(owner_worker_id) =
|
|
state.primary_worker_id_for_workdir(working_directory_id)
|
|
{
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"working directory {working_directory_id} is already assigned to worker {owner_worker_id}"
|
|
)));
|
|
}
|
|
}
|
|
let backend = state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::ExecutionBackendUnavailable {
|
|
message: "worker creation requires an execution backend".to_string(),
|
|
}
|
|
})?;
|
|
let config_bundle = state.resolve_config_bundle_ref(request.config_bundle.as_ref())?;
|
|
|
|
let worker_id = request.worker_id;
|
|
let worker_ref = WorkerRef::new(worker_id);
|
|
|
|
let durable_request = durable_create_worker_request(&request);
|
|
let record = WorkerRecord {
|
|
worker_ref: worker_ref.clone(),
|
|
worker_id: worker_id.clone(),
|
|
status: WorkerStatus::Stopped,
|
|
worker_state: None,
|
|
workspace_id: scope.map(|scope| scope.workspace_id.clone()),
|
|
request: durable_request,
|
|
run_generation: 1,
|
|
execution_bound: true,
|
|
restore_intent: WorkerRestoreIntent::Explicit,
|
|
working_directory: None,
|
|
execution_handle: None,
|
|
internal_workers: BTreeMap::new(),
|
|
};
|
|
state.workers.insert(worker_id, record);
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
let spawn_request = WorkerExecutionSpawnRequest {
|
|
worker_ref: worker_ref.clone(),
|
|
run_generation: 1,
|
|
request,
|
|
workspace_scope: scope.cloned(),
|
|
context: self.execution_context(worker_ref.clone()),
|
|
working_directory: None,
|
|
config_bundle,
|
|
};
|
|
(backend, worker_ref, spawn_request)
|
|
};
|
|
|
|
let spawn_result = backend.spawn_worker(spawn_request);
|
|
let (handle, initial_worker_state, working_directory) = match spawn_result {
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle,
|
|
worker_state,
|
|
working_directory,
|
|
} => (handle, worker_state, working_directory),
|
|
WorkerExecutionSpawnResult::Rejected(result)
|
|
| WorkerExecutionSpawnResult::Errored(result) => {
|
|
self.rollback_failed_create(&worker_ref)?;
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
});
|
|
}
|
|
};
|
|
|
|
if let Some(mut initial_input) = {
|
|
let state = self.lock()?;
|
|
state.worker(&worker_ref)?.request.initial_input.clone()
|
|
} {
|
|
let expected_submission_id = initial_input
|
|
.submission_request_id
|
|
.clone()
|
|
.filter(|request_id| !request_id.trim().is_empty())
|
|
.unwrap_or_else(|| Uuid::now_v7().to_string());
|
|
initial_input.submission_request_id = Some(expected_submission_id.clone());
|
|
let dispatch_result = backend.dispatch_input(&handle, initial_input.clone());
|
|
if !dispatch_result.is_accepted() {
|
|
self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?;
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: dispatch_result.operation,
|
|
outcome: dispatch_result.outcome,
|
|
message: dispatch_result.message_or_default(),
|
|
result: dispatch_result,
|
|
});
|
|
}
|
|
let has_durable_acceptance = dispatch_result
|
|
.submission
|
|
.as_ref()
|
|
.is_some_and(|ack| ack.submission_request_id == expected_submission_id);
|
|
if !has_durable_acceptance {
|
|
self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?;
|
|
let result = WorkerExecutionResult::rejected(
|
|
WorkerExecutionOperation::Input,
|
|
"execution backend accepted initial input without a durable submission acknowledgement",
|
|
);
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
});
|
|
}
|
|
let detail = match self.commit_created_worker(
|
|
&worker_ref,
|
|
handle.clone(),
|
|
initial_worker_state.clone(),
|
|
working_directory,
|
|
dispatch_result,
|
|
) {
|
|
Ok(detail) => detail,
|
|
Err(error) => {
|
|
self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?;
|
|
return Err(error);
|
|
}
|
|
};
|
|
if let Err(error) = self.record_input_observation(&worker_ref, initial_input) {
|
|
self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?;
|
|
return Err(error);
|
|
}
|
|
Ok(detail)
|
|
} else {
|
|
match self.commit_created_worker(
|
|
&worker_ref,
|
|
handle.clone(),
|
|
initial_worker_state,
|
|
working_directory,
|
|
WorkerExecutionResult::accepted(WorkerExecutionOperation::Spawn),
|
|
) {
|
|
Ok(detail) => Ok(detail),
|
|
Err(error) => {
|
|
self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?;
|
|
Err(error)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// List Workers known to this Runtime.
|
|
pub fn list_workers(&self) -> Result<Vec<WorkerSummary>, RuntimeError> {
|
|
let state = self.lock()?;
|
|
Ok(state.workers.values().map(WorkerRecord::summary).collect())
|
|
}
|
|
|
|
pub fn subscribe_event_selector(
|
|
&self,
|
|
selector: EventSubscriptionSelector,
|
|
) -> Result<RuntimeEventSelectorSubscription, RuntimeError> {
|
|
self.subscribe_event_selector_for_workspace(None, selector)
|
|
}
|
|
|
|
pub fn subscribe_event_selector_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
selector: EventSubscriptionSelector,
|
|
) -> Result<RuntimeEventSelectorSubscription, RuntimeError> {
|
|
self.subscribe_event_selector_for_workspace(Some(scope), selector)
|
|
}
|
|
|
|
fn subscribe_event_selector_for_workspace(
|
|
&self,
|
|
scope: Option<&RuntimeWorkspaceScope>,
|
|
selector: EventSubscriptionSelector,
|
|
) -> Result<RuntimeEventSelectorSubscription, RuntimeError> {
|
|
selector.validate().map_err(subscription_validation_error)?;
|
|
if matches!(
|
|
selector,
|
|
EventSubscriptionSelector::WorkerProtocol { .. }
|
|
| EventSubscriptionSelector::WorkspaceWorkers
|
|
| EventSubscriptionSelector::WorkspaceWorkdirs
|
|
) {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Runtime event subscriptions support only runtime_workers and worker_lifecycle selectors"
|
|
.to_string(),
|
|
));
|
|
}
|
|
|
|
let mut state = self.lock()?;
|
|
if let Some(scope) = scope {
|
|
state.ensure_workspace_owner_for_existing_workers(scope)?;
|
|
state.persist_runtime_snapshot()?;
|
|
}
|
|
let snapshot = state.subscription_snapshot(scope, &selector)?;
|
|
let snapshot_revision = state.subscription_revision;
|
|
let subscription_id = state.next_event_subscription_id;
|
|
state.next_event_subscription_id = state.next_event_subscription_id.saturating_add(1);
|
|
let (sender, receiver) = mpsc::channel(SUBSCRIPTION_QUEUE_CAPACITY);
|
|
let lagged = Arc::new(AtomicBool::new(false));
|
|
state.subscriptions.insert(
|
|
subscription_id,
|
|
SubscriptionSink {
|
|
selector: selector.clone(),
|
|
workspace_id: scope.map(|scope| scope.workspace_id.clone()),
|
|
sender,
|
|
lagged: lagged.clone(),
|
|
},
|
|
);
|
|
|
|
Ok(RuntimeEventSelectorSubscription {
|
|
subscription_id,
|
|
selector,
|
|
snapshot_revision,
|
|
snapshot,
|
|
receiver,
|
|
lagged,
|
|
runtime: Arc::downgrade(&self.inner),
|
|
})
|
|
}
|
|
|
|
/// List stopped Workers known to this Runtime.
|
|
pub fn list_stopped_workers(&self) -> Result<Vec<WorkerSummary>, RuntimeError> {
|
|
let state = self.lock()?;
|
|
Ok(state
|
|
.workers
|
|
.values()
|
|
.filter(|worker| worker.status == WorkerStatus::Stopped)
|
|
.map(WorkerRecord::summary)
|
|
.collect())
|
|
}
|
|
|
|
/// List Workers visible to a workspace-scoped Runtime authorization context.
|
|
pub fn list_workers_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
) -> Result<Vec<WorkerSummary>, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let visible_workspace = state.ensure_workspace_owner_for_existing_workers(scope)?;
|
|
if !visible_workspace {
|
|
return Ok(Vec::new());
|
|
}
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(state
|
|
.workers
|
|
.values()
|
|
.filter(|worker| worker.belongs_to_workspace(&scope.workspace_id))
|
|
.map(WorkerRecord::summary)
|
|
.collect())
|
|
}
|
|
|
|
/// List stopped Workers visible to a workspace-scoped Runtime authorization context.
|
|
pub fn list_stopped_workers_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
) -> Result<Vec<WorkerSummary>, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let visible_workspace = state.ensure_workspace_owner_for_existing_workers(scope)?;
|
|
if !visible_workspace {
|
|
return Ok(Vec::new());
|
|
}
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(state
|
|
.workers
|
|
.values()
|
|
.filter(|worker| {
|
|
worker.status == WorkerStatus::Stopped
|
|
&& worker.belongs_to_workspace(&scope.workspace_id)
|
|
})
|
|
.map(WorkerRecord::summary)
|
|
.collect())
|
|
}
|
|
|
|
/// Fetch Worker detail through a workspace-scoped Runtime authorization context.
|
|
pub fn worker_detail_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if !worker.belongs_to_workspace(&scope.workspace_id) {
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
});
|
|
}
|
|
state.ensure_workspace_owner(scope, true)?;
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(state.worker(worker_ref)?.detail())
|
|
}
|
|
|
|
/// Fetch Worker detail. The supplied [`WorkerRef`] must match this Runtime.
|
|
pub fn worker_detail(&self, worker_ref: &WorkerRef) -> Result<WorkerDetail, RuntimeError> {
|
|
let state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
Ok(worker.detail())
|
|
}
|
|
|
|
fn ensure_worker_in_workspace(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<(), RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if !worker.belongs_to_workspace(&scope.workspace_id) {
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
});
|
|
}
|
|
state.ensure_workspace_owner(scope, true)?;
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Replace the Workspace API identity binding persisted for a Worker.
|
|
pub fn replace_worker_workspace_api_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
workspace_api: WorkspaceApiRef,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
if workspace_api.workspace_id != scope.workspace_id {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"Workspace API scope `{}` does not match authorized workspace `{}`",
|
|
workspace_api.workspace_id, scope.workspace_id
|
|
)));
|
|
}
|
|
self.replace_worker_workspace_api(worker_ref, workspace_api)
|
|
}
|
|
|
|
pub fn replace_worker_workspace_api(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
workspace_api: WorkspaceApiRef,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
let previous_workspace_api = {
|
|
let state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if let Some(existing) = worker.request.workspace_api.as_ref()
|
|
&& (existing.workspace_id != workspace_api.workspace_id
|
|
|| existing.base_url.trim_end_matches('/')
|
|
!= workspace_api.base_url.trim_end_matches('/'))
|
|
{
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Workspace API replacement cannot change Worker Workspace identity or base URL"
|
|
.to_string(),
|
|
));
|
|
}
|
|
worker.request.workspace_api.clone()
|
|
};
|
|
|
|
{
|
|
let mut state = self.lock()?;
|
|
state.worker_mut(worker_ref)?.request.workspace_api = Some(workspace_api);
|
|
if let Err(error) = state.persist_runtime_snapshot() {
|
|
state.worker_mut(worker_ref)?.request.workspace_api = previous_workspace_api;
|
|
return Err(error);
|
|
}
|
|
}
|
|
|
|
let state = self.lock()?;
|
|
Ok(state.worker(worker_ref)?.detail())
|
|
}
|
|
|
|
/// Attach a live execution through a workspace-scoped Runtime authorization context.
|
|
pub fn restore_worker_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.restore_worker(worker_ref)
|
|
}
|
|
|
|
/// Attach a live execution to a persisted Worker definition.
|
|
///
|
|
/// Current liveness is never read from disk. If a handle is already
|
|
/// present this is idempotent; otherwise the configured backend is tried.
|
|
pub fn restore_worker(&self, worker_ref: &WorkerRef) -> Result<WorkerDetail, RuntimeError> {
|
|
let (backend, request) = {
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
let (worker_request, previous_working_directory, run_generation) = {
|
|
let worker = state.worker(worker_ref)?;
|
|
if worker.execution_handle.is_some() {
|
|
return Ok(worker.detail());
|
|
}
|
|
(
|
|
worker.request.clone(),
|
|
worker.working_directory.clone(),
|
|
worker.run_generation.saturating_add(1).max(1),
|
|
)
|
|
};
|
|
let backend = state.execution_backend.clone().ok_or_else(|| {
|
|
RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
message: "runtime has no execution backend".to_string(),
|
|
}
|
|
})?;
|
|
{
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
worker.run_generation = run_generation;
|
|
worker.execution_bound = true;
|
|
worker.restore_intent = WorkerRestoreIntent::Automatic;
|
|
}
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
let workspace_scope = worker_request.workspace_api.as_ref().and_then(|api| {
|
|
state
|
|
.workspace_owners
|
|
.get(&api.workspace_id)
|
|
.map(|server_id| RuntimeWorkspaceScope::new(&api.workspace_id, server_id))
|
|
});
|
|
let request = WorkerExecutionRestoreRequest {
|
|
worker_ref: worker_ref.clone(),
|
|
run_generation,
|
|
request: worker_request,
|
|
workspace_scope,
|
|
context: self.execution_context(worker_ref.clone()),
|
|
previous_working_directory,
|
|
working_directory: None,
|
|
config_bundle: None,
|
|
};
|
|
(backend, request)
|
|
};
|
|
|
|
match backend.restore_worker(request) {
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle,
|
|
worker_state,
|
|
working_directory,
|
|
} => {
|
|
self.commit_restored_worker_execution(
|
|
worker_ref,
|
|
handle,
|
|
worker_state,
|
|
WorkerStatus::Idle,
|
|
working_directory,
|
|
)?;
|
|
self.worker_detail(worker_ref)
|
|
}
|
|
WorkerExecutionSpawnResult::Rejected(result)
|
|
| WorkerExecutionSpawnResult::Errored(result) => {
|
|
#[cfg(feature = "fs-store")]
|
|
{
|
|
self.lock()?
|
|
.record_restore_failure(worker_ref, result.clone())?;
|
|
}
|
|
Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
})
|
|
}
|
|
}
|
|
}
|
|
|
|
fn ensure_worker_execution(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> {
|
|
let state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if worker.execution_handle.is_some() {
|
|
return Ok(());
|
|
}
|
|
let message = if worker.status == WorkerStatus::Stopped {
|
|
"stopped worker requires an explicit restore"
|
|
} else {
|
|
"worker has no live execution handle"
|
|
};
|
|
Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id,
|
|
message: message.to_string(),
|
|
})
|
|
}
|
|
|
|
/// Accept input into a Worker through a workspace-scoped Runtime authorization context.
|
|
pub fn send_input_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
input: WorkerInput,
|
|
) -> Result<WorkerInteractionAck, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.send_input(worker_ref, input)
|
|
}
|
|
|
|
/// Accept input into a Worker.
|
|
pub fn send_input(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
mut input: WorkerInput,
|
|
) -> Result<WorkerInteractionAck, RuntimeError> {
|
|
validate_worker_input(&input)?;
|
|
let expected_submission_id =
|
|
if matches!(input.kind, WorkerInputKind::User | WorkerInputKind::Notify) {
|
|
let submission_id = input
|
|
.submission_request_id
|
|
.clone()
|
|
.filter(|request_id| !request_id.trim().is_empty())
|
|
.unwrap_or_else(|| Uuid::now_v7().to_string());
|
|
input.submission_request_id = Some(submission_id.clone());
|
|
Some(submission_id)
|
|
} else {
|
|
None
|
|
};
|
|
self.ensure_worker_execution(worker_ref)?;
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if !worker.status.is_active() {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} is not running",
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
let backend = state.execution_backend.clone();
|
|
let handle = worker.execution_handle.clone();
|
|
match (backend, handle) {
|
|
(Some(backend), Some(handle)) => (backend, handle),
|
|
_ => {
|
|
return Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
message: "worker has no live execution handle".to_string(),
|
|
});
|
|
}
|
|
}
|
|
};
|
|
|
|
let dispatch_result = backend.dispatch_input(&handle, input.clone());
|
|
if !dispatch_result.is_accepted() {
|
|
self.record_execution_result(worker_ref, dispatch_result.clone())?;
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: dispatch_result.operation,
|
|
outcome: dispatch_result.outcome,
|
|
message: dispatch_result.message_or_default(),
|
|
result: dispatch_result,
|
|
});
|
|
}
|
|
if let Some(expected_submission_id) = expected_submission_id
|
|
&& dispatch_result
|
|
.submission
|
|
.as_ref()
|
|
.is_none_or(|ack| ack.submission_request_id != expected_submission_id)
|
|
{
|
|
let result = WorkerExecutionResult::rejected(
|
|
WorkerExecutionOperation::Input,
|
|
"execution backend did not acknowledge the committed Runtime submission request id",
|
|
);
|
|
self.record_execution_result(worker_ref, result.clone())?;
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
});
|
|
}
|
|
|
|
let submission = dispatch_result.submission.clone();
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
if let Some(snapshot) = dispatch_result.worker_state.as_ref() {
|
|
let _ = worker.apply_worker_state(snapshot);
|
|
}
|
|
let status = worker.status;
|
|
#[cfg(feature = "ws-server")]
|
|
if let Some(payload) = input_protocol_event(&input) {
|
|
state.push_worker_observation_event(worker_ref.clone(), payload);
|
|
}
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
|
|
Ok(WorkerInteractionAck {
|
|
worker_ref: worker_ref.clone(),
|
|
status,
|
|
submission,
|
|
})
|
|
}
|
|
|
|
/// Store a client-local file in the owning Worker session before input submit.
|
|
pub fn upload_worker_file_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
file_name: &str,
|
|
media_type: &str,
|
|
content: &[u8],
|
|
) -> Result<protocol::UploadedFileRef, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.upload_worker_file(worker_ref, file_name, media_type, content)
|
|
}
|
|
|
|
pub fn upload_worker_file(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
file_name: &str,
|
|
media_type: &str,
|
|
content: &[u8],
|
|
) -> Result<protocol::UploadedFileRef, RuntimeError> {
|
|
self.upload_worker_file_inner(worker_ref, file_name, media_type, content, None)
|
|
}
|
|
|
|
pub fn upload_worker_file_with_context_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
file_name: &str,
|
|
media_type: &str,
|
|
content: &[u8],
|
|
context: &session_store::UploadedFileUploadContext,
|
|
) -> Result<protocol::UploadedFileRef, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.upload_worker_file_inner(worker_ref, file_name, media_type, content, Some(context))
|
|
}
|
|
|
|
pub fn upload_worker_file_with_context(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
file_name: &str,
|
|
media_type: &str,
|
|
content: &[u8],
|
|
context: &session_store::UploadedFileUploadContext,
|
|
) -> Result<protocol::UploadedFileRef, RuntimeError> {
|
|
self.upload_worker_file_inner(worker_ref, file_name, media_type, content, Some(context))
|
|
}
|
|
|
|
fn upload_worker_file_inner(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
file_name: &str,
|
|
media_type: &str,
|
|
content: &[u8],
|
|
context: Option<&session_store::UploadedFileUploadContext>,
|
|
) -> Result<protocol::UploadedFileRef, RuntimeError> {
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
match (
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
) {
|
|
(Some(backend), Some(handle)) => (backend, handle),
|
|
_ => {
|
|
return Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
message: "worker has no live execution handle".to_string(),
|
|
});
|
|
}
|
|
}
|
|
};
|
|
backend
|
|
.upload_file(&handle, file_name, media_type, content, context)
|
|
.map_err(|result| RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
})
|
|
}
|
|
|
|
/// Delete an unsubmitted uploaded file from the owning Worker session.
|
|
pub fn delete_worker_uploaded_file_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
artifact_id: &str,
|
|
) -> Result<(), RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.delete_worker_uploaded_file(worker_ref, artifact_id)
|
|
}
|
|
|
|
pub fn delete_worker_uploaded_file(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
artifact_id: &str,
|
|
) -> Result<(), RuntimeError> {
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
match (
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
) {
|
|
(Some(backend), Some(handle)) => (backend, handle),
|
|
_ => {
|
|
return Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
message: "worker has no live execution handle".to_string(),
|
|
});
|
|
}
|
|
}
|
|
};
|
|
let result = backend.delete_uploaded_file(&handle, artifact_id);
|
|
if result.is_accepted() {
|
|
Ok(())
|
|
} else {
|
|
Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
})
|
|
}
|
|
}
|
|
|
|
/// Return live completion entries through a workspace-scoped Runtime authorization context.
|
|
pub fn worker_completions_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
kind: protocol::CompletionKind,
|
|
prefix: &str,
|
|
) -> Result<Vec<protocol::CompletionEntry>, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.worker_completions(worker_ref, kind, prefix)
|
|
}
|
|
|
|
/// Return live completion entries for the Worker composer.
|
|
pub fn worker_completions(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
kind: protocol::CompletionKind,
|
|
prefix: &str,
|
|
) -> Result<Vec<protocol::CompletionEntry>, RuntimeError> {
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
(
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
)
|
|
};
|
|
let Some((backend, handle)) = backend.zip(handle) else {
|
|
return Ok(Vec::new());
|
|
};
|
|
Ok(backend.worker_completions(&handle, kind, prefix))
|
|
}
|
|
|
|
/// Accept a protocol method through a workspace-scoped Runtime authorization context.
|
|
pub fn send_protocol_method_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
method: Method,
|
|
) -> Result<Vec<Event>, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.send_protocol_method(worker_ref, method)
|
|
}
|
|
|
|
/// Accept a protocol method for a Worker through a Backend/runtime transport.
|
|
///
|
|
/// Most methods are delivered to the execution backend unchanged. Methods with
|
|
/// direct same-connection replies in the local socket protocol return those
|
|
/// events from this function so WebSocket transports can write them back to the
|
|
/// requesting client without rebroadcasting them.
|
|
pub fn send_protocol_method(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
method: Method,
|
|
) -> Result<Vec<Event>, RuntimeError> {
|
|
if let Method::ListCompletions { kind, prefix } = method {
|
|
let entries = self.worker_completions(worker_ref, kind, &prefix)?;
|
|
return Ok(vec![Event::Completions { kind, entries }]);
|
|
}
|
|
if matches!(&method, Method::Shutdown { .. }) {
|
|
self.stop_worker(worker_ref, Some("worker protocol shutdown".to_string()))?;
|
|
return Ok(Vec::new());
|
|
}
|
|
self.ensure_worker_execution(worker_ref)?;
|
|
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if !worker.status.is_active() {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} is not running",
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
let backend = state.execution_backend.clone();
|
|
let handle = worker.execution_handle.clone();
|
|
match (backend, handle) {
|
|
(Some(backend), Some(handle)) => (backend, handle),
|
|
_ => {
|
|
return Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
message: "worker has no live execution handle".to_string(),
|
|
});
|
|
}
|
|
}
|
|
};
|
|
|
|
let dispatch_result = backend.dispatch_method(&handle, method);
|
|
if !dispatch_result.is_accepted() {
|
|
self.record_execution_result(worker_ref, dispatch_result.clone())?;
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: dispatch_result.operation,
|
|
outcome: dispatch_result.outcome,
|
|
message: dispatch_result.message_or_default(),
|
|
result: dispatch_result,
|
|
});
|
|
}
|
|
|
|
self.record_execution_result(worker_ref, dispatch_result)?;
|
|
Ok(Vec::new())
|
|
}
|
|
|
|
fn commit_created_worker(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
handle: WorkerExecutionHandle,
|
|
initial_worker_state: protocol::WorkerStateSnapshot,
|
|
working_directory: Option<CatalogWorkingDirectoryStatus>,
|
|
result: WorkerExecutionResult,
|
|
) -> Result<WorkerDetail, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let detail = {
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
worker.execution_handle = Some(handle);
|
|
worker.execution_bound = true;
|
|
worker.status = WorkerStatus::Idle;
|
|
let _ = worker.apply_worker_state(&initial_worker_state);
|
|
if let Some(snapshot) = result.worker_state.as_ref() {
|
|
let _ = worker.apply_worker_state(snapshot);
|
|
}
|
|
worker.restore_intent = restore_intent_for_status(worker.status);
|
|
worker.working_directory = working_directory;
|
|
worker.detail()
|
|
};
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(detail)
|
|
}
|
|
|
|
fn cleanup_connected_failed_create(
|
|
&self,
|
|
backend: &WorkerExecutionBackendRef,
|
|
worker_ref: &WorkerRef,
|
|
handle: &WorkerExecutionHandle,
|
|
) -> Result<(), RuntimeError> {
|
|
let stop_result = backend.stop_worker(handle);
|
|
if stop_result.is_accepted() {
|
|
return self.rollback_failed_create(worker_ref);
|
|
}
|
|
let mut state = self.lock()?;
|
|
let record = state.worker_mut(worker_ref)?;
|
|
record.execution_handle = Some(handle.clone());
|
|
record.worker_state = stop_result.worker_state.clone();
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(())
|
|
}
|
|
|
|
fn rollback_failed_create(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
if state.workers.contains_key(&worker_ref.worker_id) {
|
|
state.delete_worker_snapshot(&worker_ref.worker_id)?;
|
|
let record = state
|
|
.workers
|
|
.remove(&worker_ref.worker_id)
|
|
.expect("Worker existence checked before failed-create rollback");
|
|
let workspace_id = record.workspace_id.clone();
|
|
if let Some(workspace_id) = workspace_id.as_deref() {
|
|
state.forget_workspace_owner_if_unused(workspace_id);
|
|
}
|
|
#[cfg(feature = "ws-server")]
|
|
state
|
|
.observation_events
|
|
.retain(|event| event.worker_ref != *worker_ref);
|
|
state.publish_worker_removed(worker_ref.worker_id, workspace_id.as_deref())?;
|
|
state.persist_runtime_snapshot()?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn record_execution_result(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
result: WorkerExecutionResult,
|
|
) -> Result<(), RuntimeError> {
|
|
// Accepted dispatch without a state snapshot is transport evidence only;
|
|
// the revisioned protocol stream remains live authority. Test/detached
|
|
// backends may return an exact full snapshot as their acknowledgement.
|
|
if !result.is_accepted() {
|
|
return Ok(());
|
|
}
|
|
let Some(snapshot) = result.worker_state else {
|
|
return Ok(());
|
|
};
|
|
let mut state = self.lock()?;
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
let applied = worker
|
|
.apply_worker_state(&snapshot)
|
|
.is_ok_and(|result| matches!(result, protocol::WorkerStateSnapshotApply::Applied));
|
|
if !applied {
|
|
return Ok(());
|
|
}
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(())
|
|
}
|
|
|
|
fn dispatch_lifecycle_to_backend(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
operation: WorkerExecutionOperation,
|
|
) -> Result<(), RuntimeError> {
|
|
let Some((backend, handle)) = ({
|
|
let state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if !worker.status.is_active() {
|
|
return Ok(());
|
|
}
|
|
match (
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
) {
|
|
(Some(backend), Some(handle)) => Some((backend, handle)),
|
|
_ => None,
|
|
}
|
|
}) else {
|
|
return Ok(());
|
|
};
|
|
|
|
let result = match operation {
|
|
WorkerExecutionOperation::Stop => backend.stop_worker(&handle),
|
|
WorkerExecutionOperation::Cancel => backend.cancel_worker(&handle),
|
|
WorkerExecutionOperation::Spawn
|
|
| WorkerExecutionOperation::Restore
|
|
| WorkerExecutionOperation::Input
|
|
| WorkerExecutionOperation::UploadFile
|
|
| WorkerExecutionOperation::DeleteUploadedFile
|
|
| WorkerExecutionOperation::ProtocolMethod => return Ok(()),
|
|
};
|
|
if result.is_accepted() {
|
|
return Ok(());
|
|
}
|
|
Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id.clone(),
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
})
|
|
}
|
|
|
|
/// Stop a Worker through a workspace-scoped Runtime authorization context.
|
|
pub fn stop_worker_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
reason: Option<String>,
|
|
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.stop_worker(worker_ref, reason)
|
|
}
|
|
|
|
/// Stop a Worker. Repeated stops are idempotent.
|
|
pub fn stop_worker(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
reason: Option<String>,
|
|
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
|
self.dispatch_lifecycle_to_backend(worker_ref, WorkerExecutionOperation::Stop)?;
|
|
let _ = reason;
|
|
self.transition_worker(worker_ref, WorkerStatus::Stopped)
|
|
}
|
|
|
|
/// Cancel a Worker through a workspace-scoped Runtime authorization context.
|
|
pub fn cancel_worker_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
reason: Option<String>,
|
|
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.cancel_worker(worker_ref, reason)
|
|
}
|
|
|
|
/// Cancel the current run while keeping the Worker session available.
|
|
pub fn cancel_worker(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
reason: Option<String>,
|
|
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
|
{
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
if state.worker(worker_ref)?.status == WorkerStatus::Stopped {
|
|
return Ok(WorkerLifecycleAck {
|
|
worker_ref: worker_ref.clone(),
|
|
status: WorkerStatus::Stopped,
|
|
worker_state: None,
|
|
});
|
|
}
|
|
}
|
|
self.dispatch_lifecycle_to_backend(worker_ref, WorkerExecutionOperation::Cancel)?;
|
|
let _ = reason;
|
|
let state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
Ok(WorkerLifecycleAck {
|
|
worker_ref: worker_ref.clone(),
|
|
status: worker.status,
|
|
worker_state: worker.worker_state.clone(),
|
|
})
|
|
}
|
|
|
|
/// Delete a non-running Worker through a workspace-scoped Runtime authorization context.
|
|
pub fn delete_worker_scoped(
|
|
&self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerDeleteResult, RuntimeError> {
|
|
self.ensure_worker_in_workspace(scope, worker_ref)?;
|
|
self.delete_worker(worker_ref)
|
|
}
|
|
|
|
/// Delete a non-running Worker from Runtime state and persisted Worker storage.
|
|
pub fn delete_worker(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerDeleteResult, RuntimeError> {
|
|
let operation_lock = self.worker_operation_lock(worker_ref.worker_id)?;
|
|
let _operation_guard = operation_lock
|
|
.lock()
|
|
.map_err(|_| RuntimeError::StatePoisoned)?;
|
|
let (backend, execution_handle) = {
|
|
let state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if worker.status.is_active() {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} is running and must be stopped before deletion",
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
(
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
)
|
|
};
|
|
if let Some(handle) = execution_handle {
|
|
let backend = backend.ok_or_else(|| RuntimeError::ExecutionBackendUnavailable {
|
|
message: "Worker deletion requires its execution backend to confirm shutdown"
|
|
.to_string(),
|
|
})?;
|
|
let result = backend.stop_worker(&handle);
|
|
if !result.is_accepted() {
|
|
return Err(RuntimeError::WorkerExecutionRejected {
|
|
worker_id: worker_ref.worker_id,
|
|
operation: result.operation,
|
|
outcome: result.outcome,
|
|
message: result.message_or_default(),
|
|
result,
|
|
});
|
|
}
|
|
}
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if worker.status.is_active() {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker {} became active before deletion",
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
state.delete_worker_snapshot(&worker_ref.worker_id)?;
|
|
let removed = state.workers.remove(&worker_ref.worker_id).ok_or_else(|| {
|
|
RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
}
|
|
})?;
|
|
let removed_workspace_id = removed.workspace_id.clone();
|
|
if let Some(workspace_id) = removed_workspace_id.as_deref() {
|
|
state.forget_workspace_owner_if_unused(workspace_id);
|
|
}
|
|
#[cfg(feature = "ws-server")]
|
|
state
|
|
.observation_events
|
|
.retain(|event| event.worker_ref != *worker_ref);
|
|
state.publish_worker_removed(worker_ref.worker_id, removed_workspace_id.as_deref())?;
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(WorkerDeleteResult {
|
|
worker_id: removed.worker_id,
|
|
deleted: true,
|
|
})
|
|
}
|
|
|
|
/// Cursor pointing after the current worker-scoped protocol observation event.
|
|
#[cfg(feature = "ws-server")]
|
|
pub fn worker_observation_cursor_now(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerObservationCursor, RuntimeError> {
|
|
let state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let sequence = state
|
|
.observation_events
|
|
.iter()
|
|
.rev()
|
|
.find(|event| &event.worker_ref == worker_ref)
|
|
.map(|event| event.sequence)
|
|
.unwrap_or(0);
|
|
Ok(WorkerObservationCursor::new(sequence))
|
|
}
|
|
|
|
/// Build the current Worker Snapshot event used as the first observation frame.
|
|
#[cfg(feature = "ws-server")]
|
|
pub fn worker_observation_snapshot(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<protocol::Event, RuntimeError> {
|
|
let (backend, handle) = {
|
|
let state = self.lock()?;
|
|
let worker = state.worker(worker_ref)?;
|
|
(
|
|
state.execution_backend.clone(),
|
|
worker.execution_handle.clone(),
|
|
)
|
|
};
|
|
if let (Some(backend), Some(handle)) = (backend, handle) {
|
|
if let Some(snapshot) = backend.worker_snapshot(&handle) {
|
|
return Ok(snapshot);
|
|
}
|
|
}
|
|
Err(RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id: worker_ref.worker_id,
|
|
message: "authoritative Worker snapshot is unavailable".to_string(),
|
|
})
|
|
}
|
|
|
|
/// Replay retained worker-scoped protocol observation events after a cursor.
|
|
#[cfg(feature = "ws-server")]
|
|
pub fn read_worker_observation_events(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
cursor: WorkerObservationCursor,
|
|
) -> Result<Vec<WorkerObservationEvent>, RuntimeError> {
|
|
let state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
state.validate_worker_observation_cursor(worker_ref, cursor)?;
|
|
Ok(state
|
|
.observation_events
|
|
.iter()
|
|
.filter(|event| &event.worker_ref == worker_ref && event.sequence > cursor.sequence)
|
|
.cloned()
|
|
.collect())
|
|
}
|
|
|
|
/// Subscribe to live protocol observation events.
|
|
#[cfg(feature = "ws-server")]
|
|
pub fn subscribe_worker_observation(
|
|
&self,
|
|
) -> Result<broadcast::Receiver<WorkerObservationEvent>, RuntimeError> {
|
|
Ok(self.lock()?.observation_tx.subscribe())
|
|
}
|
|
|
|
/// Append a Worker protocol event to the observation bus.
|
|
#[cfg(feature = "ws-server")]
|
|
pub fn observe_worker_event(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
payload: protocol::Event,
|
|
) -> Result<WorkerObservationEvent, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
let worker_state_changed =
|
|
state.project_protocol_event_to_worker_state(worker_ref, &payload);
|
|
let activity_changed = state.project_internal_worker_activity(worker_ref, &payload);
|
|
if worker_state_changed || activity_changed {
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
}
|
|
if worker_state_changed {
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
}
|
|
let event = state.push_worker_observation_event(worker_ref.clone(), payload);
|
|
Ok(event)
|
|
}
|
|
|
|
/// Snapshot current diagnostics.
|
|
pub fn diagnostics(&self) -> Result<Vec<RuntimeDiagnostic>, RuntimeError> {
|
|
Ok(self.lock()?.diagnostics.clone())
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn record_input_observation(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
input: WorkerInput,
|
|
) -> Result<(), RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
if let Some(payload) = input_protocol_event(&input) {
|
|
state.push_worker_observation_event(worker_ref.clone(), payload);
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(not(feature = "ws-server"))]
|
|
fn record_input_observation(
|
|
&self,
|
|
_worker_ref: &WorkerRef,
|
|
_input: WorkerInput,
|
|
) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
fn transition_worker(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
status: WorkerStatus,
|
|
) -> Result<WorkerLifecycleAck, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
state.ensure_running()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
worker.status = status;
|
|
worker.worker_state = None;
|
|
worker.restore_intent = restore_intent_for_status(status);
|
|
worker.execution_handle = None;
|
|
worker.internal_workers.clear();
|
|
let status = worker.status;
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(WorkerLifecycleAck {
|
|
worker_ref: worker_ref.clone(),
|
|
status,
|
|
worker_state: None,
|
|
})
|
|
}
|
|
|
|
fn install_execution_backend(
|
|
&self,
|
|
backend: Arc<dyn WorkerExecutionBackend>,
|
|
) -> Result<(), RuntimeError> {
|
|
let backend = WorkerExecutionBackendRef::new(backend)?;
|
|
let mut state = self.lock()?;
|
|
state.execution_backend = Some(backend);
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn execution_context(&self, worker_ref: WorkerRef) -> crate::execution::WorkerExecutionContext {
|
|
let runtime = self.clone();
|
|
crate::execution::WorkerExecutionContext::new(
|
|
worker_ref,
|
|
Arc::new(move |worker_ref, payload| runtime.observe_worker_event(&worker_ref, payload)),
|
|
)
|
|
}
|
|
|
|
#[cfg(not(feature = "ws-server"))]
|
|
fn execution_context(&self, worker_ref: WorkerRef) -> crate::execution::WorkerExecutionContext {
|
|
crate::execution::WorkerExecutionContext::new(worker_ref)
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn restore_persisted_worker_executions(&self) -> Result<(), RuntimeError> {
|
|
#[derive(Clone)]
|
|
struct RestoreCandidate {
|
|
worker_ref: WorkerRef,
|
|
request: CreateWorkerRequest,
|
|
run_generation: u64,
|
|
previous_working_directory: Option<CatalogWorkingDirectoryStatus>,
|
|
config_bundle: Option<ConfigBundle>,
|
|
}
|
|
|
|
let candidates = {
|
|
let mut state = self.lock()?;
|
|
if state.execution_backend.is_none() {
|
|
return Ok(());
|
|
}
|
|
let worker_ids = state
|
|
.workers
|
|
.values()
|
|
.filter(|worker| {
|
|
worker.execution_handle.is_none()
|
|
&& worker.execution_bound
|
|
&& worker.status.is_active()
|
|
&& worker.restore_intent == WorkerRestoreIntent::Automatic
|
|
})
|
|
.map(|worker| worker.worker_id)
|
|
.collect::<Vec<_>>();
|
|
let mut candidates = Vec::with_capacity(worker_ids.len());
|
|
for worker_id in worker_ids {
|
|
let (worker_ref, request, previous_working_directory, run_generation) = {
|
|
let worker = state
|
|
.workers
|
|
.get(&worker_id)
|
|
.expect("collected Worker exists");
|
|
(
|
|
worker.worker_ref.clone(),
|
|
worker.request.clone(),
|
|
worker.working_directory.clone(),
|
|
worker.run_generation.saturating_add(1).max(1),
|
|
)
|
|
};
|
|
state
|
|
.workers
|
|
.get_mut(&worker_id)
|
|
.expect("collected Worker exists")
|
|
.run_generation = run_generation;
|
|
state.persist_worker(&worker_id)?;
|
|
candidates.push(RestoreCandidate {
|
|
worker_ref,
|
|
request,
|
|
run_generation,
|
|
previous_working_directory,
|
|
config_bundle: None,
|
|
});
|
|
}
|
|
candidates
|
|
};
|
|
|
|
for candidate in candidates {
|
|
let (backend, workspace_scope) = {
|
|
let state = self.lock()?;
|
|
let workspace_scope = candidate.request.workspace_api.as_ref().and_then(|api| {
|
|
state
|
|
.workspace_owners
|
|
.get(&api.workspace_id)
|
|
.map(|server_id| RuntimeWorkspaceScope::new(&api.workspace_id, server_id))
|
|
});
|
|
(state.execution_backend.clone(), workspace_scope)
|
|
};
|
|
let Some(backend) = backend else {
|
|
return Ok(());
|
|
};
|
|
let request = WorkerExecutionRestoreRequest {
|
|
worker_ref: candidate.worker_ref.clone(),
|
|
run_generation: candidate.run_generation,
|
|
request: candidate.request,
|
|
workspace_scope,
|
|
context: self.execution_context(candidate.worker_ref.clone()),
|
|
previous_working_directory: candidate.previous_working_directory,
|
|
working_directory: None,
|
|
config_bundle: candidate.config_bundle,
|
|
};
|
|
match backend.restore_worker(request) {
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle,
|
|
worker_state,
|
|
working_directory,
|
|
} => self.commit_restored_worker_execution(
|
|
&candidate.worker_ref,
|
|
handle,
|
|
worker_state,
|
|
WorkerStatus::Idle,
|
|
working_directory,
|
|
)?,
|
|
WorkerExecutionSpawnResult::Rejected(result)
|
|
| WorkerExecutionSpawnResult::Errored(result) => {
|
|
let mut state = self.lock()?;
|
|
state.record_restore_failure(&candidate.worker_ref, result)?;
|
|
}
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn commit_restored_worker_execution(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
handle: WorkerExecutionHandle,
|
|
worker_state: protocol::WorkerStateSnapshot,
|
|
status: WorkerStatus,
|
|
working_directory: Option<CatalogWorkingDirectoryStatus>,
|
|
) -> Result<(), RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
state.ensure_worker_ref(worker_ref)?;
|
|
{
|
|
let worker = state.worker_mut(worker_ref)?;
|
|
worker.execution_handle = Some(handle);
|
|
worker.execution_bound = true;
|
|
worker.status = status;
|
|
let _ = worker.apply_worker_state(&worker_state);
|
|
worker.restore_intent = restore_intent_for_status(worker.status);
|
|
worker.working_directory = working_directory;
|
|
}
|
|
state.publish_worker_upsert(worker_ref.worker_id)?;
|
|
state.persist_runtime_snapshot()?;
|
|
state.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Bind the Backend registry identity once. Retention evidence fails closed
|
|
/// until the Runtime host supplies this trusted configuration.
|
|
pub fn bind_runtime_identity(&self, runtime_id: &str) -> Result<(), RuntimeError> {
|
|
if runtime_id.trim().is_empty() || runtime_id.len() > 160 {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Runtime identity must be non-empty and bounded".to_string(),
|
|
));
|
|
}
|
|
let mut state = self.lock()?;
|
|
match state.runtime_identity.as_deref() {
|
|
Some(current) if current == runtime_id => Ok(()),
|
|
Some(_) => Err(RuntimeError::InvalidRequest(
|
|
"Runtime identity is already bound".to_string(),
|
|
)),
|
|
None => {
|
|
state.runtime_identity = Some(runtime_id.to_string());
|
|
Ok(())
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Read canonical aggregate facts needed by a Backend removal plan.
|
|
#[cfg(feature = "fs-store")]
|
|
pub fn worker_retention_inventory(
|
|
&self,
|
|
workspace_id: &str,
|
|
worker_ref: &WorkerRef,
|
|
) -> Result<WorkerRetentionInventory, RuntimeError> {
|
|
let state = self.lock()?;
|
|
let runtime_id = state.runtime_identity.as_deref().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Runtime identity is not bound for Worker retention".to_string(),
|
|
)
|
|
})?;
|
|
let worker = state.worker(worker_ref)?;
|
|
if worker.workspace_id.as_deref() != Some(workspace_id) {
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
});
|
|
}
|
|
let store = state.fs_store().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Worker retention archive authority requires an fs-backed Runtime".to_string(),
|
|
)
|
|
})?;
|
|
FsWorkerRetentionProvider::new(store.runtime_dir()).inventory(
|
|
workspace_id,
|
|
runtime_id,
|
|
worker.worker_id,
|
|
worker.run_generation,
|
|
)
|
|
}
|
|
|
|
/// Enumerate host-authoritative Runtime inventory for Backend orphan
|
|
/// reconciliation. Runtime identity and Workspace scope are derived here,
|
|
/// not accepted in a diagnostic payload.
|
|
#[cfg(feature = "fs-store")]
|
|
pub fn list_worker_retention_inventory(
|
|
&self,
|
|
workspace_id: &str,
|
|
) -> Result<WorkerRetentionInventorySnapshot, RuntimeError> {
|
|
let state = self.lock()?;
|
|
let runtime_id = state.runtime_identity.as_deref().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Runtime identity is not bound for Worker retention".to_string(),
|
|
)
|
|
})?;
|
|
let store = state.fs_store().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Worker retention inventory requires an fs-backed Runtime".to_string(),
|
|
)
|
|
})?;
|
|
let provider = FsWorkerRetentionProvider::new(store.runtime_dir());
|
|
provider.snapshot(workspace_id, runtime_id)
|
|
}
|
|
|
|
/// Execute a Backend-resolved retention plan. Only stopped Workers are
|
|
/// eligible. Provider receipt lookup happens before live lookup so exact
|
|
/// retries converge after aggregate removal.
|
|
#[cfg(feature = "fs-store")]
|
|
pub fn execute_worker_retention(
|
|
&self,
|
|
request: &WorkerRetentionExecutionRequest,
|
|
) -> Result<WorkerRetentionExecutionResult, RuntimeError> {
|
|
let mut state = self.lock()?;
|
|
let runtime_id = state.runtime_identity.clone().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Runtime identity is not bound for Worker retention".to_string(),
|
|
)
|
|
})?;
|
|
if request.source_runtime_id != runtime_id {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Worker retention Runtime identity mismatch".to_string(),
|
|
));
|
|
}
|
|
let store = state.fs_store().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Worker retention execution requires an fs-backed Runtime".to_string(),
|
|
)
|
|
})?;
|
|
let provider = FsWorkerRetentionProvider::new(store.runtime_dir());
|
|
if let Some(completed) = provider.completed_for(request)? {
|
|
state.workers.remove(&request.worker_id);
|
|
state.persist_runtime_snapshot()?;
|
|
return Ok(completed);
|
|
}
|
|
let Some(worker) = state.workers.get(&request.worker_id) else {
|
|
// Recover a pending receipt after a crash between aggregate removal
|
|
// and final receipt/Runtime catalog commit.
|
|
return provider.recover_after_source_removal(request);
|
|
};
|
|
if worker.workspace_id.as_deref() != Some(request.workspace_id.as_str()) {
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: request.worker_id,
|
|
});
|
|
}
|
|
if worker.status != WorkerStatus::Stopped {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Worker retention requires a stopped Worker".to_string(),
|
|
));
|
|
}
|
|
if worker.run_generation != request.expected_run_generation {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"Worker retention plan expected generation {}, current generation is {}",
|
|
request.expected_run_generation, worker.run_generation
|
|
)));
|
|
}
|
|
let result = provider.execute(request)?;
|
|
state.workers.remove(&request.worker_id);
|
|
state.persist_runtime_snapshot()?;
|
|
Ok(result)
|
|
}
|
|
|
|
fn worker_operation_lock(&self, worker_id: WorkerId) -> Result<Arc<Mutex<()>>, RuntimeError> {
|
|
let mut operations = self
|
|
.worker_operations
|
|
.lock()
|
|
.map_err(|_| RuntimeError::StatePoisoned)?;
|
|
Ok(operations
|
|
.entry(worker_id)
|
|
.or_insert_with(|| Arc::new(Mutex::new(())))
|
|
.clone())
|
|
}
|
|
|
|
fn lock(&self) -> Result<MutexGuard<'_, RuntimeState>, RuntimeError> {
|
|
self.inner.lock().map_err(|_| RuntimeError::StatePoisoned)
|
|
}
|
|
}
|
|
|
|
#[cfg_attr(not(feature = "fs-store"), allow(dead_code))]
|
|
#[derive(Clone, Debug)]
|
|
enum RuntimePersistence {
|
|
Memory,
|
|
#[cfg(feature = "fs-store")]
|
|
Fs(FsRuntimeStore),
|
|
}
|
|
|
|
#[derive(Debug)]
|
|
struct SubscriptionSink {
|
|
selector: EventSubscriptionSelector,
|
|
workspace_id: Option<String>,
|
|
sender: mpsc::Sender<RuntimeSubscriptionUpdate>,
|
|
lagged: Arc<AtomicBool>,
|
|
}
|
|
|
|
#[derive(Clone)]
|
|
struct BackendResourceClientRef(Arc<dyn BackendResourceClient>);
|
|
|
|
impl std::fmt::Debug for BackendResourceClientRef {
|
|
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
formatter.write_str("BackendResourceClientRef(..)")
|
|
}
|
|
}
|
|
|
|
#[derive(Debug)]
|
|
struct RuntimeState {
|
|
display_name: Option<String>,
|
|
backend: RuntimeBackendKind,
|
|
/// Backend-bound stable identity used for cross-boundary retention evidence.
|
|
/// It is configured once by the Runtime host and never model input.
|
|
runtime_identity: Option<String>,
|
|
#[cfg_attr(not(feature = "fs-store"), allow(dead_code))]
|
|
persistence: RuntimePersistence,
|
|
status: RuntimeStatus,
|
|
execution_backend: Option<WorkerExecutionBackendRef>,
|
|
backend_resource_client: Option<BackendResourceClientRef>,
|
|
workspace_backend_resource_clients: BTreeMap<String, BackendResourceClientRef>,
|
|
#[cfg(feature = "fs-store")]
|
|
next_diagnostic_id: u64,
|
|
workers: BTreeMap<WorkerId, WorkerRecord>,
|
|
workspace_owners: BTreeMap<String, String>,
|
|
config_bundles: BTreeMap<String, ConfigBundle>,
|
|
workspace_config_latest: BTreeMap<String, ConfigBundleRef>,
|
|
workspace_config_fetch_gates: BTreeMap<String, Arc<Mutex<()>>>,
|
|
diagnostics: Vec<RuntimeDiagnostic>,
|
|
subscription_revision: u64,
|
|
worker_subject_revisions: BTreeMap<WorkerId, u64>,
|
|
next_event_subscription_id: u64,
|
|
subscriptions: BTreeMap<u64, SubscriptionSink>,
|
|
#[cfg(feature = "ws-server")]
|
|
next_observation_sequence: u64,
|
|
#[cfg(feature = "ws-server")]
|
|
observation_events: VecDeque<WorkerObservationEvent>,
|
|
#[cfg(feature = "ws-server")]
|
|
observation_tx: broadcast::Sender<WorkerObservationEvent>,
|
|
}
|
|
|
|
impl RuntimeState {
|
|
fn new(display_name: Option<String>) -> Self {
|
|
Self {
|
|
display_name,
|
|
backend: RuntimeBackendKind::Memory,
|
|
runtime_identity: None,
|
|
persistence: RuntimePersistence::Memory,
|
|
status: RuntimeStatus::Running,
|
|
execution_backend: None,
|
|
backend_resource_client: None,
|
|
workspace_backend_resource_clients: BTreeMap::new(),
|
|
#[cfg(feature = "fs-store")]
|
|
next_diagnostic_id: 1,
|
|
workers: BTreeMap::new(),
|
|
workspace_owners: BTreeMap::new(),
|
|
config_bundles: BTreeMap::new(),
|
|
workspace_config_latest: BTreeMap::new(),
|
|
workspace_config_fetch_gates: BTreeMap::new(),
|
|
diagnostics: Vec::new(),
|
|
subscription_revision: 0,
|
|
worker_subject_revisions: BTreeMap::new(),
|
|
next_event_subscription_id: 1,
|
|
subscriptions: BTreeMap::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
next_observation_sequence: 1,
|
|
#[cfg(feature = "ws-server")]
|
|
observation_events: VecDeque::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
observation_tx: broadcast::channel(256).0,
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn new_fs_backed(display_name: Option<String>, store: FsRuntimeStore) -> Self {
|
|
Self {
|
|
display_name,
|
|
backend: RuntimeBackendKind::FsStore,
|
|
runtime_identity: None,
|
|
persistence: RuntimePersistence::Fs(store),
|
|
status: RuntimeStatus::Running,
|
|
execution_backend: None,
|
|
backend_resource_client: None,
|
|
workspace_backend_resource_clients: BTreeMap::new(),
|
|
#[cfg(feature = "fs-store")]
|
|
next_diagnostic_id: 1,
|
|
workers: BTreeMap::new(),
|
|
workspace_owners: BTreeMap::new(),
|
|
config_bundles: BTreeMap::new(),
|
|
workspace_config_latest: BTreeMap::new(),
|
|
workspace_config_fetch_gates: BTreeMap::new(),
|
|
diagnostics: Vec::new(),
|
|
subscription_revision: 0,
|
|
worker_subject_revisions: BTreeMap::new(),
|
|
next_event_subscription_id: 1,
|
|
subscriptions: BTreeMap::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
next_observation_sequence: 1,
|
|
#[cfg(feature = "ws-server")]
|
|
observation_events: VecDeque::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
observation_tx: broadcast::channel(256).0,
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn from_persisted(
|
|
persisted: PersistedRuntimeState,
|
|
store: FsRuntimeStore,
|
|
) -> Result<Self, RuntimeError> {
|
|
let mut workers = BTreeMap::new();
|
|
let diagnostics = persisted.diagnostics;
|
|
let next_diagnostic_id = persisted.next_diagnostic_id;
|
|
for (worker_id, worker) in persisted.workers {
|
|
let run_generation = worker.execution.last_run_generation;
|
|
workers.insert(
|
|
worker_id,
|
|
WorkerRecord {
|
|
worker_ref: worker.worker_ref,
|
|
worker_id: worker.worker_id,
|
|
status: worker.status,
|
|
worker_state: None,
|
|
workspace_id: worker.workspace_id,
|
|
request: worker.request,
|
|
run_generation,
|
|
execution_bound: worker.execution.binding.is_some(),
|
|
restore_intent: worker.execution.restore_intent,
|
|
working_directory: worker.working_directory,
|
|
execution_handle: None,
|
|
internal_workers: BTreeMap::new(),
|
|
},
|
|
);
|
|
}
|
|
|
|
Ok(Self {
|
|
display_name: persisted.display_name,
|
|
backend: RuntimeBackendKind::FsStore,
|
|
runtime_identity: None,
|
|
persistence: RuntimePersistence::Fs(store),
|
|
status: persisted.status,
|
|
execution_backend: None,
|
|
backend_resource_client: None,
|
|
workspace_backend_resource_clients: BTreeMap::new(),
|
|
next_diagnostic_id,
|
|
workers,
|
|
config_bundles: BTreeMap::new(),
|
|
workspace_config_latest: BTreeMap::new(),
|
|
workspace_config_fetch_gates: BTreeMap::new(),
|
|
workspace_owners: persisted.workspace_owners,
|
|
diagnostics,
|
|
subscription_revision: 0,
|
|
worker_subject_revisions: BTreeMap::new(),
|
|
next_event_subscription_id: 1,
|
|
subscriptions: BTreeMap::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
next_observation_sequence: 1,
|
|
#[cfg(feature = "ws-server")]
|
|
observation_events: VecDeque::new(),
|
|
#[cfg(feature = "ws-server")]
|
|
observation_tx: broadcast::channel(256).0,
|
|
})
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn persisted_state(&self) -> PersistedRuntimeState {
|
|
PersistedRuntimeState {
|
|
display_name: self.display_name.clone(),
|
|
status: self.status,
|
|
next_diagnostic_id: self.next_diagnostic_id,
|
|
workers: self
|
|
.workers
|
|
.iter()
|
|
.map(|(worker_id, worker)| (worker_id.clone(), worker.persisted_record()))
|
|
.collect(),
|
|
workspace_owners: self.workspace_owners.clone(),
|
|
diagnostics: self.diagnostics.clone(),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn fs_store(&self) -> Option<&FsRuntimeStore> {
|
|
match &self.persistence {
|
|
RuntimePersistence::Memory => None,
|
|
RuntimePersistence::Fs(store) => Some(store),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn persist_runtime_snapshot(&self) -> Result<(), RuntimeError> {
|
|
if let Some(store) = self.fs_store() {
|
|
store.write_runtime_snapshot(&self.persisted_state())?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn persist_worker(&self, worker_id: &WorkerId) -> Result<(), RuntimeError> {
|
|
if let Some(store) = self.fs_store() {
|
|
let worker =
|
|
self.workers
|
|
.get(worker_id)
|
|
.ok_or_else(|| RuntimeError::WorkerNotFound {
|
|
worker_id: *worker_id,
|
|
})?;
|
|
store.write_worker_snapshot(&worker.persisted_record())?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn delete_worker_snapshot(&self, worker_id: &WorkerId) -> Result<(), RuntimeError> {
|
|
if let Some(store) = self.fs_store() {
|
|
store.delete_worker_snapshot(worker_id)?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn persist_workers(&self) -> Result<(), RuntimeError> {
|
|
if self.fs_store().is_some() {
|
|
for worker_id in self.workers.keys() {
|
|
self.persist_worker(worker_id)?;
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(not(feature = "fs-store"))]
|
|
fn persist_runtime_snapshot(&self) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(not(feature = "fs-store"))]
|
|
fn persist_worker(&self, _worker_id: &WorkerId) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(not(feature = "fs-store"))]
|
|
fn delete_worker_snapshot(&self, _worker_id: &WorkerId) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(not(feature = "fs-store"))]
|
|
fn persist_workers(&self) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
fn ensure_running(&self) -> Result<(), RuntimeError> {
|
|
if self.status == RuntimeStatus::Stopped {
|
|
Err(RuntimeError::RuntimeStopped)
|
|
} else {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
fn check_config_bundle_ref(
|
|
&self,
|
|
reference: &ConfigBundleRef,
|
|
) -> Result<ConfigBundleAvailability, RuntimeError> {
|
|
validate_config_bundle_ref(reference)?;
|
|
let bundle = self.config_bundles.get(&reference.id).ok_or_else(|| {
|
|
RuntimeError::ConfigBundleMissing {
|
|
bundle_id: reference.id.clone(),
|
|
}
|
|
})?;
|
|
if bundle.metadata.digest != reference.digest {
|
|
return Err(RuntimeError::ConfigBundleDigestMismatch {
|
|
bundle_id: reference.id.clone(),
|
|
expected_digest: reference.digest.clone(),
|
|
actual_digest: bundle.metadata.digest.clone(),
|
|
});
|
|
}
|
|
Ok(ConfigBundleAvailability {
|
|
reference: reference.clone(),
|
|
summary: bundle.summary(),
|
|
})
|
|
}
|
|
|
|
fn resolve_config_bundle_ref(
|
|
&self,
|
|
reference: Option<&ConfigBundleRef>,
|
|
) -> Result<Option<ConfigBundle>, RuntimeError> {
|
|
let Some(reference) = reference else {
|
|
return Ok(None);
|
|
};
|
|
self.check_config_bundle_ref(reference)?;
|
|
Ok(self.config_bundles.get(&reference.id).cloned())
|
|
}
|
|
|
|
fn validate_worker_config_boundary(
|
|
&self,
|
|
_request: &CreateWorkerRequest,
|
|
) -> Result<(), RuntimeError> {
|
|
Ok(())
|
|
}
|
|
|
|
fn ensure_worker_ref(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> {
|
|
if !self.workers.contains_key(&worker_ref.worker_id) {
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
});
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn ensure_workspace_owner(
|
|
&mut self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
claim_if_missing: bool,
|
|
) -> Result<bool, RuntimeError> {
|
|
if scope.workspace_id.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Runtime auth workspace_id must not be empty".to_string(),
|
|
));
|
|
}
|
|
if scope.server_id.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Runtime auth server_id must not be empty".to_string(),
|
|
));
|
|
}
|
|
match self.workspace_owners.get(&scope.workspace_id) {
|
|
Some(owner_server_id) if owner_server_id == &scope.server_id => Ok(true),
|
|
Some(owner_server_id) => Err(RuntimeError::WorkspaceOwnerMismatch {
|
|
workspace_id: scope.workspace_id.clone(),
|
|
owner_server_id: owner_server_id.clone(),
|
|
requester_server_id: scope.server_id.clone(),
|
|
}),
|
|
None if claim_if_missing => {
|
|
self.workspace_owners
|
|
.insert(scope.workspace_id.clone(), scope.server_id.clone());
|
|
Ok(true)
|
|
}
|
|
None => Ok(false),
|
|
}
|
|
}
|
|
|
|
fn ensure_workspace_owner_for_existing_workers(
|
|
&mut self,
|
|
scope: &RuntimeWorkspaceScope,
|
|
) -> Result<bool, RuntimeError> {
|
|
let has_workspace_worker = self
|
|
.workers
|
|
.values()
|
|
.any(|worker| worker.belongs_to_workspace(&scope.workspace_id));
|
|
self.ensure_workspace_owner(scope, has_workspace_worker)
|
|
}
|
|
|
|
fn forget_workspace_owner_if_unused(&mut self, workspace_id: &str) -> bool {
|
|
if self
|
|
.workers
|
|
.values()
|
|
.any(|worker| worker.belongs_to_workspace(workspace_id))
|
|
{
|
|
return false;
|
|
}
|
|
self.workspace_owners.remove(workspace_id).is_some()
|
|
}
|
|
|
|
fn worker(&self, worker_ref: &WorkerRef) -> Result<&WorkerRecord, RuntimeError> {
|
|
self.ensure_worker_ref(worker_ref)?;
|
|
self.workers
|
|
.get(&worker_ref.worker_id)
|
|
.ok_or_else(|| RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
})
|
|
}
|
|
|
|
fn worker_mut(&mut self, worker_ref: &WorkerRef) -> Result<&mut WorkerRecord, RuntimeError> {
|
|
self.ensure_worker_ref(worker_ref)?;
|
|
self.workers
|
|
.get_mut(&worker_ref.worker_id)
|
|
.ok_or_else(|| RuntimeError::WorkerNotFound {
|
|
worker_id: worker_ref.worker_id,
|
|
})
|
|
}
|
|
|
|
fn subscription_snapshot(
|
|
&self,
|
|
scope: Option<&RuntimeWorkspaceScope>,
|
|
selector: &EventSubscriptionSelector,
|
|
) -> Result<SubscriptionSnapshot, RuntimeError> {
|
|
let workers = match selector {
|
|
EventSubscriptionSelector::RuntimeWorkers => self
|
|
.workers
|
|
.values()
|
|
.filter(|worker| {
|
|
scope.is_none_or(|scope| worker.belongs_to_workspace(&scope.workspace_id))
|
|
})
|
|
.map(|worker| self.subscription_worker(worker))
|
|
.collect::<Result<Vec<_>, _>>()?,
|
|
EventSubscriptionSelector::WorkerLifecycle { worker_ids } => {
|
|
let mut selected = Vec::with_capacity(worker_ids.as_slice().len());
|
|
for worker_id in worker_ids.as_slice() {
|
|
let runtime_worker_id = WorkerId::parse(worker_id.as_str()).ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(format!(
|
|
"worker_lifecycle selector contains invalid Runtime Worker id {worker_id}"
|
|
))
|
|
})?;
|
|
let worker = self.workers.get(&runtime_worker_id).ok_or(
|
|
RuntimeError::WorkerNotFound {
|
|
worker_id: runtime_worker_id,
|
|
},
|
|
)?;
|
|
if scope.is_some_and(|scope| !worker.belongs_to_workspace(&scope.workspace_id))
|
|
{
|
|
return Err(RuntimeError::WorkerNotFound {
|
|
worker_id: runtime_worker_id,
|
|
});
|
|
}
|
|
selected.push(self.subscription_worker(worker)?);
|
|
}
|
|
selected
|
|
}
|
|
EventSubscriptionSelector::WorkerProtocol { .. }
|
|
| EventSubscriptionSelector::WorkspaceWorkers
|
|
| EventSubscriptionSelector::WorkspaceWorkdirs => {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"selector is not produced by the Runtime lifecycle subscription".to_string(),
|
|
));
|
|
}
|
|
};
|
|
Ok(SubscriptionSnapshot::Workers { workers })
|
|
}
|
|
|
|
fn subscription_worker(
|
|
&self,
|
|
worker: &WorkerRecord,
|
|
) -> Result<SubscriptionWorker, RuntimeError> {
|
|
let worker_id = SubscriptionWorkerId::new(worker.worker_id.to_string())
|
|
.map_err(subscription_validation_error)?;
|
|
let repository_id = worker
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|working_directory| working_directory.summary.repository_id.clone());
|
|
let working_directory_id = worker
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|working_directory| {
|
|
SubscriptionWorkdirId::new(working_directory.summary.working_directory_id.clone())
|
|
.map_err(subscription_validation_error)
|
|
})
|
|
.transpose()?;
|
|
let profile = match &worker.request.profile {
|
|
ProfileSelector::Builtin(name) | ProfileSelector::Named(name) => Some(name.clone()),
|
|
};
|
|
Ok(SubscriptionWorker {
|
|
worker_id,
|
|
runtime_id: None,
|
|
resource_key: None,
|
|
subject_revision: self
|
|
.worker_subject_revisions
|
|
.get(&worker.worker_id)
|
|
.copied()
|
|
.unwrap_or(0),
|
|
worker_state: worker.worker_state.clone(),
|
|
state: subscription_worker_state(worker.status),
|
|
has_running_internal_workers: worker
|
|
.internal_workers
|
|
.values()
|
|
.any(|worker| worker.status == protocol::WorkerStatus::Running),
|
|
workspace_id: worker.workspace_id.clone(),
|
|
display_name: worker.request.display_name.clone(),
|
|
profile,
|
|
repository_id,
|
|
repository_key: None,
|
|
working_directory_id,
|
|
})
|
|
}
|
|
|
|
fn publish_worker_upsert(&mut self, worker_id: WorkerId) -> Result<(), RuntimeError> {
|
|
self.subscription_revision = self.subscription_revision.saturating_add(1);
|
|
let subject_revision = {
|
|
let revision = self.worker_subject_revisions.entry(worker_id).or_insert(0);
|
|
*revision = revision.saturating_add(1);
|
|
*revision
|
|
};
|
|
let worker = self
|
|
.workers
|
|
.get(&worker_id)
|
|
.ok_or(RuntimeError::WorkerNotFound { worker_id })?;
|
|
let workspace_id = worker.workspace_id.clone();
|
|
let mut projected = self.subscription_worker(worker)?;
|
|
projected.subject_revision = subject_revision;
|
|
let projected_worker_id = projected.worker_id.clone();
|
|
self.deliver_worker_subscription_update(
|
|
&projected_worker_id,
|
|
workspace_id.as_deref(),
|
|
RuntimeSubscriptionUpdate {
|
|
subject_revision,
|
|
payload: SubscriptionEventPayload::WorkerUpserted { worker: projected },
|
|
},
|
|
);
|
|
Ok(())
|
|
}
|
|
|
|
fn publish_worker_removed(
|
|
&mut self,
|
|
worker_id: WorkerId,
|
|
workspace_id: Option<&str>,
|
|
) -> Result<(), RuntimeError> {
|
|
self.subscription_revision = self.subscription_revision.saturating_add(1);
|
|
let subject_revision = {
|
|
let revision = self.worker_subject_revisions.entry(worker_id).or_insert(0);
|
|
*revision = revision.saturating_add(1);
|
|
*revision
|
|
};
|
|
let worker_id = SubscriptionWorkerId::new(worker_id.to_string())
|
|
.map_err(subscription_validation_error)?;
|
|
self.deliver_worker_subscription_update(
|
|
&worker_id,
|
|
workspace_id,
|
|
RuntimeSubscriptionUpdate {
|
|
subject_revision,
|
|
payload: SubscriptionEventPayload::WorkerRemoved {
|
|
worker_id: worker_id.clone(),
|
|
runtime_id: None,
|
|
},
|
|
},
|
|
);
|
|
Ok(())
|
|
}
|
|
|
|
fn deliver_worker_subscription_update(
|
|
&mut self,
|
|
worker_id: &SubscriptionWorkerId,
|
|
workspace_id: Option<&str>,
|
|
update: RuntimeSubscriptionUpdate,
|
|
) {
|
|
let mut closed = Vec::new();
|
|
for (subscription_id, sink) in &self.subscriptions {
|
|
if sink
|
|
.workspace_id
|
|
.as_deref()
|
|
.is_some_and(|expected| workspace_id != Some(expected))
|
|
{
|
|
continue;
|
|
}
|
|
let selected = match &sink.selector {
|
|
EventSubscriptionSelector::RuntimeWorkers => true,
|
|
EventSubscriptionSelector::WorkerLifecycle { worker_ids } => {
|
|
worker_ids.contains(worker_id)
|
|
}
|
|
EventSubscriptionSelector::WorkerProtocol { .. }
|
|
| EventSubscriptionSelector::WorkspaceWorkers
|
|
| EventSubscriptionSelector::WorkspaceWorkdirs => false,
|
|
};
|
|
if !selected {
|
|
continue;
|
|
}
|
|
match sink.sender.try_send(update.clone()) {
|
|
Ok(()) => {}
|
|
Err(mpsc::error::TrySendError::Full(_)) => {
|
|
sink.lagged.store(true, Ordering::Release);
|
|
closed.push(*subscription_id);
|
|
}
|
|
Err(mpsc::error::TrySendError::Closed(_)) => closed.push(*subscription_id),
|
|
}
|
|
}
|
|
for subscription_id in closed {
|
|
self.subscriptions.remove(&subscription_id);
|
|
}
|
|
}
|
|
|
|
fn primary_worker_id_for_workdir(&self, working_directory_id: &str) -> Option<WorkerId> {
|
|
self.workers.values().find_map(|worker| {
|
|
if worker
|
|
.working_directory
|
|
.as_ref()
|
|
.is_some_and(|binding| binding.summary.working_directory_id == working_directory_id)
|
|
|| requested_primary_workdir_id(&worker.request) == Some(working_directory_id)
|
|
{
|
|
Some(worker.worker_id)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn record_restore_failure(
|
|
&mut self,
|
|
worker_ref: &WorkerRef,
|
|
result: WorkerExecutionResult,
|
|
) -> Result<(), RuntimeError> {
|
|
let message = result
|
|
.message
|
|
.clone()
|
|
.unwrap_or_else(|| "worker execution restore failed".to_string());
|
|
let diagnostic_id = self.next_diagnostic_id;
|
|
self.next_diagnostic_id += 1;
|
|
let workspace_id = self
|
|
.workers
|
|
.get(&worker_ref.worker_id)
|
|
.and_then(|worker| worker.workspace_id.as_deref())
|
|
.unwrap_or("<unscoped>");
|
|
eprintln!(
|
|
"yoi-runtime: Worker execution restore failed: diagnostic_id={diagnostic_id} runtime_id={} workspace_id={workspace_id} worker_id={} operation={:?} outcome={:?}: {message}",
|
|
self.runtime_identity.as_deref().unwrap_or("<unbound>"),
|
|
worker_ref.worker_id,
|
|
result.operation,
|
|
result.outcome,
|
|
);
|
|
self.diagnostics.push(RuntimeDiagnostic {
|
|
id: diagnostic_id,
|
|
severity: DiagnosticSeverity::Warning,
|
|
code: "worker_execution_restore_failed".to_string(),
|
|
message: format!(
|
|
"worker {} execution restore failed: {message}",
|
|
worker_ref.worker_id
|
|
),
|
|
worker_ref: Some(worker_ref.clone()),
|
|
});
|
|
let worker = self.worker_mut(worker_ref)?;
|
|
worker.execution_handle = None;
|
|
worker.status = WorkerStatus::Stopped;
|
|
worker.restore_intent = WorkerRestoreIntent::Explicit;
|
|
worker.internal_workers.clear();
|
|
self.publish_worker_upsert(worker_ref.worker_id)?;
|
|
self.persist_runtime_snapshot()?;
|
|
self.persist_worker(&worker_ref.worker_id)?;
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn validate_worker_observation_cursor(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
cursor: WorkerObservationCursor,
|
|
) -> Result<(), RuntimeError> {
|
|
if let Some(first) = self
|
|
.observation_events
|
|
.iter()
|
|
.find(|event| &event.worker_ref == worker_ref)
|
|
{
|
|
if cursor.sequence != 0 && cursor.sequence < first.sequence {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker observation cursor {} is expired for worker {}",
|
|
cursor.encode(),
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
}
|
|
if cursor.sequence >= self.next_observation_sequence {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"worker observation cursor {} is unknown for worker {}",
|
|
cursor.encode(),
|
|
worker_ref.worker_id
|
|
)));
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn push_worker_observation_event(
|
|
&mut self,
|
|
worker_ref: WorkerRef,
|
|
payload: protocol::Event,
|
|
) -> WorkerObservationEvent {
|
|
const MAX_OBSERVATION_BACKLOG: usize = 1024;
|
|
|
|
let sequence = self.next_observation_sequence;
|
|
self.next_observation_sequence += 1;
|
|
let event = WorkerObservationEvent::new(sequence, worker_ref, payload);
|
|
self.observation_events.push_back(event.clone());
|
|
while self.observation_events.len() > MAX_OBSERVATION_BACKLOG {
|
|
self.observation_events.pop_front();
|
|
}
|
|
let _ = self.observation_tx.send(event.clone());
|
|
event
|
|
}
|
|
|
|
fn internal_worker_snapshot_statuses(
|
|
statuses: &mut BTreeMap<String, InternalWorkerActivity>,
|
|
snapshot: &protocol::InternalWorkerSnapshot,
|
|
) {
|
|
statuses.insert(
|
|
snapshot.worker.session_id.clone(),
|
|
InternalWorkerActivity {
|
|
status: snapshot.status,
|
|
parent_session_id: snapshot.worker.parent_session_id.clone(),
|
|
},
|
|
);
|
|
for child in &snapshot.internal_workers {
|
|
Self::internal_worker_snapshot_statuses(statuses, child);
|
|
}
|
|
}
|
|
|
|
fn remove_internal_worker_subtree(
|
|
statuses: &mut BTreeMap<String, InternalWorkerActivity>,
|
|
root_session_id: &str,
|
|
) {
|
|
let mut removed = vec![root_session_id.to_string()];
|
|
while let Some(parent_session_id) = removed.pop() {
|
|
let children = statuses
|
|
.iter()
|
|
.filter_map(|(session_id, worker)| {
|
|
(worker.parent_session_id.as_deref() == Some(parent_session_id.as_str()))
|
|
.then(|| session_id.clone())
|
|
})
|
|
.collect::<Vec<_>>();
|
|
statuses.remove(&parent_session_id);
|
|
removed.extend(children);
|
|
}
|
|
}
|
|
|
|
fn project_internal_worker_event(
|
|
statuses: &mut BTreeMap<String, InternalWorkerActivity>,
|
|
worker: &protocol::InternalWorkerRef,
|
|
event: &protocol::Event,
|
|
) {
|
|
match event {
|
|
protocol::Event::Snapshot {
|
|
state,
|
|
internal_workers,
|
|
..
|
|
} => {
|
|
Self::remove_internal_worker_subtree(statuses, &worker.session_id);
|
|
statuses.insert(
|
|
worker.session_id.clone(),
|
|
InternalWorkerActivity {
|
|
status: state.catalog_status(),
|
|
parent_session_id: worker.parent_session_id.clone(),
|
|
},
|
|
);
|
|
for child in internal_workers {
|
|
Self::internal_worker_snapshot_statuses(statuses, child);
|
|
}
|
|
}
|
|
protocol::Event::InternalWorker {
|
|
worker: nested_worker,
|
|
event,
|
|
..
|
|
} => Self::project_internal_worker_event(statuses, nested_worker, event),
|
|
protocol::Event::WorkerState { snapshot }
|
|
| protocol::Event::CommandAcknowledged {
|
|
acknowledgement:
|
|
protocol::WorkerCommandAcknowledgement {
|
|
state: snapshot, ..
|
|
},
|
|
} => {
|
|
statuses.insert(
|
|
worker.session_id.clone(),
|
|
InternalWorkerActivity {
|
|
status: snapshot.catalog_status(),
|
|
parent_session_id: worker.parent_session_id.clone(),
|
|
},
|
|
);
|
|
}
|
|
_ => {}
|
|
}
|
|
}
|
|
|
|
fn update_internal_worker_activity(
|
|
statuses: &mut BTreeMap<String, InternalWorkerActivity>,
|
|
event: &protocol::Event,
|
|
) -> bool {
|
|
let was_running = statuses
|
|
.values()
|
|
.any(|worker| worker.status == protocol::WorkerStatus::Running);
|
|
match event {
|
|
protocol::Event::Snapshot {
|
|
internal_workers, ..
|
|
} => {
|
|
statuses.clear();
|
|
for child in internal_workers {
|
|
Self::internal_worker_snapshot_statuses(statuses, child);
|
|
}
|
|
}
|
|
protocol::Event::InternalWorker { worker, event, .. } => {
|
|
Self::project_internal_worker_event(statuses, worker, event);
|
|
}
|
|
_ => {}
|
|
}
|
|
let is_running = statuses
|
|
.values()
|
|
.any(|worker| worker.status == protocol::WorkerStatus::Running);
|
|
was_running != is_running
|
|
}
|
|
|
|
fn project_internal_worker_activity(
|
|
&mut self,
|
|
worker_ref: &WorkerRef,
|
|
event: &protocol::Event,
|
|
) -> bool {
|
|
let Some(worker) = self.workers.get_mut(&worker_ref.worker_id) else {
|
|
return false;
|
|
};
|
|
Self::update_internal_worker_activity(&mut worker.internal_workers, event)
|
|
}
|
|
|
|
fn project_protocol_event_to_worker_state(
|
|
&mut self,
|
|
worker_ref: &WorkerRef,
|
|
event: &protocol::Event,
|
|
) -> bool {
|
|
let Some(worker) = self.workers.get_mut(&worker_ref.worker_id) else {
|
|
return false;
|
|
};
|
|
let incoming = match event {
|
|
protocol::Event::WorkerState { snapshot }
|
|
| protocol::Event::Snapshot {
|
|
state: snapshot, ..
|
|
}
|
|
| protocol::Event::CommandAcknowledged {
|
|
acknowledgement:
|
|
protocol::WorkerCommandAcknowledgement {
|
|
state: snapshot, ..
|
|
},
|
|
} => snapshot,
|
|
_ => return false,
|
|
};
|
|
match worker.apply_worker_state(incoming) {
|
|
Ok(protocol::WorkerStateSnapshotApply::Applied) => true,
|
|
Ok(
|
|
protocol::WorkerStateSnapshotApply::Duplicate
|
|
| protocol::WorkerStateSnapshotApply::Stale,
|
|
)
|
|
| Err(_) => false,
|
|
}
|
|
}
|
|
}
|
|
|
|
#[derive(Debug, Clone)]
|
|
struct InternalWorkerActivity {
|
|
status: protocol::WorkerStatus,
|
|
parent_session_id: Option<String>,
|
|
}
|
|
|
|
#[derive(Debug)]
|
|
struct WorkerRecord {
|
|
worker_ref: WorkerRef,
|
|
worker_id: WorkerId,
|
|
status: WorkerStatus,
|
|
worker_state: Option<protocol::WorkerStateSnapshot>,
|
|
workspace_id: Option<String>,
|
|
request: CreateWorkerRequest,
|
|
run_generation: u64,
|
|
execution_bound: bool,
|
|
restore_intent: WorkerRestoreIntent,
|
|
working_directory: Option<CatalogWorkingDirectoryStatus>,
|
|
execution_handle: Option<WorkerExecutionHandle>,
|
|
internal_workers: BTreeMap<String, InternalWorkerActivity>,
|
|
}
|
|
|
|
impl WorkerRecord {
|
|
fn apply_worker_state(
|
|
&mut self,
|
|
incoming: &protocol::WorkerStateSnapshot,
|
|
) -> Result<protocol::WorkerStateSnapshotApply, protocol::WorkerStateSnapshotConflict> {
|
|
match self.worker_state.as_mut() {
|
|
Some(current) => protocol::apply_worker_state_snapshot(current, incoming),
|
|
None => {
|
|
self.worker_state = Some(incoming.clone());
|
|
Ok(protocol::WorkerStateSnapshotApply::Applied)
|
|
}
|
|
}
|
|
}
|
|
|
|
fn belongs_to_workspace(&self, workspace_id: &str) -> bool {
|
|
self.workspace_id.as_deref() == Some(workspace_id)
|
|
}
|
|
|
|
fn summary(&self) -> WorkerSummary {
|
|
WorkerSummary {
|
|
worker_ref: self.worker_ref.clone(),
|
|
worker_id: self.worker_id,
|
|
status: self.status,
|
|
worker_state: self.worker_state.clone(),
|
|
workspace_id: self.workspace_id.clone(),
|
|
working_directory: self.working_directory.clone(),
|
|
profile: self.request.profile.clone(),
|
|
display_name: self.request.display_name.clone(),
|
|
profile_source: self.request.profile_source.reference(),
|
|
config_bundle: self.request.config_bundle.clone(),
|
|
}
|
|
}
|
|
|
|
fn detail(&self) -> WorkerDetail {
|
|
WorkerDetail {
|
|
worker_ref: self.worker_ref.clone(),
|
|
worker_id: self.worker_id,
|
|
status: self.status,
|
|
worker_state: self.worker_state.clone(),
|
|
workspace_id: self.workspace_id.clone(),
|
|
working_directory: self.working_directory.clone(),
|
|
profile: self.request.profile.clone(),
|
|
display_name: self.request.display_name.clone(),
|
|
profile_source: self.request.profile_source.reference(),
|
|
config_bundle: self.request.config_bundle.clone(),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn persisted_record(&self) -> PersistedWorkerRecord {
|
|
PersistedWorkerRecord {
|
|
worker_ref: self.worker_ref.clone(),
|
|
worker_id: self.worker_id.clone(),
|
|
request: self.request.clone(),
|
|
status: self.status,
|
|
execution: PersistedWorkerExecution {
|
|
last_run_generation: self.run_generation,
|
|
binding: self
|
|
.execution_bound
|
|
.then_some(PersistedWorkerExecutionBinding {
|
|
run_generation: self.run_generation,
|
|
}),
|
|
restore_intent: self.restore_intent,
|
|
},
|
|
workspace_id: self.workspace_id.clone(),
|
|
working_directory: self.working_directory.clone(),
|
|
}
|
|
}
|
|
}
|
|
|
|
fn restore_intent_for_status(status: WorkerStatus) -> WorkerRestoreIntent {
|
|
if status.is_active() {
|
|
WorkerRestoreIntent::Automatic
|
|
} else {
|
|
WorkerRestoreIntent::Explicit
|
|
}
|
|
}
|
|
|
|
fn repository_resource_error(error: BackendResourceError) -> RuntimeError {
|
|
let (code, message) = match error {
|
|
BackendResourceError::Expired => (
|
|
"repository_access_credential_expired",
|
|
"Repository access credential lease expired",
|
|
),
|
|
BackendResourceError::Unauthorized { .. } => (
|
|
"repository_access_credential_unauthorized",
|
|
"Repository access credential lease was rejected",
|
|
),
|
|
BackendResourceError::MissingResource => (
|
|
"repository_access_credential_unavailable",
|
|
"Repository access credential lease is unavailable or already consumed",
|
|
),
|
|
BackendResourceError::Timeout => (
|
|
"repository_access_resource_fetch_timeout",
|
|
"Timed out while fetching Repository SSH access from Workspace Backend",
|
|
),
|
|
BackendResourceError::Transport { .. } => (
|
|
"repository_access_provider_unavailable",
|
|
"Repository access credential provider is unavailable",
|
|
),
|
|
BackendResourceError::UnsupportedKind
|
|
| BackendResourceError::Oversized { .. }
|
|
| BackendResourceError::DigestMismatch { .. }
|
|
| BackendResourceError::ContentTypeMismatch { .. }
|
|
| BackendResourceError::InvalidResponse { .. } => (
|
|
"repository_access_credential_invalid",
|
|
"Repository access credential response is invalid",
|
|
),
|
|
};
|
|
RuntimeError::WorkingDirectory(
|
|
crate::working_directory::WorkingDirectoryDiagnostic::rejected(code, message),
|
|
)
|
|
}
|
|
|
|
fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest {
|
|
let mut durable = request.clone();
|
|
if let Some(working_directory) = durable.working_directory_request.as_mut()
|
|
&& let Some(materialization) = working_directory.materialization.as_mut()
|
|
{
|
|
materialization.ssh = None;
|
|
}
|
|
durable
|
|
}
|
|
|
|
fn requested_primary_workdir_id(request: &CreateWorkerRequest) -> Option<&str> {
|
|
request
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|claim| claim.working_directory_id.as_str())
|
|
.or_else(|| {
|
|
request
|
|
.working_directory_request
|
|
.as_ref()
|
|
.and_then(|request| request.backend_workdir_id.as_deref())
|
|
})
|
|
}
|
|
|
|
fn validate_create_worker_request(request: &CreateWorkerRequest) -> Result<(), RuntimeError> {
|
|
if request.create_fingerprint.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"create_fingerprint must not be empty".to_string(),
|
|
));
|
|
}
|
|
match &request.profile_source {
|
|
crate::catalog::ProfileSourceArchiveSource::Embedded { archive } => {
|
|
archive.verify().map_err(|err| {
|
|
RuntimeError::InvalidRequest(format!("profile_source archive is invalid: {err}"))
|
|
})?;
|
|
}
|
|
crate::catalog::ProfileSourceArchiveSource::WorkspaceConfig { archive } => {
|
|
if archive.digest.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"profile_source.archive.digest must not be empty".to_string(),
|
|
));
|
|
}
|
|
}
|
|
}
|
|
if let Some(input) = &request.initial_input {
|
|
if input.kind != WorkerInputKind::User {
|
|
return Err(RuntimeError::InvalidInitialInputKind {
|
|
kind: format!("{:?}", input.kind),
|
|
});
|
|
}
|
|
let has_segments = input
|
|
.segments
|
|
.as_ref()
|
|
.is_some_and(|segments| !segments.is_empty());
|
|
if input.content.trim().is_empty() && !has_segments {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"initial_input.content must not be empty".to_string(),
|
|
));
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn worker_execution_failure_code(result: &WorkerExecutionResult) -> &'static str {
|
|
let message = result.message.as_deref().unwrap_or_default();
|
|
if message.contains("timed out waiting for durable Worker Submit acceptance") {
|
|
"input_commit_timeout"
|
|
} else if message.contains("worker adapter task did not complete within") {
|
|
"adapter_task_timeout"
|
|
} else if message.contains("failed to send Worker method") {
|
|
"worker_method_send_failed"
|
|
} else if message.contains("worker rejected Submit") {
|
|
"worker_submit_rejected"
|
|
} else if message.contains("before durable acceptance") {
|
|
"worker_failed_before_input_commit"
|
|
} else if message.contains("event stream closed") {
|
|
"worker_event_stream_closed"
|
|
} else {
|
|
"worker_execution_failed"
|
|
}
|
|
}
|
|
|
|
fn runtime_worker_create_failure_fields(
|
|
error: &RuntimeError,
|
|
) -> (&'static str, Option<String>, Option<String>) {
|
|
match error {
|
|
RuntimeError::RuntimeStopped => ("runtime_stopped", None, None),
|
|
RuntimeError::InvalidInitialInputKind { .. } => ("invalid_initial_input_kind", None, None),
|
|
RuntimeError::WorkerNotFound { .. } => ("worker_not_found", None, None),
|
|
RuntimeError::WorkerExecutionUnavailable { .. } => {
|
|
("worker_execution_unavailable", None, None)
|
|
}
|
|
RuntimeError::ExecutionBackendUnavailable { .. } => {
|
|
("execution_backend_unavailable", None, None)
|
|
}
|
|
RuntimeError::WorkerExecutionRejected {
|
|
operation, outcome, ..
|
|
} => (
|
|
"worker_execution_rejected",
|
|
Some(format!("{operation:?}")),
|
|
Some(format!("{outcome:?}")),
|
|
),
|
|
RuntimeError::LimitTooLarge { .. } => ("limit_too_large", None, None),
|
|
RuntimeError::InvalidRequest(_) => ("invalid_request", None, None),
|
|
RuntimeError::WorkspaceOwnerMismatch { .. } => ("workspace_owner_mismatch", None, None),
|
|
RuntimeError::WorkingDirectory(_) => ("working_directory", None, None),
|
|
RuntimeError::ConfigBundleMissing { .. } => ("config_bundle_missing", None, None),
|
|
RuntimeError::ConfigBundleDigestMismatch { .. } => {
|
|
("config_bundle_digest_mismatch", None, None)
|
|
}
|
|
RuntimeError::InvalidProfileSelector { .. } => ("invalid_profile_selector", None, None),
|
|
RuntimeError::UnsupportedConfigDeclaration { .. } => {
|
|
("unsupported_config_declaration", None, None)
|
|
}
|
|
RuntimeError::StoreIo { .. } => ("store_io", None, None),
|
|
RuntimeError::StoreMissing { .. } => ("store_missing", None, None),
|
|
RuntimeError::StoreCorrupt { .. } => ("store_corrupt", None, None),
|
|
RuntimeError::StatePoisoned => ("state_poisoned", None, None),
|
|
}
|
|
}
|
|
|
|
fn write_runtime_worker_create_failure(
|
|
worker_id: WorkerId,
|
|
workspace_id: Option<&str>,
|
|
error: &RuntimeError,
|
|
elapsed: std::time::Duration,
|
|
) {
|
|
let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(error);
|
|
let execution_failure_code = match error {
|
|
RuntimeError::WorkerExecutionRejected { result, .. } => {
|
|
worker_execution_failure_code(result)
|
|
}
|
|
_ => "",
|
|
};
|
|
tracing::error!(
|
|
target: "yoi::worker_create",
|
|
event = "worker_create_failed",
|
|
component = "runtime",
|
|
workspace_id = workspace_id.unwrap_or(""),
|
|
worker_id = %worker_id,
|
|
error_kind,
|
|
operation = operation.as_deref().unwrap_or(""),
|
|
outcome = outcome.as_deref().unwrap_or(""),
|
|
execution_failure_code,
|
|
duration_ms = u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX),
|
|
"Worker creation failed"
|
|
);
|
|
}
|
|
|
|
fn validate_create_workspace_scope(
|
|
request: &CreateWorkerRequest,
|
|
workspace_id: Option<&str>,
|
|
) -> Result<(), RuntimeError> {
|
|
let Some(workspace_id) = workspace_id else {
|
|
return Ok(());
|
|
};
|
|
if workspace_id.trim().is_empty() {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Runtime auth workspace_id must not be empty".to_string(),
|
|
));
|
|
}
|
|
if let Some(request_workspace_id) = request
|
|
.workspace_api
|
|
.as_ref()
|
|
.map(|workspace_api| workspace_api.workspace_id.as_str())
|
|
{
|
|
if request_workspace_id != workspace_id {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"request workspace_id {request_workspace_id} does not match Runtime auth workspace_id {workspace_id}"
|
|
)));
|
|
}
|
|
}
|
|
let snapshot = request.memory_settings.as_ref().ok_or_else(|| {
|
|
RuntimeError::InvalidRequest(
|
|
"Workspace-scoped Worker create requires a bound Memory settings snapshot".to_string(),
|
|
)
|
|
})?;
|
|
if snapshot.workspace_id != workspace_id {
|
|
return Err(RuntimeError::InvalidRequest(format!(
|
|
"Memory settings workspace_id {} does not match Runtime auth workspace_id {workspace_id}",
|
|
snapshot.workspace_id
|
|
)));
|
|
}
|
|
if snapshot.settings_revision == 0 {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Memory settings revision must be at least 1".to_string(),
|
|
));
|
|
}
|
|
if !manifest::is_normalized_workspace_memory_language(&snapshot.language) {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"Memory settings language must be a normalized bounded UTF-8 value".to_string(),
|
|
));
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn validate_worker_input(input: &WorkerInput) -> Result<(), RuntimeError> {
|
|
let has_segments = input
|
|
.segments
|
|
.as_ref()
|
|
.is_some_and(|segments| !segments.is_empty());
|
|
if !input.kind.is_empty_content_allowed() && input.content.trim().is_empty() && !has_segments {
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"worker input content must not be empty".to_string(),
|
|
));
|
|
}
|
|
if input
|
|
.segments
|
|
.as_ref()
|
|
.is_some_and(|segments| segments.is_empty())
|
|
{
|
|
return Err(RuntimeError::InvalidRequest(
|
|
"worker input segments must not be empty".to_string(),
|
|
));
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn input_protocol_event(input: &WorkerInput) -> Option<protocol::Event> {
|
|
match input.kind {
|
|
// Submit is projected only after the Worker commits UserInput. Queued
|
|
// payloads must never become model- or client-visible history early.
|
|
WorkerInputKind::User | WorkerInputKind::Notify => None,
|
|
WorkerInputKind::Compact
|
|
| WorkerInputKind::ListRewindTargets
|
|
| WorkerInputKind::RegisterPeer => Some(protocol::Event::SystemItem {
|
|
item: serde_json::json!({
|
|
"kind": "embedded_worker_command_input",
|
|
"command": input.kind,
|
|
"content": input.content.clone(),
|
|
}),
|
|
}),
|
|
}
|
|
}
|
|
|
|
fn subscription_validation_error(error: SubscriptionValidationError) -> RuntimeError {
|
|
RuntimeError::InvalidRequest(format!("invalid event subscription: {error}"))
|
|
}
|
|
|
|
fn subscription_worker_state(status: WorkerStatus) -> SubscriptionWorkerState {
|
|
match status {
|
|
WorkerStatus::Idle => SubscriptionWorkerState::Idle,
|
|
WorkerStatus::Running => SubscriptionWorkerState::Running,
|
|
WorkerStatus::Paused => SubscriptionWorkerState::Paused,
|
|
WorkerStatus::Stopped => SubscriptionWorkerState::Stopped,
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::catalog::{
|
|
ConfigBundleRef, MaterializerKind, ProfileSelector, RepositoryMaterializationContext,
|
|
RepositorySshCredentialCandidate, RepositorySshMaterializationAccess, SensitiveString,
|
|
WorkingDirectoryClaim, WorkingDirectoryRepository, WorkingDirectoryRequest,
|
|
WorkspaceApiRef,
|
|
};
|
|
use crate::config_bundle::{
|
|
ConfigBundle, ConfigBundleMetadata, ConfigBundleProvenance, ConfigDeclaration,
|
|
ConfigDeclarationKind, ConfigProfileDescriptor,
|
|
};
|
|
use crate::execution::{
|
|
WorkerExecutionBackend, WorkerExecutionContext, WorkerExecutionHandle,
|
|
WorkerExecutionRestoreRequest,
|
|
};
|
|
use crate::working_directory::WorkingDirectoryDiagnostic;
|
|
use async_trait::async_trait;
|
|
use std::collections::BTreeMap;
|
|
#[cfg(feature = "fs-store")]
|
|
use std::sync::atomic::{AtomicU64, Ordering};
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
#[test]
|
|
fn worker_create_failure_fields_exclude_raw_error_messages() {
|
|
let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(
|
|
&RuntimeError::InvalidRequest("private-token".to_string()),
|
|
);
|
|
assert_eq!(error_kind, "invalid_request");
|
|
assert_eq!(operation, None);
|
|
assert_eq!(outcome, None);
|
|
let timeout = WorkerExecutionResult::errored(
|
|
WorkerExecutionOperation::Input,
|
|
"timed out waiting for durable Worker Submit acceptance; private-token",
|
|
);
|
|
assert_eq!(
|
|
worker_execution_failure_code(&timeout),
|
|
"input_commit_timeout"
|
|
);
|
|
}
|
|
|
|
fn test_command() -> protocol::WorkerCommandEnvelope {
|
|
protocol::WorkerCommandEnvelope {
|
|
command_id: 1,
|
|
expected_execution_generation: 1,
|
|
expected_worker_state_revision: 0,
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn repository_resource_failures_keep_typed_credential_diagnostics() {
|
|
let cases = [
|
|
(
|
|
BackendResourceError::Expired,
|
|
"repository_access_credential_expired",
|
|
),
|
|
(
|
|
BackendResourceError::MissingResource,
|
|
"repository_access_credential_unavailable",
|
|
),
|
|
(
|
|
BackendResourceError::Timeout,
|
|
"repository_access_resource_fetch_timeout",
|
|
),
|
|
(
|
|
BackendResourceError::Unauthorized {
|
|
message: "denied".to_string(),
|
|
},
|
|
"repository_access_credential_unauthorized",
|
|
),
|
|
];
|
|
for (error, expected_code) in cases {
|
|
let RuntimeError::WorkingDirectory(diagnostic) = repository_resource_error(error)
|
|
else {
|
|
panic!("Repository resource failure lost its typed diagnostic")
|
|
};
|
|
assert_eq!(diagnostic.code, expected_code);
|
|
}
|
|
}
|
|
|
|
fn internal_worker_ref(
|
|
session_id: &str,
|
|
parent_session_id: Option<&str>,
|
|
) -> protocol::InternalWorkerRef {
|
|
protocol::InternalWorkerRef {
|
|
session_id: session_id.to_string(),
|
|
parent_session_id: parent_session_id.map(str::to_string),
|
|
name: session_id.to_string(),
|
|
kind: protocol::InternalWorkerKind::SubWorker,
|
|
}
|
|
}
|
|
|
|
fn internal_worker_status_event(
|
|
worker: protocol::InternalWorkerRef,
|
|
status: protocol::WorkerStatus,
|
|
) -> protocol::Event {
|
|
protocol::Event::InternalWorker {
|
|
worker,
|
|
revision: 1,
|
|
event: Box::new(protocol::Event::WorkerState {
|
|
snapshot: status.into(),
|
|
}),
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn internal_worker_activity_tracks_running_children_independently() {
|
|
let mut activity = BTreeMap::new();
|
|
assert!(RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&internal_worker_status_event(
|
|
internal_worker_ref("child-a", None),
|
|
protocol::WorkerStatus::Running,
|
|
),
|
|
));
|
|
assert!(!RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&internal_worker_status_event(
|
|
internal_worker_ref("child-b", None),
|
|
protocol::WorkerStatus::Running,
|
|
),
|
|
));
|
|
assert!(!RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&internal_worker_status_event(
|
|
internal_worker_ref("child-a", None),
|
|
protocol::WorkerStatus::Idle,
|
|
),
|
|
));
|
|
assert!(RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&internal_worker_status_event(
|
|
internal_worker_ref("child-b", None),
|
|
protocol::WorkerStatus::Stopped,
|
|
),
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn nested_internal_worker_activity_reaches_the_parent_projection() {
|
|
let mut activity = BTreeMap::new();
|
|
let direct_child = internal_worker_ref("child", None);
|
|
let nested_running = protocol::Event::InternalWorker {
|
|
worker: direct_child.clone(),
|
|
revision: 1,
|
|
event: Box::new(internal_worker_status_event(
|
|
internal_worker_ref("grandchild", Some("child")),
|
|
protocol::WorkerStatus::Running,
|
|
)),
|
|
};
|
|
assert!(RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&nested_running,
|
|
));
|
|
|
|
let nested_idle = protocol::Event::InternalWorker {
|
|
worker: direct_child,
|
|
revision: 2,
|
|
event: Box::new(internal_worker_status_event(
|
|
internal_worker_ref("grandchild", Some("child")),
|
|
protocol::WorkerStatus::Idle,
|
|
)),
|
|
};
|
|
assert!(RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&nested_idle,
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn parent_snapshot_replaces_stale_internal_worker_activity() {
|
|
let mut activity = BTreeMap::new();
|
|
RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&internal_worker_status_event(
|
|
internal_worker_ref("child-a", None),
|
|
protocol::WorkerStatus::Running,
|
|
),
|
|
);
|
|
let snapshot = protocol::Event::Snapshot {
|
|
session: protocol::SessionSnapshot {
|
|
pending_submissions: protocol::PendingSubmissionsSnapshot::default(),
|
|
entries: Vec::new(),
|
|
},
|
|
greeting: protocol::Greeting {
|
|
worker_name: "parent".to_string(),
|
|
cwd: "/tmp".to_string(),
|
|
provider: "test".to_string(),
|
|
model: "test".to_string(),
|
|
scope_summary: String::new(),
|
|
tools: Vec::new(),
|
|
context_window: 0,
|
|
context_tokens: 0,
|
|
},
|
|
state: protocol::WorkerStatus::Idle.into(),
|
|
in_flight: protocol::InFlightSnapshot::default(),
|
|
internal_workers: Vec::new(),
|
|
};
|
|
assert!(RuntimeState::update_internal_worker_activity(
|
|
&mut activity,
|
|
&snapshot,
|
|
));
|
|
assert!(activity.is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn runtime_identity_binding_is_immutable_and_host_owned() {
|
|
let runtime = Runtime::new_memory();
|
|
runtime.bind_runtime_identity("runtime-a").unwrap();
|
|
runtime.bind_runtime_identity("runtime-a").unwrap();
|
|
assert!(runtime.bind_runtime_identity("runtime-b").is_err());
|
|
assert_eq!(
|
|
runtime.lock().unwrap().runtime_identity.as_deref(),
|
|
Some("runtime-a")
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn typed_segments_allow_empty_flat_content() {
|
|
let input = WorkerInput {
|
|
kind: WorkerInputKind::User,
|
|
content: String::new(),
|
|
submission_request_id: None,
|
|
segments: Some(vec![protocol::Segment::Flow {
|
|
selector: "builtin:coder-review".to_string(),
|
|
}]),
|
|
};
|
|
assert!(validate_worker_input(&input).is_ok());
|
|
}
|
|
|
|
#[test]
|
|
fn empty_user_input_without_segments_is_rejected() {
|
|
let input = WorkerInput {
|
|
kind: WorkerInputKind::User,
|
|
content: String::new(),
|
|
submission_request_id: None,
|
|
segments: Some(Vec::new()),
|
|
};
|
|
assert!(matches!(
|
|
validate_worker_input(&input),
|
|
Err(RuntimeError::InvalidRequest(_))
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn typed_flow_segments_allow_empty_initial_flat_content() {
|
|
let mut request = task_request("flow");
|
|
request.initial_input = Some(WorkerInput {
|
|
kind: WorkerInputKind::User,
|
|
content: String::new(),
|
|
submission_request_id: None,
|
|
segments: Some(vec![protocol::Segment::Flow {
|
|
selector: "builtin:coder-review".to_string(),
|
|
}]),
|
|
});
|
|
assert!(validate_create_worker_request(&request).is_ok());
|
|
}
|
|
|
|
fn test_profile_source_archive() -> crate::profile_archive::ProfileSourceArchive {
|
|
crate::profile_archive::ProfileSourceArchive::build(
|
|
crate::profile_archive::ProfileSourceArchiveInput {
|
|
id: "test-profile-source".to_string(),
|
|
entrypoints: BTreeMap::from([(
|
|
"builtin:coder".to_string(),
|
|
"profiles/coder.dcdl".to_string(),
|
|
)]),
|
|
imports: BTreeMap::new(),
|
|
sources: BTreeMap::from([("profiles/coder.dcdl".to_string(), "{}".to_string())]),
|
|
},
|
|
)
|
|
.unwrap()
|
|
}
|
|
|
|
fn task_request(_objective: &str) -> CreateWorkerRequest {
|
|
let profile = ProfileSelector::Builtin("builtin:coder".to_string());
|
|
let bundle = test_bundle_for_profile(profile.clone());
|
|
CreateWorkerRequest {
|
|
worker_id: WorkerId::now_v7(),
|
|
create_fingerprint: "test-create".to_string(),
|
|
profile,
|
|
display_name: None,
|
|
profile_source: crate::catalog::ProfileSourceArchiveSource::Embedded {
|
|
archive: test_profile_source_archive(),
|
|
},
|
|
config_bundle: Some(ConfigBundleRef {
|
|
id: bundle.metadata.id,
|
|
digest: bundle.metadata.digest,
|
|
}),
|
|
initial_input: None,
|
|
working_directory_request: None,
|
|
working_directory: None,
|
|
worker_observation_enabled: false,
|
|
worker_observation_grants: Vec::new(),
|
|
workspace_api: None,
|
|
memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot {
|
|
workspace_id: "local".to_string(),
|
|
settings_revision: 1,
|
|
language: "English".to_string(),
|
|
}),
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn durable_worker_request_omits_repository_credentials() {
|
|
let mut request = task_request("worker-secret-redaction");
|
|
request.working_directory_request = Some(WorkingDirectoryRequest {
|
|
repository: WorkingDirectoryRepository {
|
|
id: "repository-1".to_string(),
|
|
provider: "git".to_string(),
|
|
source: workspace_api::RepositorySource {
|
|
kind: workspace_api::RepositorySourceKind::Ssh,
|
|
uri: "ssh://git@example.test/repo.git".to_string(),
|
|
},
|
|
source_revision: 1,
|
|
source_fingerprint: "sha256:source".to_string(),
|
|
selector: None,
|
|
},
|
|
materializer: MaterializerKind::RuntimeGitClone,
|
|
backend_workdir_id: Some("working-directory-1".to_string()),
|
|
materialization: Some(RepositoryMaterializationContext {
|
|
workspace_id: "workspace-1".to_string(),
|
|
runtime_id: "runtime-1".to_string(),
|
|
operation_id: "operation-1".to_string(),
|
|
config_revision: 1,
|
|
config_projection_digest: "sha256:projection".to_string(),
|
|
ssh: Some(RepositorySshMaterializationAccess {
|
|
credential_candidates: vec![RepositorySshCredentialCandidate {
|
|
credential_id: "credential-1".to_string(),
|
|
credential_revision: 1,
|
|
private_key: SensitiveString::new("private-key-bytes"),
|
|
}],
|
|
host_trust_id: "host-trust-1".to_string(),
|
|
host_trust_revision: 1,
|
|
access: workspace_api::RepositoryAccessMode::ReadOnly,
|
|
expires_at_epoch_seconds: u64::MAX,
|
|
repository_id: "repository-1".to_string(),
|
|
repository_source_fingerprint: "sha256:source".to_string(),
|
|
repository_uri: "ssh://git@example.test/repo.git".to_string(),
|
|
secret_resource: repository_resource_handle(),
|
|
known_hosts_entry: SensitiveString::new("known-hosts-entry"),
|
|
}),
|
|
}),
|
|
});
|
|
|
|
let durable = durable_create_worker_request(&request);
|
|
|
|
assert!(
|
|
request
|
|
.working_directory_request
|
|
.as_ref()
|
|
.and_then(|working_directory| working_directory.materialization.as_ref())
|
|
.and_then(|materialization| materialization.ssh.as_ref())
|
|
.is_some()
|
|
);
|
|
assert!(
|
|
durable
|
|
.working_directory_request
|
|
.as_ref()
|
|
.and_then(|working_directory| working_directory.materialization.as_ref())
|
|
.and_then(|materialization| materialization.ssh.as_ref())
|
|
.is_none()
|
|
);
|
|
let serialized = serde_json::to_string(&durable).unwrap();
|
|
assert!(!serialized.contains("private-key-bytes"));
|
|
assert!(!serialized.contains("known-hosts-entry"));
|
|
}
|
|
|
|
fn repository_resource_handle() -> crate::resource::BackendResourceHandle {
|
|
crate::resource::BackendResourceHandle {
|
|
kind: crate::resource::BackendResourceKind::RepositorySshAccess,
|
|
workspace_id: "workspace-1".to_string(),
|
|
scope_id: Some("repository-ssh-access".to_string()),
|
|
runtime_id: Some("runtime-1".to_string()),
|
|
worker_id: None,
|
|
resource_id: "repository-access-1".to_string(),
|
|
digest: "opaque:repository-access-1".to_string(),
|
|
operation: crate::resource::BackendResourceOperation::FetchOnce,
|
|
expires_at_unix_seconds: i64::MAX,
|
|
nonce: "repository-access-1".to_string(),
|
|
revision: "1".to_string(),
|
|
generation: None,
|
|
max_bytes: crate::resource::DEFAULT_REPOSITORY_SSH_ACCESS_MAX_BYTES,
|
|
content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE.to_string(),
|
|
redaction: crate::resource::ResourceRedactionPolicy::RuntimeInternalOnly,
|
|
audit_correlation_id: "repository-access-1".to_string(),
|
|
profile_source_graph: None,
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn repository_access_resource_is_fetched_before_provider_authorization() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend
|
|
.repository_access_available
|
|
.store(true, Ordering::SeqCst);
|
|
runtime.bind_runtime_identity("runtime-1").unwrap();
|
|
let handle = repository_resource_handle();
|
|
runtime
|
|
.install_workspace_backend_resource_client(
|
|
"workspace-1",
|
|
Arc::new(TestRepositoryResourceClient {
|
|
response: Mutex::new(Some(crate::resource::BackendResourceFetchResponse {
|
|
kind: crate::resource::BackendResourceKind::RepositorySshAccess,
|
|
resource_id: handle.resource_id.clone(),
|
|
digest: handle.digest.clone(),
|
|
content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE
|
|
.to_string(),
|
|
bytes: serde_json::to_vec(&RepositorySshAccessSecret {
|
|
credential_candidates: vec![
|
|
crate::resource::RepositorySshAccessSecretCandidate {
|
|
credential_id: "credential-1".to_string(),
|
|
credential_revision: 1,
|
|
private_key: "private-key-bytes-1".to_string(),
|
|
},
|
|
crate::resource::RepositorySshAccessSecretCandidate {
|
|
credential_id: "credential-2".to_string(),
|
|
credential_revision: 3,
|
|
private_key: "private-key-bytes-2".to_string(),
|
|
},
|
|
],
|
|
known_hosts_entry: "known-hosts-entry".to_string(),
|
|
})
|
|
.unwrap(),
|
|
audit_correlation_id: handle.audit_correlation_id.clone(),
|
|
})),
|
|
}),
|
|
)
|
|
.unwrap();
|
|
let request = WorkingDirectoryRepositoryAccessRequest {
|
|
working_directory_id: "working-directory-1".to_string(),
|
|
materialization: RepositoryMaterializationContext {
|
|
workspace_id: "workspace-1".to_string(),
|
|
runtime_id: "runtime-1".to_string(),
|
|
operation_id: "operation-1".to_string(),
|
|
config_revision: 1,
|
|
config_projection_digest: "sha256:projection".to_string(),
|
|
ssh: Some(RepositorySshMaterializationAccess {
|
|
credential_candidates: vec![
|
|
RepositorySshCredentialCandidate {
|
|
credential_id: "credential-1".to_string(),
|
|
credential_revision: 1,
|
|
private_key: SensitiveString::default(),
|
|
},
|
|
RepositorySshCredentialCandidate {
|
|
credential_id: "credential-2".to_string(),
|
|
credential_revision: 3,
|
|
private_key: SensitiveString::default(),
|
|
},
|
|
],
|
|
host_trust_id: "host-trust-1".to_string(),
|
|
host_trust_revision: 1,
|
|
access: workspace_api::RepositoryAccessMode::ReadOnly,
|
|
expires_at_epoch_seconds: u64::MAX,
|
|
repository_id: "repository-1".to_string(),
|
|
repository_source_fingerprint: "sha256:source".to_string(),
|
|
repository_uri: "ssh://git@example.test/repo.git".to_string(),
|
|
secret_resource: handle,
|
|
known_hosts_entry: SensitiveString::default(),
|
|
}),
|
|
},
|
|
};
|
|
|
|
let replay = request.clone();
|
|
runtime
|
|
.authorize_working_directory_repository_access_from_resource(request)
|
|
.await
|
|
.unwrap();
|
|
assert!(
|
|
runtime
|
|
.authorize_working_directory_repository_access_from_resource(replay)
|
|
.await
|
|
.is_err()
|
|
);
|
|
|
|
let accesses = backend.repository_accesses.lock().unwrap();
|
|
assert_eq!(accesses.len(), 1);
|
|
let access = accesses[0].materialization.ssh.as_ref().unwrap();
|
|
assert_eq!(access.credential_candidates.len(), 2);
|
|
assert_eq!(
|
|
access.credential_candidates[0].credential_id,
|
|
"credential-1"
|
|
);
|
|
assert_eq!(
|
|
access.credential_candidates[1].credential_id,
|
|
"credential-2"
|
|
);
|
|
assert_eq!(
|
|
access.credential_candidates[0].private_key.expose(),
|
|
"private-key-bytes-1"
|
|
);
|
|
assert_eq!(
|
|
access.credential_candidates[1].private_key.expose(),
|
|
"private-key-bytes-2"
|
|
);
|
|
assert_eq!(access.known_hosts_entry.expose(), "known-hosts-entry");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn working_directory_create_fetches_repository_access_before_provider_call() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend
|
|
.repository_access_available
|
|
.store(true, Ordering::SeqCst);
|
|
runtime.bind_runtime_identity("runtime-1").unwrap();
|
|
let handle = repository_resource_handle();
|
|
runtime
|
|
.install_workspace_backend_resource_client(
|
|
"workspace-1",
|
|
Arc::new(TestRepositoryResourceClient {
|
|
response: Mutex::new(Some(crate::resource::BackendResourceFetchResponse {
|
|
kind: crate::resource::BackendResourceKind::RepositorySshAccess,
|
|
resource_id: handle.resource_id.clone(),
|
|
digest: handle.digest.clone(),
|
|
content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE
|
|
.to_string(),
|
|
bytes: serde_json::to_vec(&RepositorySshAccessSecret {
|
|
credential_candidates: vec![
|
|
crate::resource::RepositorySshAccessSecretCandidate {
|
|
credential_id: "credential-1".to_string(),
|
|
credential_revision: 1,
|
|
private_key: "create-private-key-bytes-1".to_string(),
|
|
},
|
|
crate::resource::RepositorySshAccessSecretCandidate {
|
|
credential_id: "credential-2".to_string(),
|
|
credential_revision: 3,
|
|
private_key: "create-private-key-bytes-2".to_string(),
|
|
},
|
|
],
|
|
known_hosts_entry: "create-known-hosts-entry".to_string(),
|
|
})
|
|
.unwrap(),
|
|
audit_correlation_id: handle.audit_correlation_id.clone(),
|
|
})),
|
|
}),
|
|
)
|
|
.unwrap();
|
|
let request = WorkingDirectoryRequest {
|
|
repository: WorkingDirectoryRepository {
|
|
id: "repository-1".to_string(),
|
|
provider: "git".to_string(),
|
|
source: workspace_api::RepositorySource {
|
|
kind: workspace_api::RepositorySourceKind::Ssh,
|
|
uri: "ssh://git@example.test/repo.git".to_string(),
|
|
},
|
|
source_revision: 1,
|
|
source_fingerprint: "sha256:source".to_string(),
|
|
selector: None,
|
|
},
|
|
materializer: MaterializerKind::RuntimeGitClone,
|
|
backend_workdir_id: Some("working-directory-1".to_string()),
|
|
materialization: Some(RepositoryMaterializationContext {
|
|
workspace_id: "workspace-1".to_string(),
|
|
runtime_id: "runtime-1".to_string(),
|
|
operation_id: "operation-create".to_string(),
|
|
config_revision: 1,
|
|
config_projection_digest: "sha256:projection".to_string(),
|
|
ssh: Some(RepositorySshMaterializationAccess {
|
|
credential_candidates: vec![
|
|
RepositorySshCredentialCandidate {
|
|
credential_id: "credential-1".to_string(),
|
|
credential_revision: 1,
|
|
private_key: SensitiveString::default(),
|
|
},
|
|
RepositorySshCredentialCandidate {
|
|
credential_id: "credential-2".to_string(),
|
|
credential_revision: 3,
|
|
private_key: SensitiveString::default(),
|
|
},
|
|
],
|
|
host_trust_id: "host-trust-1".to_string(),
|
|
host_trust_revision: 1,
|
|
access: workspace_api::RepositoryAccessMode::ReadOnly,
|
|
expires_at_epoch_seconds: u64::MAX,
|
|
repository_id: "repository-1".to_string(),
|
|
repository_source_fingerprint: "sha256:source".to_string(),
|
|
repository_uri: "ssh://git@example.test/repo.git".to_string(),
|
|
secret_resource: handle,
|
|
known_hosts_entry: SensitiveString::default(),
|
|
}),
|
|
}),
|
|
};
|
|
|
|
assert!(
|
|
runtime
|
|
.create_working_directory_from_resource(request)
|
|
.await
|
|
.is_err()
|
|
);
|
|
|
|
let requests = backend.working_directory_requests.lock().unwrap();
|
|
let access = requests[0]
|
|
.materialization
|
|
.as_ref()
|
|
.and_then(|materialization| materialization.ssh.as_ref())
|
|
.unwrap();
|
|
assert_eq!(
|
|
access.credential_candidates[0].private_key.expose(),
|
|
"create-private-key-bytes-1"
|
|
);
|
|
assert_eq!(
|
|
access.credential_candidates[1].private_key.expose(),
|
|
"create-private-key-bytes-2"
|
|
);
|
|
assert_eq!(
|
|
access.known_hosts_entry.expose(),
|
|
"create-known-hosts-entry"
|
|
);
|
|
}
|
|
|
|
fn scoped_task_request(objective: &str, workspace_id: &str) -> CreateWorkerRequest {
|
|
let mut request = task_request(objective);
|
|
request.workspace_api = Some(WorkspaceApiRef {
|
|
workspace_id: workspace_id.to_string(),
|
|
base_url: format!("https://workspace.example/{workspace_id}"),
|
|
});
|
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
|
workspace_id: workspace_id.to_string(),
|
|
settings_revision: 1,
|
|
language: "English".to_string(),
|
|
});
|
|
request
|
|
}
|
|
|
|
#[test]
|
|
fn workspace_create_requires_matching_normalized_memory_settings_snapshot() {
|
|
let mut request = scoped_task_request("memory-snapshot", "workspace-a");
|
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok());
|
|
request.memory_settings.as_mut().unwrap().language = "Français".to_string();
|
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok());
|
|
|
|
request.memory_settings = None;
|
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
|
|
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
|
workspace_id: "workspace-b".to_string(),
|
|
settings_revision: 1,
|
|
language: "English".to_string(),
|
|
});
|
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
|
|
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
|
workspace_id: "workspace-a".to_string(),
|
|
settings_revision: 2,
|
|
language: " english ".to_string(),
|
|
});
|
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
|
}
|
|
|
|
fn scope(workspace_id: &str, server_id: &str) -> RuntimeWorkspaceScope {
|
|
RuntimeWorkspaceScope::new(workspace_id, server_id)
|
|
}
|
|
|
|
fn test_bundle_for_profile(profile: ProfileSelector) -> ConfigBundle {
|
|
ConfigBundle {
|
|
metadata: ConfigBundleMetadata {
|
|
id: "bundle-1".to_string(),
|
|
digest: String::new(),
|
|
revision: "rev-1".to_string(),
|
|
workspace_id: "workspace-1".to_string(),
|
|
created_at: "2026-06-26T00:00:00Z".to_string(),
|
|
provenance: ConfigBundleProvenance {
|
|
source: "workspace-backend".to_string(),
|
|
detail: Some("profile-sync".to_string()),
|
|
},
|
|
},
|
|
profiles: vec![ConfigProfileDescriptor {
|
|
selector: profile,
|
|
label: Some("Coder".to_string()),
|
|
}],
|
|
declarations: vec![ConfigDeclaration {
|
|
kind: ConfigDeclarationKind::CapabilityGrant,
|
|
name: "read".to_string(),
|
|
reference: "capability:read".to_string(),
|
|
}],
|
|
prompt_catalog: None,
|
|
profile_source_archive: None,
|
|
profile_source_archive_handle: None,
|
|
}
|
|
.with_computed_digest()
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct TestExecutionBackend {
|
|
dispatch_result: Mutex<Option<WorkerExecutionResult>>,
|
|
stop_result: Mutex<Option<WorkerExecutionResult>>,
|
|
restore_result: Mutex<Option<WorkerExecutionSpawnResult>>,
|
|
restore_count: Mutex<u64>,
|
|
run_generations: Mutex<Vec<u64>>,
|
|
config_bundles: Mutex<Vec<Option<ConfigBundle>>>,
|
|
workspace_config_fetches: Mutex<Vec<WorkspaceConfigFetchRequest>>,
|
|
workspace_config_results: Mutex<Vec<WorkspaceConfigFetchResult>>,
|
|
contexts: Mutex<BTreeMap<WorkerId, WorkerExecutionContext>>,
|
|
dispatched_inputs: Mutex<Vec<WorkerInput>>,
|
|
repository_accesses: Mutex<Vec<WorkingDirectoryRepositoryAccessRequest>>,
|
|
repository_access_available: AtomicBool,
|
|
working_directory_requests: Mutex<Vec<WorkingDirectoryRequest>>,
|
|
preserve_submission_acknowledgement_id: AtomicBool,
|
|
#[cfg(feature = "ws-server")]
|
|
snapshots: Mutex<BTreeMap<WorkerId, protocol::Event>>,
|
|
}
|
|
|
|
impl TestExecutionBackend {
|
|
fn set_dispatch_result(&self, result: WorkerExecutionResult) {
|
|
*self.dispatch_result.lock().unwrap() = Some(result);
|
|
}
|
|
|
|
fn set_stop_result(&self, result: WorkerExecutionResult) {
|
|
*self.stop_result.lock().unwrap() = Some(result);
|
|
}
|
|
|
|
fn preserve_submission_acknowledgement_id(&self) {
|
|
self.preserve_submission_acknowledgement_id
|
|
.store(true, Ordering::SeqCst);
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn set_worker_snapshot(&self, worker_ref: &WorkerRef, snapshot: protocol::Event) {
|
|
self.snapshots
|
|
.lock()
|
|
.unwrap()
|
|
.insert(worker_ref.worker_id.clone(), snapshot);
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn publish_text_delta(
|
|
&self,
|
|
worker_ref: &WorkerRef,
|
|
text: &str,
|
|
) -> Result<crate::observation::WorkerObservationEvent, RuntimeError> {
|
|
let contexts = self.contexts.lock().unwrap();
|
|
let context = contexts.get(&worker_ref.worker_id).expect("context stored");
|
|
context.publish_protocol_event(protocol::Event::TextDelta { text: text.into() })
|
|
}
|
|
}
|
|
|
|
impl WorkerExecutionBackend for TestExecutionBackend {
|
|
fn backend_id(&self) -> &str {
|
|
"test-execution-backend"
|
|
}
|
|
|
|
fn fetch_workspace_config(
|
|
&self,
|
|
request: WorkspaceConfigFetchRequest,
|
|
) -> Result<WorkspaceConfigFetchResult, String> {
|
|
self.workspace_config_fetches.lock().unwrap().push(request);
|
|
let mut results = self.workspace_config_results.lock().unwrap();
|
|
if results.is_empty() {
|
|
return Err("no Workspace Config fetch result configured".to_string());
|
|
}
|
|
Ok(results.remove(0))
|
|
}
|
|
|
|
fn create_working_directory(
|
|
&self,
|
|
request: &WorkingDirectoryRequest,
|
|
) -> Result<CatalogWorkingDirectoryStatus, WorkingDirectoryDiagnostic> {
|
|
self.working_directory_requests
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.clone());
|
|
Err(WorkingDirectoryDiagnostic::rejected(
|
|
"working_directory_unsupported",
|
|
"Worker execution backend does not support working directory materialization",
|
|
))
|
|
}
|
|
|
|
fn authorize_working_directory_repository_access(
|
|
&self,
|
|
request: &WorkingDirectoryRepositoryAccessRequest,
|
|
) -> Result<(), WorkingDirectoryDiagnostic> {
|
|
self.repository_accesses
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.clone());
|
|
if self.repository_access_available.load(Ordering::SeqCst) {
|
|
Ok(())
|
|
} else {
|
|
Err(WorkingDirectoryDiagnostic::rejected(
|
|
"working_directory_repository_access_unsupported",
|
|
"Worker execution backend does not support Repository access authorization",
|
|
))
|
|
}
|
|
}
|
|
|
|
fn spawn_worker(&self, request: WorkerExecutionSpawnRequest) -> WorkerExecutionSpawnResult {
|
|
self.run_generations
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.run_generation);
|
|
self.config_bundles
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.config_bundle.clone());
|
|
self.contexts
|
|
.lock()
|
|
.unwrap()
|
|
.insert(request.worker_ref.worker_id.clone(), request.context);
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()),
|
|
worker_state: protocol::WorkerStateSnapshot {
|
|
execution_generation: request.run_generation,
|
|
..protocol::WorkerStatus::Idle.into()
|
|
},
|
|
working_directory: request
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|binding| binding.status()),
|
|
}
|
|
}
|
|
|
|
fn restore_worker(
|
|
&self,
|
|
request: WorkerExecutionRestoreRequest,
|
|
) -> WorkerExecutionSpawnResult {
|
|
*self.restore_count.lock().unwrap() += 1;
|
|
self.run_generations
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.run_generation);
|
|
self.config_bundles
|
|
.lock()
|
|
.unwrap()
|
|
.push(request.config_bundle.clone());
|
|
if let Some(result) = self.restore_result.lock().unwrap().clone() {
|
|
return result;
|
|
}
|
|
self.contexts
|
|
.lock()
|
|
.unwrap()
|
|
.insert(request.worker_ref.worker_id.clone(), request.context);
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()),
|
|
worker_state: protocol::WorkerStateSnapshot {
|
|
execution_generation: request.run_generation,
|
|
..protocol::WorkerStatus::Idle.into()
|
|
},
|
|
working_directory: request
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|binding| binding.status()),
|
|
}
|
|
}
|
|
|
|
fn dispatch_input(
|
|
&self,
|
|
_handle: &WorkerExecutionHandle,
|
|
input: WorkerInput,
|
|
) -> WorkerExecutionResult {
|
|
let submission_id = input.submission_request_id.clone();
|
|
self.dispatched_inputs.lock().unwrap().push(input);
|
|
let mut result = self
|
|
.dispatch_result
|
|
.lock()
|
|
.unwrap()
|
|
.clone()
|
|
.unwrap_or_else(|| {
|
|
WorkerExecutionResult::accepted_submission(
|
|
WorkerExecutionOperation::Input,
|
|
"request-test",
|
|
"test-submission",
|
|
protocol::SubmissionDisposition::Started,
|
|
)
|
|
});
|
|
if !self
|
|
.preserve_submission_acknowledgement_id
|
|
.load(Ordering::SeqCst)
|
|
&& let (Some(ack), Some(submission_id)) =
|
|
(result.submission.as_mut(), submission_id)
|
|
{
|
|
ack.submission_request_id = submission_id;
|
|
}
|
|
result
|
|
}
|
|
|
|
fn stop_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult {
|
|
self.stop_result
|
|
.lock()
|
|
.unwrap()
|
|
.take()
|
|
.unwrap_or_else(|| WorkerExecutionResult::accepted(WorkerExecutionOperation::Stop))
|
|
}
|
|
|
|
fn cancel_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult {
|
|
WorkerExecutionResult::accepted(WorkerExecutionOperation::Cancel)
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
fn worker_snapshot(&self, handle: &WorkerExecutionHandle) -> Option<protocol::Event> {
|
|
self.snapshots
|
|
.lock()
|
|
.unwrap()
|
|
.get(&handle.worker_ref().worker_id)
|
|
.cloned()
|
|
}
|
|
}
|
|
|
|
struct TestRepositoryResourceClient {
|
|
response: Mutex<Option<crate::resource::BackendResourceFetchResponse>>,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl BackendResourceClient for TestRepositoryResourceClient {
|
|
async fn fetch_resource(
|
|
&self,
|
|
_request: BackendResourceFetchRequest,
|
|
) -> Result<crate::resource::BackendResourceFetchResponse, BackendResourceError> {
|
|
self.response
|
|
.lock()
|
|
.unwrap()
|
|
.take()
|
|
.ok_or(BackendResourceError::MissingResource)
|
|
}
|
|
}
|
|
|
|
fn runtime_with_backend() -> Runtime {
|
|
let runtime = Runtime::with_execution_backend(
|
|
RuntimeOptions::default(),
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
runtime
|
|
}
|
|
|
|
fn runtime_and_backend() -> (Runtime, Arc<TestExecutionBackend>) {
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime =
|
|
Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
(runtime, backend)
|
|
}
|
|
|
|
fn test_bundle() -> ConfigBundle {
|
|
test_bundle_for_profile(ProfileSelector::Builtin("builtin:coder".to_string()))
|
|
}
|
|
|
|
fn bundled_task_request(objective: &str, bundle: &ConfigBundle) -> CreateWorkerRequest {
|
|
let mut request = task_request(objective);
|
|
request.config_bundle = Some(ConfigBundleRef {
|
|
id: bundle.metadata.id.clone(),
|
|
digest: bundle.metadata.digest.clone(),
|
|
});
|
|
request
|
|
}
|
|
|
|
fn receive_subscription_update(
|
|
subscription: &mut RuntimeEventSelectorSubscription,
|
|
) -> Result<RuntimeSubscriptionUpdate, RuntimeSubscriptionRecvError> {
|
|
tokio::runtime::Builder::new_current_thread()
|
|
.build()
|
|
.unwrap()
|
|
.block_on(subscription.recv())
|
|
}
|
|
|
|
#[test]
|
|
fn runtime_worker_subscription_has_gap_free_snapshot_and_live_updates() {
|
|
let runtime = runtime_with_backend();
|
|
let mut subscription = runtime
|
|
.subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers)
|
|
.unwrap();
|
|
assert_eq!(subscription.snapshot_revision(), 0);
|
|
let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else {
|
|
panic!("runtime_workers must return a Worker snapshot");
|
|
};
|
|
assert!(workers.is_empty());
|
|
|
|
let created = runtime.create_worker(task_request("live")).unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert_eq!(update.subject_revision, 1);
|
|
match update.payload {
|
|
SubscriptionEventPayload::WorkerUpserted { worker } => {
|
|
assert_eq!(worker.worker_id.as_str(), created.worker_id.to_string());
|
|
assert_eq!(worker.subject_revision, 1);
|
|
assert_eq!(worker.state, SubscriptionWorkerState::Idle);
|
|
}
|
|
payload => panic!("unexpected subscription payload: {payload:?}"),
|
|
}
|
|
|
|
runtime
|
|
.observe_worker_event(
|
|
&created.worker_ref,
|
|
internal_worker_status_event(
|
|
internal_worker_ref("child-live", None),
|
|
protocol::WorkerStatus::Running,
|
|
),
|
|
)
|
|
.unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert_eq!(update.subject_revision, 2);
|
|
match update.payload {
|
|
SubscriptionEventPayload::WorkerUpserted { worker } => {
|
|
assert_eq!(worker.state, SubscriptionWorkerState::Idle);
|
|
assert!(worker.has_running_internal_workers);
|
|
}
|
|
payload => panic!("unexpected subscription payload: {payload:?}"),
|
|
}
|
|
|
|
runtime
|
|
.observe_worker_event(
|
|
&created.worker_ref,
|
|
internal_worker_status_event(
|
|
internal_worker_ref("child-live", None),
|
|
protocol::WorkerStatus::Idle,
|
|
),
|
|
)
|
|
.unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert_eq!(update.subject_revision, 3);
|
|
match update.payload {
|
|
SubscriptionEventPayload::WorkerUpserted { worker } => {
|
|
assert_eq!(worker.state, SubscriptionWorkerState::Idle);
|
|
assert!(!worker.has_running_internal_workers);
|
|
}
|
|
payload => panic!("unexpected subscription payload: {payload:?}"),
|
|
}
|
|
|
|
runtime
|
|
.observe_worker_event(
|
|
&created.worker_ref,
|
|
internal_worker_status_event(
|
|
internal_worker_ref("child-live", None),
|
|
protocol::WorkerStatus::Running,
|
|
),
|
|
)
|
|
.unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert_eq!(update.subject_revision, 4);
|
|
match update.payload {
|
|
SubscriptionEventPayload::WorkerUpserted { worker } => {
|
|
assert_eq!(worker.state, SubscriptionWorkerState::Idle);
|
|
assert!(worker.has_running_internal_workers);
|
|
}
|
|
payload => panic!("unexpected subscription payload: {payload:?}"),
|
|
}
|
|
|
|
runtime.stop_worker(&created.worker_ref, None).unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert_eq!(update.subject_revision, 5);
|
|
match update.payload {
|
|
SubscriptionEventPayload::WorkerUpserted { worker } => {
|
|
assert_eq!(worker.worker_id.as_str(), created.worker_id.to_string());
|
|
assert_eq!(worker.state, SubscriptionWorkerState::Stopped);
|
|
assert!(!worker.has_running_internal_workers);
|
|
}
|
|
payload => panic!("unexpected subscription payload: {payload:?}"),
|
|
}
|
|
runtime.stop_runtime().unwrap();
|
|
}
|
|
|
|
#[test]
|
|
fn worker_lifecycle_subscription_delivers_only_selected_workers() {
|
|
let runtime = runtime_with_backend();
|
|
let first = runtime.create_worker(task_request("first")).unwrap();
|
|
let second = runtime.create_worker(task_request("second")).unwrap();
|
|
let selected_id = SubscriptionWorkerId::new(first.worker_id.to_string()).unwrap();
|
|
let mut subscription = runtime
|
|
.subscribe_event_selector(EventSubscriptionSelector::WorkerLifecycle {
|
|
worker_ids: protocol::subscription::SubscriptionWorkerIds::new([
|
|
selected_id.clone()
|
|
])
|
|
.unwrap(),
|
|
})
|
|
.unwrap();
|
|
let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else {
|
|
panic!("worker_lifecycle must return a Worker snapshot");
|
|
};
|
|
assert_eq!(workers.len(), 1);
|
|
assert_eq!(workers[0].worker_id, selected_id);
|
|
|
|
runtime.stop_worker(&second.worker_ref, None).unwrap();
|
|
assert!(matches!(
|
|
subscription.receiver.try_recv(),
|
|
Err(mpsc::error::TryRecvError::Empty)
|
|
));
|
|
|
|
runtime.stop_worker(&first.worker_ref, None).unwrap();
|
|
let update = receive_subscription_update(&mut subscription).unwrap();
|
|
assert!(matches!(
|
|
update.payload,
|
|
SubscriptionEventPayload::WorkerUpserted { worker }
|
|
if worker.worker_id == selected_id
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn scoped_runtime_worker_subscription_hides_other_workspaces() {
|
|
let runtime = runtime_with_backend();
|
|
let workspace_a = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
scoped_task_request("a", "workspace-a"),
|
|
)
|
|
.unwrap();
|
|
let workspace_b = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-b", "server-b"),
|
|
scoped_task_request("b", "workspace-b"),
|
|
)
|
|
.unwrap();
|
|
let mut subscription = runtime
|
|
.subscribe_event_selector_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
EventSubscriptionSelector::RuntimeWorkers,
|
|
)
|
|
.unwrap();
|
|
let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else {
|
|
panic!("runtime_workers must return a Worker snapshot");
|
|
};
|
|
assert_eq!(workers.len(), 1);
|
|
assert_eq!(
|
|
workers[0].worker_id.as_str(),
|
|
workspace_a.worker_id.to_string()
|
|
);
|
|
|
|
runtime.stop_worker(&workspace_b.worker_ref, None).unwrap();
|
|
assert!(matches!(
|
|
subscription.receiver.try_recv(),
|
|
Err(mpsc::error::TryRecvError::Empty)
|
|
));
|
|
runtime.stop_worker(&workspace_a.worker_ref, None).unwrap();
|
|
assert!(receive_subscription_update(&mut subscription).is_ok());
|
|
}
|
|
|
|
#[test]
|
|
fn lagged_runtime_subscription_closes_without_blocking_mutation() {
|
|
let runtime = runtime_with_backend();
|
|
let created = runtime.create_worker(task_request("lag")).unwrap();
|
|
let mut subscription = runtime
|
|
.subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers)
|
|
.unwrap();
|
|
{
|
|
let mut state = runtime.lock().unwrap();
|
|
for _ in 0..=SUBSCRIPTION_QUEUE_CAPACITY {
|
|
state.publish_worker_upsert(created.worker_id).unwrap();
|
|
}
|
|
}
|
|
assert_eq!(subscription.receiver.len(), SUBSCRIPTION_QUEUE_CAPACITY);
|
|
assert!(matches!(
|
|
receive_subscription_update(&mut subscription),
|
|
Err(RuntimeSubscriptionRecvError::Lagged)
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn dropping_runtime_subscription_releases_producer_state() {
|
|
let runtime = runtime_with_backend();
|
|
let subscription = runtime
|
|
.subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers)
|
|
.unwrap();
|
|
assert_eq!(runtime.lock().unwrap().subscriptions.len(), 1);
|
|
drop(subscription);
|
|
assert!(runtime.lock().unwrap().subscriptions.is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn scoped_worker_access_hides_other_workspace_workers() {
|
|
let runtime = runtime_with_backend();
|
|
let workspace_a = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
scoped_task_request("a", "workspace-a"),
|
|
)
|
|
.unwrap();
|
|
let workspace_b = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-b", "server-b"),
|
|
scoped_task_request("b", "workspace-b"),
|
|
)
|
|
.unwrap();
|
|
|
|
assert_eq!(workspace_a.workspace_id.as_deref(), Some("workspace-a"));
|
|
assert_eq!(workspace_b.workspace_id.as_deref(), Some("workspace-b"));
|
|
assert_eq!(
|
|
runtime
|
|
.list_workers_scoped(&scope("workspace-a", "server-a"))
|
|
.unwrap()
|
|
.into_iter()
|
|
.map(|worker| worker.worker_ref)
|
|
.collect::<Vec<_>>(),
|
|
vec![workspace_a.worker_ref.clone()]
|
|
);
|
|
assert_eq!(
|
|
runtime
|
|
.list_workers_scoped(&scope("workspace-b", "server-b"))
|
|
.unwrap()
|
|
.into_iter()
|
|
.map(|worker| worker.worker_ref)
|
|
.collect::<Vec<_>>(),
|
|
vec![workspace_b.worker_ref.clone()]
|
|
);
|
|
|
|
let detail_error = runtime
|
|
.worker_detail_scoped(&scope("workspace-a", "server-a"), &workspace_b.worker_ref)
|
|
.unwrap_err();
|
|
assert!(matches!(detail_error, RuntimeError::WorkerNotFound { .. }));
|
|
|
|
let input_error = runtime
|
|
.send_input_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
&workspace_b.worker_ref,
|
|
WorkerInput::user("cross workspace"),
|
|
)
|
|
.unwrap_err();
|
|
assert!(matches!(input_error, RuntimeError::WorkerNotFound { .. }));
|
|
|
|
let protocol_error = runtime
|
|
.send_protocol_method_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
&workspace_b.worker_ref,
|
|
Method::Shutdown {
|
|
command: test_command(),
|
|
},
|
|
)
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
protocol_error,
|
|
RuntimeError::WorkerNotFound { .. }
|
|
));
|
|
assert_eq!(
|
|
runtime
|
|
.worker_detail(&workspace_b.worker_ref)
|
|
.unwrap()
|
|
.status,
|
|
WorkerStatus::Idle
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn workspace_api_replacement_updates_persisted_request() {
|
|
let (runtime, _backend) = runtime_and_backend();
|
|
let scope = scope("workspace-a", "server-a");
|
|
let worker = runtime
|
|
.create_worker_scoped(
|
|
&scope,
|
|
scoped_task_request("repair credential", "workspace-a"),
|
|
)
|
|
.unwrap();
|
|
let replacement = WorkspaceApiRef {
|
|
workspace_id: "workspace-a".to_string(),
|
|
base_url: "https://workspace.example/workspace-a/".to_string(),
|
|
};
|
|
|
|
runtime
|
|
.replace_worker_workspace_api_scoped(&scope, &worker.worker_ref, replacement.clone())
|
|
.unwrap();
|
|
|
|
let state = runtime.lock().unwrap();
|
|
assert_eq!(
|
|
state
|
|
.worker(&worker.worker_ref)
|
|
.unwrap()
|
|
.request
|
|
.workspace_api,
|
|
Some(replacement)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn workspace_owner_binding_rejects_other_backend_and_forgets_after_last_worker_delete() {
|
|
let runtime = runtime_with_backend();
|
|
let server_a = scope("workspace-a", "server-a");
|
|
let server_b = scope("workspace-a", "server-b");
|
|
let first = runtime
|
|
.create_worker_scoped(&server_a, scoped_task_request("first", "workspace-a"))
|
|
.unwrap();
|
|
let second = runtime
|
|
.create_worker_scoped(&server_a, scoped_task_request("second", "workspace-a"))
|
|
.unwrap();
|
|
|
|
let create_error = runtime
|
|
.create_worker_scoped(&server_b, scoped_task_request("stolen", "workspace-a"))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
create_error,
|
|
RuntimeError::WorkspaceOwnerMismatch { .. }
|
|
));
|
|
let read_error = runtime
|
|
.worker_detail_scoped(&server_b, &first.worker_ref)
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
read_error,
|
|
RuntimeError::WorkspaceOwnerMismatch { .. }
|
|
));
|
|
|
|
runtime.stop_worker(&first.worker_ref, None).unwrap();
|
|
runtime.stop_worker(&second.worker_ref, None).unwrap();
|
|
runtime
|
|
.delete_worker_scoped(&server_a, &first.worker_ref)
|
|
.unwrap();
|
|
let still_owned_error = runtime
|
|
.create_worker_scoped(&server_b, scoped_task_request("still owned", "workspace-a"))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
still_owned_error,
|
|
RuntimeError::WorkspaceOwnerMismatch { .. }
|
|
));
|
|
|
|
runtime
|
|
.delete_worker_scoped(&server_a, &second.worker_ref)
|
|
.unwrap();
|
|
let rebound = runtime
|
|
.create_worker_scoped(&server_b, scoped_task_request("rebound", "workspace-a"))
|
|
.unwrap();
|
|
assert_eq!(rebound.workspace_id.as_deref(), Some("workspace-a"));
|
|
}
|
|
|
|
#[test]
|
|
fn scoped_create_rejects_request_workspace_mismatch() {
|
|
let runtime = runtime_with_backend();
|
|
let error = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
scoped_task_request("mismatch", "workspace-b"),
|
|
)
|
|
.unwrap_err();
|
|
assert!(matches!(error, RuntimeError::InvalidRequest(_)));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn unscoped_legacy_workers_are_hidden_from_scoped_access() {
|
|
let runtime = runtime_with_backend();
|
|
let legacy = runtime.create_worker(task_request("legacy")).unwrap();
|
|
|
|
assert!(legacy.workspace_id.is_none());
|
|
assert!(
|
|
runtime
|
|
.list_workers_scoped(&scope("workspace-a", "server-a"))
|
|
.unwrap()
|
|
.is_empty()
|
|
);
|
|
let detail_error = runtime
|
|
.worker_detail_scoped(&scope("workspace-a", "server-a"), &legacy.worker_ref)
|
|
.unwrap_err();
|
|
assert!(matches!(detail_error, RuntimeError::WorkerNotFound { .. }));
|
|
}
|
|
|
|
#[test]
|
|
fn stopped_worker_scoped_list_uses_workspace_boundary() {
|
|
let runtime = runtime_with_backend();
|
|
let workspace_a = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
scoped_task_request("a", "workspace-a"),
|
|
)
|
|
.unwrap();
|
|
let workspace_b = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-b", "server-b"),
|
|
scoped_task_request("b", "workspace-b"),
|
|
)
|
|
.unwrap();
|
|
runtime.stop_worker(&workspace_a.worker_ref, None).unwrap();
|
|
runtime.stop_worker(&workspace_b.worker_ref, None).unwrap();
|
|
|
|
let stopped = runtime
|
|
.list_stopped_workers_scoped(&scope("workspace-a", "server-a"))
|
|
.unwrap();
|
|
assert_eq!(stopped.len(), 1);
|
|
assert_eq!(stopped[0].worker_ref, workspace_a.worker_ref);
|
|
}
|
|
|
|
#[test]
|
|
fn create_list_and_detail_preserve_runtime_local_worker_authority() {
|
|
let runtime = runtime_with_backend();
|
|
let detail = runtime.create_worker(task_request("implement v0")).unwrap();
|
|
|
|
assert_eq!(detail.status, WorkerStatus::Idle);
|
|
assert_eq!(detail.config_bundle.as_ref().unwrap().id, "bundle-1");
|
|
|
|
let list = runtime.list_workers().unwrap();
|
|
assert_eq!(list.len(), 1);
|
|
assert_eq!(list[0].worker_ref, detail.worker_ref);
|
|
assert_eq!(list[0].config_bundle, detail.config_bundle);
|
|
|
|
let fetched = runtime.worker_detail(&detail.worker_ref).unwrap();
|
|
assert_eq!(fetched.worker_id, detail.worker_id);
|
|
assert_eq!(fetched.profile, detail.profile);
|
|
}
|
|
|
|
#[test]
|
|
fn stopped_worker_list_excludes_alive_and_cancelled_workers() {
|
|
let runtime = runtime_with_backend();
|
|
let alive = runtime.create_worker(task_request("alive")).unwrap();
|
|
let stopped = runtime.create_worker(task_request("stopped")).unwrap();
|
|
let cancelled = runtime.create_worker(task_request("cancelled")).unwrap();
|
|
|
|
runtime.stop_worker(&stopped.worker_ref, None).unwrap();
|
|
runtime.cancel_worker(&cancelled.worker_ref, None).unwrap();
|
|
|
|
let candidates = runtime.list_stopped_workers().unwrap();
|
|
assert_eq!(candidates.len(), 1);
|
|
assert_eq!(candidates[0].worker_ref, stopped.worker_ref);
|
|
assert_eq!(candidates[0].status, WorkerStatus::Stopped);
|
|
assert_ne!(candidates[0].worker_ref, alive.worker_ref);
|
|
assert_ne!(candidates[0].worker_ref, cancelled.worker_ref);
|
|
}
|
|
|
|
#[test]
|
|
fn synced_config_bundle_is_stored_checked_and_used_for_worker_creation() {
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime =
|
|
Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap();
|
|
let mut bundle = test_bundle();
|
|
bundle.prompt_catalog = Some(
|
|
worker::EffectivePromptCatalog::new(
|
|
BTreeMap::from([("default".to_string(), "workspace prompt".to_string())]),
|
|
7,
|
|
"schema",
|
|
"toolchain",
|
|
)
|
|
.unwrap(),
|
|
);
|
|
bundle = bundle.with_computed_digest();
|
|
let availability = runtime.store_config_bundle(bundle.clone()).unwrap();
|
|
assert_eq!(availability.reference.id, "bundle-1");
|
|
assert_eq!(availability.reference.digest, bundle.metadata.digest);
|
|
let mut conflicting_bundle = bundle.clone();
|
|
conflicting_bundle.profiles[0].label = Some("conflicting".to_string());
|
|
conflicting_bundle = conflicting_bundle.with_computed_digest();
|
|
assert!(matches!(
|
|
runtime.store_config_bundle(conflicting_bundle),
|
|
Err(RuntimeError::ConfigBundleDigestMismatch { .. })
|
|
));
|
|
|
|
let listed = runtime.list_config_bundles().unwrap();
|
|
assert_eq!(listed.len(), 1);
|
|
assert_eq!(listed[0].id, "bundle-1");
|
|
|
|
let checked = runtime
|
|
.check_config_bundle(&availability.reference)
|
|
.unwrap();
|
|
assert_eq!(checked.summary.digest, availability.summary.digest);
|
|
|
|
let detail = runtime
|
|
.create_worker(bundled_task_request("synced", &bundle))
|
|
.unwrap();
|
|
assert_eq!(detail.config_bundle, Some(availability.reference));
|
|
assert_eq!(
|
|
backend.config_bundles.lock().unwrap().as_slice(),
|
|
&[Some(bundle.clone())]
|
|
);
|
|
|
|
runtime.stop_worker(&detail.worker_ref, None).unwrap();
|
|
runtime.restore_worker(&detail.worker_ref).unwrap();
|
|
assert_eq!(
|
|
backend.config_bundles.lock().unwrap().as_slice(),
|
|
&[Some(bundle), None]
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn remote_create_refreshes_workspace_config_and_revalidates_cached_etag() {
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime =
|
|
Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap();
|
|
let mut bundle = test_bundle();
|
|
let archive = test_profile_source_archive();
|
|
bundle.profile_source_archive = Some(archive.clone());
|
|
bundle = bundle.with_computed_digest();
|
|
backend.workspace_config_results.lock().unwrap().extend([
|
|
WorkspaceConfigFetchResult::Modified(bundle.clone()),
|
|
WorkspaceConfigFetchResult::NotModified,
|
|
]);
|
|
let archive = archive.reference;
|
|
let request = |objective: &str| {
|
|
let mut request = bundled_task_request(objective, &bundle);
|
|
request.profile_source = ProfileSourceArchiveSource::WorkspaceConfig {
|
|
archive: archive.clone(),
|
|
};
|
|
request.workspace_api = Some(WorkspaceApiRef {
|
|
workspace_id: "workspace-test".to_string(),
|
|
base_url: "https://workspace.example".to_string(),
|
|
});
|
|
request
|
|
};
|
|
|
|
let first = request("first refresh");
|
|
let second = request("cached refresh");
|
|
runtime.create_worker(first).unwrap();
|
|
runtime.create_worker(second.clone()).unwrap();
|
|
runtime.create_worker(second).unwrap();
|
|
|
|
let fetches = backend.workspace_config_fetches.lock().unwrap();
|
|
assert_eq!(fetches.len(), 2);
|
|
assert_eq!(fetches[0].cached, None);
|
|
assert_eq!(
|
|
fetches[1].cached,
|
|
Some(ConfigBundleRef {
|
|
id: bundle.metadata.id.clone(),
|
|
digest: bundle.metadata.digest.clone(),
|
|
})
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn restore_does_not_require_recorded_config_bundle() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
let bundle = test_bundle();
|
|
let detail = runtime
|
|
.create_worker(bundled_task_request("missing-on-restore", &bundle))
|
|
.unwrap();
|
|
runtime.stop_worker(&detail.worker_ref, None).unwrap();
|
|
runtime.lock().unwrap().config_bundles.clear();
|
|
runtime.restore_worker(&detail.worker_ref).unwrap();
|
|
assert_eq!(
|
|
backend.config_bundles.lock().unwrap().as_slice(),
|
|
&[Some(bundle), None]
|
|
);
|
|
|
|
let (runtime, backend) = runtime_and_backend();
|
|
let bundle = test_bundle();
|
|
let detail = runtime
|
|
.create_worker(bundled_task_request("mismatch-on-restore", &bundle))
|
|
.unwrap();
|
|
runtime.stop_worker(&detail.worker_ref, None).unwrap();
|
|
let mut replacement = bundle.clone();
|
|
replacement.profiles[0].label = Some("replacement".to_string());
|
|
replacement = replacement.with_computed_digest();
|
|
runtime
|
|
.lock()
|
|
.unwrap()
|
|
.config_bundles
|
|
.insert(replacement.metadata.id.clone(), replacement);
|
|
runtime.restore_worker(&detail.worker_ref).unwrap();
|
|
assert_eq!(
|
|
backend.config_bundles.lock().unwrap().as_slice(),
|
|
&[Some(bundle), None]
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn config_bundle_errors_are_typed() {
|
|
let runtime = Runtime::new_memory();
|
|
let bundle = test_bundle();
|
|
|
|
runtime.store_config_bundle(bundle.clone()).unwrap();
|
|
let mismatch = runtime
|
|
.check_config_bundle(&ConfigBundleRef {
|
|
id: bundle.metadata.id.clone(),
|
|
digest: "0".repeat(64),
|
|
})
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
mismatch,
|
|
RuntimeError::ConfigBundleDigestMismatch { .. }
|
|
));
|
|
|
|
let mut unsupported = test_bundle();
|
|
unsupported.declarations.push(ConfigDeclaration {
|
|
kind: ConfigDeclarationKind::Unsupported,
|
|
name: "plugin-registry".to_string(),
|
|
reference: "plugin-registry:v0".to_string(),
|
|
});
|
|
unsupported = unsupported.with_computed_digest();
|
|
let unsupported_err = runtime.store_config_bundle(unsupported).unwrap_err();
|
|
assert!(matches!(
|
|
unsupported_err,
|
|
RuntimeError::UnsupportedConfigDeclaration { .. }
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_idempotency_reuses_worker_and_rejects_different_input() {
|
|
let runtime = runtime_with_backend();
|
|
let mut request = task_request("idempotent");
|
|
request.create_fingerprint = "sha256:input-1".to_string();
|
|
request.working_directory = Some(WorkingDirectoryClaim {
|
|
working_directory_id: "workdir-idempotent".to_string(),
|
|
relative_cwd: None,
|
|
});
|
|
|
|
let first = runtime.create_worker(request.clone()).unwrap();
|
|
let workdir_count_after_first = runtime.list_working_directories().unwrap().len();
|
|
let replayed = runtime.create_worker(request.clone()).unwrap();
|
|
assert_eq!(replayed.worker_ref, first.worker_ref);
|
|
assert_eq!(runtime.list_workers().unwrap().len(), 1);
|
|
assert_eq!(
|
|
runtime.list_working_directories().unwrap().len(),
|
|
workdir_count_after_first
|
|
);
|
|
|
|
request.create_fingerprint = "sha256:different".to_string();
|
|
let error = runtime.create_worker(request).unwrap_err();
|
|
assert!(matches!(error, RuntimeError::InvalidRequest(_)));
|
|
assert_eq!(runtime.list_workers().unwrap().len(), 1);
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_rejects_system_initial_input_without_persisting_worker() {
|
|
let runtime = runtime_with_backend();
|
|
let mut request = task_request("system initial input");
|
|
request.initial_input = Some(WorkerInput::notify("role/system belongs in config bundle"));
|
|
|
|
let error = runtime.create_worker(request).unwrap_err();
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::InvalidInitialInputKind { .. }
|
|
));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_does_not_infer_state_from_started_submission_ack() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend.set_dispatch_result(WorkerExecutionResult::accepted_submission(
|
|
WorkerExecutionOperation::Input,
|
|
"request-test",
|
|
"test-submission",
|
|
protocol::SubmissionDisposition::Started,
|
|
));
|
|
let mut request = task_request("committed initial input is already idle");
|
|
request.initial_input = Some(WorkerInput::user("start the ticket"));
|
|
|
|
let detail = runtime.create_worker(request).unwrap();
|
|
|
|
assert_eq!(detail.status, WorkerStatus::Idle);
|
|
assert_eq!(
|
|
detail.worker_state.as_ref().map(|snapshot| &snapshot.state),
|
|
Some(&protocol::WorkerState::Idle)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn restored_worker_exposes_the_backend_initial_state_snapshot() {
|
|
let (runtime, _) = runtime_and_backend();
|
|
let created = runtime
|
|
.create_worker(task_request("restore initial state"))
|
|
.unwrap();
|
|
runtime.stop_worker(&created.worker_ref, None).unwrap();
|
|
|
|
let restored = runtime.restore_worker(&created.worker_ref).unwrap();
|
|
|
|
let worker_state = restored
|
|
.worker_state
|
|
.expect("restored Worker must expose its initial state");
|
|
assert_eq!(worker_state.execution_generation, 2);
|
|
assert_eq!(worker_state.state, protocol::WorkerState::Idle);
|
|
}
|
|
|
|
#[test]
|
|
fn runtime_applies_only_newer_worker_state_snapshots() {
|
|
let (runtime, _) = runtime_and_backend();
|
|
let detail = runtime
|
|
.create_worker(task_request("state ordering"))
|
|
.unwrap();
|
|
let running = protocol::WorkerStateSnapshot {
|
|
execution_generation: 7,
|
|
revision: 3,
|
|
state: protocol::WorkerState::Busy(protocol::WorkerBusyState::Run(
|
|
protocol::WorkerRunState::Running,
|
|
)),
|
|
last_command_id: 2,
|
|
};
|
|
assert!({
|
|
let mut state = runtime.lock().unwrap();
|
|
state.project_protocol_event_to_worker_state(
|
|
&detail.worker_ref,
|
|
&protocol::Event::WorkerState {
|
|
snapshot: running.clone(),
|
|
},
|
|
)
|
|
});
|
|
assert_eq!(
|
|
runtime
|
|
.worker_detail(&detail.worker_ref)
|
|
.unwrap()
|
|
.worker_state,
|
|
Some(running.clone())
|
|
);
|
|
|
|
assert!({
|
|
let mut state = runtime.lock().unwrap();
|
|
!state.project_protocol_event_to_worker_state(
|
|
&detail.worker_ref,
|
|
&protocol::Event::WorkerState {
|
|
snapshot: protocol::WorkerStateSnapshot {
|
|
revision: 2,
|
|
state: protocol::WorkerState::Idle,
|
|
..running.clone()
|
|
},
|
|
},
|
|
)
|
|
});
|
|
assert!({
|
|
let mut state = runtime.lock().unwrap();
|
|
!state.project_protocol_event_to_worker_state(
|
|
&detail.worker_ref,
|
|
&protocol::Event::WorkerState {
|
|
snapshot: protocol::WorkerStateSnapshot {
|
|
state: protocol::WorkerState::Idle,
|
|
..running.clone()
|
|
},
|
|
},
|
|
)
|
|
});
|
|
let after = runtime.worker_detail(&detail.worker_ref).unwrap();
|
|
assert_eq!(after.status, WorkerStatus::Idle);
|
|
assert_eq!(after.worker_state, Some(running));
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_rejects_mismatched_submission_acknowledgement() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend.preserve_submission_acknowledgement_id();
|
|
backend.set_dispatch_result(WorkerExecutionResult::accepted_submission(
|
|
WorkerExecutionOperation::Input,
|
|
"request-test",
|
|
"forged-submission",
|
|
protocol::SubmissionDisposition::Started,
|
|
));
|
|
let mut request = task_request("mismatched initial input commit ack");
|
|
request.initial_input = Some(WorkerInput::user("start the ticket"));
|
|
|
|
let error = runtime.create_worker(request).unwrap_err();
|
|
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::WorkerExecutionRejected {
|
|
outcome: crate::execution::WorkerExecutionOutcome::Rejected,
|
|
..
|
|
}
|
|
));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_rejects_initial_input_without_durable_submission_acknowledgement() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend.set_dispatch_result(WorkerExecutionResult::accepted(
|
|
WorkerExecutionOperation::Input,
|
|
));
|
|
let mut request = task_request("missing initial input commit ack");
|
|
request.initial_input = Some(WorkerInput::user("start the ticket"));
|
|
|
|
let error = runtime.create_worker(request).unwrap_err();
|
|
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::WorkerExecutionRejected {
|
|
outcome: crate::execution::WorkerExecutionOutcome::Rejected,
|
|
..
|
|
}
|
|
));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_without_execution_backend_is_rejected_and_not_persisted() {
|
|
let runtime = Runtime::new_memory();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let error = runtime
|
|
.create_worker(task_request("no backend"))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::ExecutionBackendUnavailable { .. }
|
|
));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn create_worker_without_execution_backend_is_rejected_before_persisting_worker() {
|
|
let runtime = Runtime::new_memory();
|
|
let error = runtime
|
|
.create_worker(task_request("missing backend"))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::ExecutionBackendUnavailable { .. }
|
|
));
|
|
assert!(runtime.list_workers().unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn connected_backend_busy_dispatch_is_typed_and_not_transcribed() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
backend.set_dispatch_result(WorkerExecutionResult::busy(
|
|
WorkerExecutionOperation::Input,
|
|
"worker is already running",
|
|
));
|
|
let detail = runtime.create_worker(task_request("busy")).unwrap();
|
|
|
|
let err = runtime
|
|
.send_input(&detail.worker_ref, WorkerInput::user("wait"))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
err,
|
|
RuntimeError::WorkerExecutionRejected {
|
|
outcome: crate::execution::WorkerExecutionOutcome::Busy,
|
|
..
|
|
}
|
|
));
|
|
let refreshed = runtime.worker_detail(&detail.worker_ref).unwrap();
|
|
assert_eq!(refreshed.status, WorkerStatus::Idle);
|
|
#[cfg(feature = "ws-server")]
|
|
assert!(
|
|
runtime
|
|
.read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap()
|
|
.is_empty()
|
|
);
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
#[test]
|
|
fn backend_protocol_publish_hook_writes_observation_bus() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
let detail = runtime.create_worker(task_request("observe")).unwrap();
|
|
|
|
let observation = backend
|
|
.publish_text_delta(&detail.worker_ref, "from backend")
|
|
.unwrap();
|
|
assert_eq!(observation.worker_ref, detail.worker_ref);
|
|
|
|
let observations = runtime
|
|
.read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap();
|
|
assert_eq!(observations.len(), 1);
|
|
assert!(matches!(
|
|
observations[0].payload,
|
|
protocol::Event::TextDelta { .. }
|
|
));
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
#[test]
|
|
fn observation_snapshot_prefers_live_backend_snapshot() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
let detail = runtime
|
|
.create_worker(task_request("observe snapshot"))
|
|
.unwrap();
|
|
let expected_entry = serde_json::json!({"kind": "restored-log-entry"});
|
|
backend.set_worker_snapshot(
|
|
&detail.worker_ref,
|
|
protocol::Event::Snapshot {
|
|
session: protocol::SessionSnapshot {
|
|
pending_submissions: protocol::PendingSubmissionsSnapshot::default(),
|
|
entries: vec![protocol::SessionSnapshotEntry {
|
|
entry_id: "restored-log-entry".to_owned(),
|
|
timestamp: 1,
|
|
provenance: protocol::SessionEntryProvenance::LegacyUnknown,
|
|
derived_from: Vec::new(),
|
|
data: protocol::SessionSnapshotEntryData::RunError {
|
|
message: expected_entry.to_string(),
|
|
},
|
|
}],
|
|
},
|
|
greeting: protocol::Greeting {
|
|
worker_name: "live-worker".to_string(),
|
|
cwd: "/tmp/live".to_string(),
|
|
provider: "test-provider".to_string(),
|
|
model: "test-model".to_string(),
|
|
scope_summary: "live snapshot".to_string(),
|
|
tools: Vec::new(),
|
|
context_window: 128,
|
|
context_tokens: 64,
|
|
},
|
|
state: protocol::WorkerStatus::Running.into(),
|
|
in_flight: protocol::InFlightSnapshot {
|
|
blocks: Vec::new(),
|
|
commands: Vec::new(),
|
|
},
|
|
internal_workers: Vec::new(),
|
|
},
|
|
);
|
|
|
|
let snapshot = runtime
|
|
.worker_observation_snapshot(&detail.worker_ref)
|
|
.unwrap();
|
|
match snapshot {
|
|
protocol::Event::Snapshot {
|
|
session,
|
|
greeting,
|
|
state,
|
|
..
|
|
} => {
|
|
assert_eq!(session.entries.len(), 1);
|
|
assert_eq!(session.entries[0].entry_id, "restored-log-entry");
|
|
assert_eq!(greeting.worker_name, "live-worker");
|
|
assert_eq!(state.catalog_status(), protocol::WorkerStatus::Running);
|
|
}
|
|
other => panic!("expected snapshot, got {other:?}"),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
#[test]
|
|
fn observation_snapshot_fails_closed_when_backend_snapshot_is_unavailable() {
|
|
let runtime = runtime_with_backend();
|
|
let detail = runtime
|
|
.create_worker(task_request("snapshot unavailable"))
|
|
.unwrap();
|
|
|
|
assert!(matches!(
|
|
runtime
|
|
.worker_observation_snapshot(&detail.worker_ref)
|
|
.unwrap_err(),
|
|
RuntimeError::WorkerExecutionUnavailable {
|
|
worker_id,
|
|
message,
|
|
} if worker_id == detail.worker_ref.worker_id
|
|
&& message == "authoritative Worker snapshot is unavailable"
|
|
));
|
|
}
|
|
|
|
struct InputOnlyBackend;
|
|
|
|
impl WorkerExecutionBackend for InputOnlyBackend {
|
|
fn backend_id(&self) -> &str {
|
|
"input-only"
|
|
}
|
|
|
|
fn spawn_worker(&self, request: WorkerExecutionSpawnRequest) -> WorkerExecutionSpawnResult {
|
|
WorkerExecutionSpawnResult::Connected {
|
|
handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()),
|
|
worker_state: protocol::WorkerStateSnapshot {
|
|
execution_generation: request.run_generation,
|
|
..protocol::WorkerStatus::Idle.into()
|
|
},
|
|
working_directory: request
|
|
.working_directory
|
|
.as_ref()
|
|
.map(|binding| binding.status()),
|
|
}
|
|
}
|
|
|
|
fn dispatch_input(
|
|
&self,
|
|
_handle: &WorkerExecutionHandle,
|
|
input: WorkerInput,
|
|
) -> WorkerExecutionResult {
|
|
WorkerExecutionResult::accepted_submission(
|
|
WorkerExecutionOperation::Input,
|
|
"request-test",
|
|
input
|
|
.submission_request_id
|
|
.expect("Runtime submission request id"),
|
|
protocol::SubmissionDisposition::Started,
|
|
)
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn connected_backend_stop_unsupported_is_typed_rejection() {
|
|
let runtime =
|
|
Runtime::with_execution_backend(RuntimeOptions::default(), Arc::new(InputOnlyBackend))
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let detail = runtime.create_worker(task_request("no stop")).unwrap();
|
|
|
|
let err = runtime
|
|
.stop_worker(&detail.worker_ref, Some("stop".to_string()))
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
err,
|
|
RuntimeError::WorkerExecutionRejected {
|
|
outcome: crate::execution::WorkerExecutionOutcome::Unsupported,
|
|
..
|
|
}
|
|
));
|
|
assert_eq!(
|
|
runtime.worker_detail(&detail.worker_ref).unwrap().status,
|
|
WorkerStatus::Idle
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn protocol_shutdown_transitions_worker_to_stopped() {
|
|
let runtime = runtime_with_backend();
|
|
let detail = runtime
|
|
.create_worker(task_request("shutdown from client"))
|
|
.unwrap();
|
|
|
|
runtime
|
|
.send_protocol_method(
|
|
&detail.worker_ref,
|
|
Method::Shutdown {
|
|
command: test_command(),
|
|
},
|
|
)
|
|
.unwrap();
|
|
|
|
assert_eq!(
|
|
runtime.worker_detail(&detail.worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn stopped_worker_rejects_input_until_explicitly_restored() {
|
|
let (runtime, backend) = runtime_and_backend();
|
|
let detail = runtime
|
|
.create_worker(task_request("restore explicitly"))
|
|
.unwrap();
|
|
runtime
|
|
.send_protocol_method(
|
|
&detail.worker_ref,
|
|
Method::Shutdown {
|
|
command: test_command(),
|
|
},
|
|
)
|
|
.unwrap();
|
|
|
|
assert!(matches!(
|
|
runtime.send_input(&detail.worker_ref, WorkerInput::user("do not wake")),
|
|
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
|
));
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
|
|
|
runtime.restore_worker(&detail.worker_ref).unwrap();
|
|
runtime
|
|
.send_input(&detail.worker_ref, WorkerInput::user("wake up"))
|
|
.unwrap();
|
|
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
|
|
assert_eq!(*backend.run_generations.lock().unwrap(), vec![1, 2]);
|
|
let restored = runtime.worker_detail(&detail.worker_ref).unwrap();
|
|
assert_eq!(restored.status, WorkerStatus::Idle);
|
|
assert_eq!(
|
|
restored
|
|
.worker_state
|
|
.as_ref()
|
|
.map(|snapshot| (snapshot.execution_generation, &snapshot.state)),
|
|
Some((2, &protocol::WorkerState::Idle))
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn restore_does_not_redispatch_spawn_initial_submit() {
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime = Runtime::with_execution_backend(
|
|
RuntimeOptions {
|
|
..RuntimeOptions::default()
|
|
},
|
|
backend.clone(),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let mut request = task_request("flow restore");
|
|
request.initial_input = Some(WorkerInput {
|
|
kind: WorkerInputKind::User,
|
|
content: String::new(),
|
|
submission_request_id: None,
|
|
segments: Some(vec![
|
|
protocol::Segment::Flow {
|
|
selector: "builtin:coder-review".to_string(),
|
|
},
|
|
protocol::Segment::text("Implement Ticket 00001"),
|
|
]),
|
|
});
|
|
let detail = runtime.create_worker(request).unwrap();
|
|
assert_eq!(backend.dispatched_inputs.lock().unwrap().len(), 1);
|
|
|
|
runtime
|
|
.stop_worker(&detail.worker_ref, Some("restore test".to_string()))
|
|
.unwrap();
|
|
runtime.restore_worker(&detail.worker_ref).unwrap();
|
|
|
|
assert_eq!(
|
|
backend.dispatched_inputs.lock().unwrap().len(),
|
|
1,
|
|
"restore must continue durable Worker state without replaying spawn initial input"
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn send_input_dispatches_segment_only_flow_submission() {
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime = Runtime::with_execution_backend(
|
|
RuntimeOptions {
|
|
..RuntimeOptions::default()
|
|
},
|
|
backend.clone(),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let detail = runtime.create_worker(task_request("flow segment")).unwrap();
|
|
let input = WorkerInput {
|
|
kind: WorkerInputKind::User,
|
|
content: String::new(),
|
|
submission_request_id: None,
|
|
segments: Some(vec![protocol::Segment::Flow {
|
|
selector: "builtin:coder-review".to_string(),
|
|
}]),
|
|
};
|
|
|
|
runtime
|
|
.send_input(&detail.worker_ref, input.clone())
|
|
.unwrap();
|
|
|
|
let dispatched = backend.dispatched_inputs.lock().unwrap();
|
|
assert_eq!(dispatched.len(), 1);
|
|
assert_eq!(dispatched[0].kind, input.kind);
|
|
assert_eq!(dispatched[0].content, input.content);
|
|
assert_eq!(dispatched[0].segments, input.segments);
|
|
let submission_request_id = dispatched[0]
|
|
.submission_request_id
|
|
.as_deref()
|
|
.expect("Runtime submission request id");
|
|
Uuid::parse_str(submission_request_id).expect("submission request id UUID");
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
#[test]
|
|
fn notify_input_does_not_duplicate_committed_notification_observation() {
|
|
let runtime = Runtime::with_execution_backend(
|
|
RuntimeOptions {
|
|
..RuntimeOptions::default()
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let detail = runtime.create_worker(task_request("chat")).unwrap();
|
|
|
|
runtime
|
|
.send_input(&detail.worker_ref, WorkerInput::user("hello"))
|
|
.unwrap();
|
|
runtime
|
|
.send_input(&detail.worker_ref, WorkerInput::notify("note"))
|
|
.unwrap();
|
|
|
|
let observations = runtime
|
|
.read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap();
|
|
assert!(observations.is_empty());
|
|
|
|
runtime
|
|
.observe_worker_event(
|
|
&detail.worker_ref,
|
|
protocol::Event::SystemItem {
|
|
item: serde_json::json!({
|
|
"kind": "notification",
|
|
"message": "note",
|
|
"body": "[Notification] note",
|
|
}),
|
|
},
|
|
)
|
|
.unwrap();
|
|
|
|
let observations = runtime
|
|
.read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap();
|
|
assert_eq!(observations.len(), 1);
|
|
let protocol::Event::SystemItem { item } = &observations[0].payload else {
|
|
panic!("committed notification observation must be a system item");
|
|
};
|
|
assert_eq!(item["kind"], "notification");
|
|
assert!(observations.iter().all(|observation| {
|
|
!matches!(
|
|
&observation.payload,
|
|
protocol::Event::SystemItem { item }
|
|
if item["kind"] == "embedded_worker_notification"
|
|
)
|
|
}));
|
|
}
|
|
|
|
#[test]
|
|
fn stop_and_cancel_workers_keep_four_state_summary() {
|
|
let runtime = runtime_with_backend();
|
|
let stopped = runtime.create_worker(task_request("stop me")).unwrap();
|
|
let cancelled = runtime.create_worker(task_request("cancel me")).unwrap();
|
|
runtime
|
|
.send_input(&cancelled.worker_ref, WorkerInput::user("start"))
|
|
.unwrap();
|
|
|
|
let stop_ack = runtime
|
|
.stop_worker(&stopped.worker_ref, Some("done".to_string()))
|
|
.unwrap();
|
|
assert_eq!(stop_ack.status, WorkerStatus::Stopped);
|
|
|
|
let cancel_ack = runtime
|
|
.cancel_worker(&cancelled.worker_ref, Some("abort".to_string()))
|
|
.unwrap();
|
|
assert_eq!(cancel_ack.status, WorkerStatus::Idle);
|
|
|
|
let summary = runtime.summary().unwrap();
|
|
assert_eq!(summary.worker_count, 2);
|
|
assert_eq!(summary.active_worker_count, 1);
|
|
assert_eq!(summary.stopped_worker_count, 1);
|
|
assert!(serde_json::from_value::<WorkerStatus>(serde_json::json!("cancelled")).is_err());
|
|
}
|
|
|
|
#[test]
|
|
fn delete_worker_removes_stopped_worker_from_runtime() {
|
|
let runtime = runtime_with_backend();
|
|
let worker = runtime.create_worker(task_request("delete me")).unwrap();
|
|
assert!(runtime.delete_worker(&worker.worker_ref).is_err());
|
|
runtime
|
|
.stop_worker(&worker.worker_ref, Some("done".to_string()))
|
|
.unwrap();
|
|
|
|
let result = runtime.delete_worker(&worker.worker_ref).unwrap();
|
|
assert!(result.deleted);
|
|
assert_eq!(result.worker_id, worker.worker_id);
|
|
assert!(matches!(
|
|
runtime.worker_detail(&worker.worker_ref),
|
|
Err(RuntimeError::WorkerNotFound { .. })
|
|
));
|
|
let summary = runtime.summary().unwrap();
|
|
assert_eq!(summary.worker_count, 0);
|
|
}
|
|
|
|
#[test]
|
|
fn delete_worker_waits_for_execution_shutdown_before_removing_snapshot() {
|
|
let root = tempfile::tempdir().unwrap();
|
|
let store_root = root.path().join("runtime");
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: store_root.clone(),
|
|
runtime_id: "delete-shutdown-barrier".to_string(),
|
|
display_name: None,
|
|
},
|
|
backend.clone(),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
backend.set_dispatch_result(WorkerExecutionResult::errored(
|
|
WorkerExecutionOperation::Input,
|
|
"initial input failed",
|
|
));
|
|
backend.set_stop_result(WorkerExecutionResult::errored(
|
|
WorkerExecutionOperation::Stop,
|
|
"shutdown is still pending",
|
|
));
|
|
let mut request = task_request("delete shutdown barrier");
|
|
request.initial_input = Some(WorkerInput::user("start"));
|
|
let worker_id = request.worker_id;
|
|
runtime.create_worker(request).unwrap_err();
|
|
let worker_ref = WorkerRef::new(worker_id);
|
|
backend.set_stop_result(WorkerExecutionResult::errored(
|
|
WorkerExecutionOperation::Stop,
|
|
"shutdown is still pending",
|
|
));
|
|
|
|
let error = runtime.delete_worker(&worker_ref).unwrap_err();
|
|
|
|
assert!(matches!(
|
|
error,
|
|
RuntimeError::WorkerExecutionRejected {
|
|
operation: WorkerExecutionOperation::Stop,
|
|
..
|
|
}
|
|
));
|
|
assert_eq!(
|
|
runtime.worker_detail(&worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
assert!(
|
|
store_root
|
|
.join("workers")
|
|
.join(worker_id.to_string())
|
|
.exists(),
|
|
"Worker aggregate must remain until backend shutdown is confirmed"
|
|
);
|
|
|
|
assert!(runtime.delete_worker(&worker_ref).unwrap().deleted);
|
|
assert!(
|
|
!store_root
|
|
.join("workers")
|
|
.join(worker_id.to_string())
|
|
.exists()
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn stop_then_cancel_preserves_stopped_terminal_state() {
|
|
let runtime = runtime_with_backend();
|
|
let worker = runtime
|
|
.create_worker(task_request("stable stopped"))
|
|
.unwrap();
|
|
|
|
let stop_ack = runtime
|
|
.stop_worker(&worker.worker_ref, Some("done".to_string()))
|
|
.unwrap();
|
|
let cancel_ack = runtime
|
|
.cancel_worker(&worker.worker_ref, Some("late cancel".to_string()))
|
|
.unwrap();
|
|
|
|
assert_eq!(stop_ack.status, WorkerStatus::Stopped);
|
|
assert_eq!(cancel_ack.status, WorkerStatus::Stopped);
|
|
assert_eq!(
|
|
runtime.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
|
|
let summary = runtime.summary().unwrap();
|
|
assert_eq!(summary.active_worker_count, 0);
|
|
assert_eq!(summary.stopped_worker_count, 1);
|
|
}
|
|
|
|
#[test]
|
|
fn cancel_then_stop_transitions_idle_session_to_stopped() {
|
|
let runtime = runtime_with_backend();
|
|
let worker = runtime
|
|
.create_worker(task_request("cancel then stop"))
|
|
.unwrap();
|
|
runtime
|
|
.send_input(&worker.worker_ref, WorkerInput::user("start"))
|
|
.unwrap();
|
|
|
|
let cancel_ack = runtime
|
|
.cancel_worker(&worker.worker_ref, Some("abort".to_string()))
|
|
.unwrap();
|
|
let stop_ack = runtime
|
|
.stop_worker(&worker.worker_ref, Some("late stop".to_string()))
|
|
.unwrap();
|
|
|
|
assert_eq!(cancel_ack.status, WorkerStatus::Idle);
|
|
assert_eq!(stop_ack.status, WorkerStatus::Stopped);
|
|
assert_eq!(
|
|
runtime.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
|
|
let summary = runtime.summary().unwrap();
|
|
assert_eq!(summary.active_worker_count, 0);
|
|
assert_eq!(summary.stopped_worker_count, 1);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
static NEXT_FS_TEST_ROOT: AtomicU64 = AtomicU64::new(1);
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn fs_store_root(label: &str) -> std::path::PathBuf {
|
|
let sequence = NEXT_FS_TEST_ROOT.fetch_add(1, Ordering::Relaxed);
|
|
let root = std::env::temp_dir().join(format!(
|
|
"worker-runtime-fs-store-{label}-{}-{sequence}",
|
|
std::process::id()
|
|
));
|
|
let _ = std::fs::remove_dir_all(&root);
|
|
root
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
fn runtime_store(runtime: &Runtime) -> FsRuntimeStore {
|
|
let state = runtime.lock().unwrap();
|
|
match &state.persistence {
|
|
RuntimePersistence::Fs(store) => store.clone(),
|
|
RuntimePersistence::Memory => panic!("expected fs-backed runtime"),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_rejects_schema_older_than_previous_release() {
|
|
let root = fs_store_root("unsupported-old-schema");
|
|
let runtime = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let runtime_path = root.join("runtime.json");
|
|
let mut runtime_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
|
runtime_json["schema_version"] = serde_json::json!(2);
|
|
std::fs::write(
|
|
&runtime_path,
|
|
serde_json::to_vec_pretty(&runtime_json).unwrap(),
|
|
)
|
|
.unwrap();
|
|
|
|
let error = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap_err();
|
|
assert!(
|
|
error
|
|
.to_string()
|
|
.contains("unsupported Runtime store schema version 2; expected 3, 4, 5, or 6")
|
|
);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_restores_workers_without_legacy_event_or_protocol_observation_logs() {
|
|
let root = fs_store_root("restore");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: Some("filesystem runtime".to_string()),
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(
|
|
runtime.summary().unwrap().backend,
|
|
RuntimeBackendKind::FsStore
|
|
);
|
|
let transport_bundle = test_bundle();
|
|
runtime
|
|
.store_config_bundle(transport_bundle.clone())
|
|
.unwrap();
|
|
|
|
let worker = runtime.create_worker(task_request("persist me")).unwrap();
|
|
runtime
|
|
.send_input(&worker.worker_ref, WorkerInput::user("first"))
|
|
.unwrap();
|
|
runtime
|
|
.send_input(&worker.worker_ref, WorkerInput::notify("second"))
|
|
.unwrap();
|
|
runtime
|
|
.stop_worker(&worker.worker_ref, Some("finished".to_string()))
|
|
.unwrap();
|
|
let worker_store_dir = root.join("workers").join(worker.worker_id.to_string());
|
|
let worker_snapshot: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(worker_store_dir.join("worker.json")).unwrap())
|
|
.unwrap();
|
|
assert_eq!(worker_snapshot["schema_version"], serde_json::json!(6));
|
|
assert_eq!(worker_snapshot["status"], serde_json::json!("stopped"));
|
|
assert_eq!(
|
|
worker_snapshot["execution"]["last_run_generation"],
|
|
serde_json::json!(1)
|
|
);
|
|
assert_eq!(
|
|
worker_snapshot["execution"]["binding"]["run_generation"],
|
|
serde_json::json!(1)
|
|
);
|
|
assert_eq!(
|
|
worker_snapshot["execution"]["restore_intent"],
|
|
serde_json::json!("explicit")
|
|
);
|
|
assert!(!root.join("events.jsonl").exists());
|
|
std::fs::write(
|
|
worker_store_dir.join("observations.jsonl"),
|
|
b"{\"legacy\":true}\n",
|
|
)
|
|
.unwrap();
|
|
std::fs::write(root.join("events.jsonl"), b"obsolete runtime event\n").unwrap();
|
|
drop(runtime);
|
|
|
|
let restored = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap();
|
|
let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap();
|
|
assert_eq!(restored_worker.status, WorkerStatus::Stopped);
|
|
assert!(matches!(
|
|
restored.check_config_bundle(&ConfigBundleRef {
|
|
id: transport_bundle.metadata.id.clone(),
|
|
digest: transport_bundle.metadata.digest.clone(),
|
|
}),
|
|
Err(RuntimeError::ConfigBundleMissing { .. })
|
|
));
|
|
assert!(!root.join("events.jsonl").exists());
|
|
assert!(!worker_store_dir.join("observations.jsonl").exists());
|
|
#[cfg(feature = "ws-server")]
|
|
{
|
|
let observations = restored
|
|
.read_worker_observation_events(&worker.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap();
|
|
assert!(observations.is_empty());
|
|
}
|
|
|
|
#[cfg(feature = "ws-server")]
|
|
{
|
|
let observation = restored
|
|
.observe_worker_event(
|
|
&worker.worker_ref,
|
|
protocol::Event::TextDelta {
|
|
text: "restored observation bus".to_string(),
|
|
},
|
|
)
|
|
.unwrap();
|
|
assert_eq!(observation.sequence, 1);
|
|
let observations = restored
|
|
.read_worker_observation_events(&worker.worker_ref, WorkerObservationCursor::zero())
|
|
.unwrap();
|
|
assert_eq!(observations.len(), 1);
|
|
assert_eq!(observations[0].cursor, observation.cursor);
|
|
}
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_restores_workspace_scope_and_hides_legacy_workers_from_scoped_access() {
|
|
let root = fs_store_root("workspace-scope");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let scoped = runtime
|
|
.create_worker_scoped(
|
|
&scope("workspace-a", "server-a"),
|
|
scoped_task_request("persist workspace", "workspace-a"),
|
|
)
|
|
.unwrap();
|
|
let legacy = runtime
|
|
.create_worker(task_request("legacy persist"))
|
|
.unwrap();
|
|
let recoverable_legacy = runtime
|
|
.create_worker(scoped_task_request(
|
|
"legacy recoverable persist",
|
|
"workspace-b",
|
|
))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let restored = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap();
|
|
let restored_scoped = restored.worker_detail(&scoped.worker_ref).unwrap();
|
|
assert_eq!(restored_scoped.workspace_id.as_deref(), Some("workspace-a"));
|
|
assert_eq!(
|
|
restored
|
|
.list_workers_scoped(&scope("workspace-a", "server-a"))
|
|
.unwrap()
|
|
.into_iter()
|
|
.map(|worker| worker.worker_ref)
|
|
.collect::<Vec<_>>(),
|
|
vec![scoped.worker_ref.clone()]
|
|
);
|
|
let legacy_error = restored
|
|
.worker_detail_scoped(&scope("workspace-a", "server-a"), &legacy.worker_ref)
|
|
.unwrap_err();
|
|
assert!(matches!(legacy_error, RuntimeError::WorkerNotFound { .. }));
|
|
let recovered_legacy = restored
|
|
.worker_detail_scoped(
|
|
&scope("workspace-b", "server-b"),
|
|
&recoverable_legacy.worker_ref,
|
|
)
|
|
.unwrap();
|
|
assert_eq!(
|
|
recovered_legacy.workspace_id.as_deref(),
|
|
Some("workspace-b")
|
|
);
|
|
let stolen_legacy_error = restored
|
|
.worker_detail_scoped(
|
|
&scope("workspace-b", "server-c"),
|
|
&recoverable_legacy.worker_ref,
|
|
)
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
stolen_legacy_error,
|
|
RuntimeError::WorkspaceOwnerMismatch { .. }
|
|
));
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_restores_active_worker_execution_handles() {
|
|
let root = fs_store_root("execution-restore");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("restore active worker"))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let backendless = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap();
|
|
let persisted_worker = backendless.worker_detail(&worker.worker_ref).unwrap();
|
|
assert_eq!(persisted_worker.status, WorkerStatus::Idle);
|
|
drop(backendless);
|
|
|
|
let restoring_backend = Arc::new(TestExecutionBackend::default());
|
|
let restored = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
restoring_backend.clone(),
|
|
)
|
|
.unwrap();
|
|
|
|
assert_eq!(*restoring_backend.restore_count.lock().unwrap(), 1);
|
|
let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap();
|
|
assert_eq!(restored_worker.status, WorkerStatus::Idle);
|
|
restored
|
|
.send_input(&worker.worker_ref, WorkerInput::user("after restart"))
|
|
.unwrap();
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_automatically_restores_every_active_lifecycle_state() {
|
|
for status in ["idle", "running", "paused"] {
|
|
let root = fs_store_root(&format!("automatic-{status}"));
|
|
let options = crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
};
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
options.clone(),
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request(&format!("restore {status}")))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let worker_path = root
|
|
.join("workers")
|
|
.join(worker.worker_id.to_string())
|
|
.join("worker.json");
|
|
let mut worker_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
|
worker_json["status"] = serde_json::json!(status);
|
|
std::fs::write(
|
|
&worker_path,
|
|
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
|
)
|
|
.unwrap();
|
|
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let restored =
|
|
Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap();
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 1, "status={status}");
|
|
assert_eq!(
|
|
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Idle,
|
|
"status={status}"
|
|
);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_current_schema_requires_lifecycle_authority() {
|
|
let root = fs_store_root("current-schema-requires-lifecycle");
|
|
let options = crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
};
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
options.clone(),
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("missing lifecycle authority"))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let worker_path = root
|
|
.join("workers")
|
|
.join(worker.worker_id.to_string())
|
|
.join("worker.json");
|
|
let mut worker_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
|
worker_json.as_object_mut().unwrap().remove("status");
|
|
std::fs::write(
|
|
&worker_path,
|
|
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
|
)
|
|
.unwrap();
|
|
|
|
let restored = Runtime::with_fs_store_and_execution_backend(
|
|
options,
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
assert!(restored.list_workers().unwrap().is_empty());
|
|
assert!(
|
|
restored
|
|
.diagnostics()
|
|
.unwrap()
|
|
.iter()
|
|
.any(|diagnostic| diagnostic.code == "worker_snapshot_ignored")
|
|
);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_shutdown_preserves_active_worker_restore_intent() {
|
|
let root = fs_store_root("shutdown-preserves-worker");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("preserve active worker"))
|
|
.unwrap();
|
|
|
|
runtime.stop_runtime().unwrap();
|
|
|
|
let snapshot: serde_json::Value = serde_json::from_slice(
|
|
&std::fs::read(
|
|
root.join("workers")
|
|
.join(worker.worker_id.to_string())
|
|
.join("worker.json"),
|
|
)
|
|
.unwrap(),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(snapshot["status"], serde_json::json!("idle"));
|
|
assert_eq!(
|
|
snapshot["execution"]["restore_intent"],
|
|
serde_json::json!("automatic")
|
|
);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_stopped_worker_requires_explicit_restore() {
|
|
let root = fs_store_root("stopped-explicit-restore");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("explicit restore only"))
|
|
.unwrap();
|
|
runtime
|
|
.stop_worker(&worker.worker_ref, Some("operator stop".to_string()))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let restored = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
backend.clone(),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
|
assert_eq!(
|
|
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
assert!(matches!(
|
|
restored.send_input(&worker.worker_ref, WorkerInput::user("implicit restore")),
|
|
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
|
));
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
|
|
|
restored.restore_worker(&worker.worker_ref).unwrap();
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
|
|
assert_eq!(
|
|
restored.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Idle
|
|
);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_migrates_schema_v3_workers_to_stopped_explicit_restore() {
|
|
let root = fs_store_root("schema-v3-restore-intent");
|
|
let options = crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
};
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
options.clone(),
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("schema v3 worker"))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let runtime_path = root.join("runtime.json");
|
|
let worker_path = root
|
|
.join("workers")
|
|
.join(worker.worker_id.to_string())
|
|
.join("worker.json");
|
|
let mut runtime_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap();
|
|
runtime_json["schema_version"] = serde_json::json!(3);
|
|
std::fs::write(
|
|
&runtime_path,
|
|
serde_json::to_vec_pretty(&runtime_json).unwrap(),
|
|
)
|
|
.unwrap();
|
|
let mut worker_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
|
worker_json["schema_version"] = serde_json::json!(3);
|
|
worker_json.as_object_mut().unwrap().remove("status");
|
|
worker_json.as_object_mut().unwrap().remove("execution");
|
|
worker_json["run_generation"] = serde_json::json!(7);
|
|
std::fs::write(
|
|
&worker_path,
|
|
serde_json::to_vec_pretty(&worker_json).unwrap(),
|
|
)
|
|
.unwrap();
|
|
|
|
let backend = Arc::new(TestExecutionBackend::default());
|
|
let migrated =
|
|
Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap();
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 0);
|
|
assert_eq!(
|
|
migrated.worker_detail(&worker.worker_ref).unwrap().status,
|
|
WorkerStatus::Stopped
|
|
);
|
|
let migrated_json: serde_json::Value =
|
|
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
|
|
assert_eq!(migrated_json["schema_version"], serde_json::json!(6));
|
|
assert_eq!(migrated_json["status"], serde_json::json!("stopped"));
|
|
assert_eq!(
|
|
migrated_json["execution"]["last_run_generation"],
|
|
serde_json::json!(7)
|
|
);
|
|
assert_eq!(
|
|
migrated_json["execution"]["binding"],
|
|
serde_json::Value::Null
|
|
);
|
|
assert_eq!(
|
|
migrated_json["execution"]["restore_intent"],
|
|
serde_json::json!("explicit")
|
|
);
|
|
assert!(matches!(
|
|
migrated.send_input(&worker.worker_ref, WorkerInput::notify("do not restore")),
|
|
Err(RuntimeError::WorkerExecutionUnavailable { .. })
|
|
));
|
|
|
|
migrated.restore_worker(&worker.worker_ref).unwrap();
|
|
assert_eq!(*backend.restore_count.lock().unwrap(), 1);
|
|
assert_eq!(backend.run_generations.lock().unwrap().as_slice(), &[8]);
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_stops_worker_and_reports_when_execution_restore_fails() {
|
|
let root = fs_store_root("execution-restore-failed");
|
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
|
let worker = runtime
|
|
.create_worker(task_request("restore failure"))
|
|
.unwrap();
|
|
drop(runtime);
|
|
|
|
let restoring_backend = Arc::new(TestExecutionBackend::default());
|
|
*restoring_backend.restore_result.lock().unwrap() =
|
|
Some(WorkerExecutionSpawnResult::Errored(
|
|
WorkerExecutionResult::errored(WorkerExecutionOperation::Restore, "restore boom"),
|
|
));
|
|
let restored = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
restoring_backend.clone(),
|
|
)
|
|
.unwrap();
|
|
|
|
assert_eq!(*restoring_backend.restore_count.lock().unwrap(), 1);
|
|
let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap();
|
|
assert_eq!(restored_worker.status, WorkerStatus::Stopped);
|
|
assert!(
|
|
restored
|
|
.diagnostics()
|
|
.unwrap()
|
|
.iter()
|
|
.any(
|
|
|diagnostic| diagnostic.code == "worker_execution_restore_failed"
|
|
&& diagnostic.worker_ref.as_ref() == Some(&worker.worker_ref)
|
|
)
|
|
);
|
|
let err = restored
|
|
.send_input(
|
|
&worker.worker_ref,
|
|
WorkerInput::user("after failed restore"),
|
|
)
|
|
.unwrap_err();
|
|
assert!(matches!(
|
|
err,
|
|
RuntimeError::WorkerExecutionUnavailable { .. }
|
|
));
|
|
|
|
let _ = std::fs::remove_dir_all(root);
|
|
}
|
|
|
|
#[cfg(feature = "fs-store")]
|
|
#[test]
|
|
fn fs_store_reports_corrupt_and_missing_data() {
|
|
let corrupt_root = fs_store_root("corrupt");
|
|
let corrupt_runtime = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: corrupt_root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap();
|
|
let corrupt_store = runtime_store(&corrupt_runtime);
|
|
std::fs::write(
|
|
corrupt_store.runtime_dir().join("runtime.json"),
|
|
b"not json",
|
|
)
|
|
.unwrap();
|
|
drop(corrupt_runtime);
|
|
let err = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: corrupt_root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.unwrap_err();
|
|
assert!(matches!(err, RuntimeError::StoreCorrupt { .. }));
|
|
let _ = std::fs::remove_dir_all(corrupt_root);
|
|
|
|
let missing_root = fs_store_root("missing");
|
|
let missing_runtime = Runtime::with_fs_store_and_execution_backend(
|
|
crate::fs_store::FsRuntimeStoreOptions {
|
|
root: missing_root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
},
|
|
Arc::new(TestExecutionBackend::default()),
|
|
)
|
|
.unwrap();
|
|
missing_runtime.store_config_bundle(test_bundle()).unwrap();
|
|
missing_runtime
|
|
.create_worker(task_request("missing worker snapshot"))
|
|
.unwrap();
|
|
let missing_store = runtime_store(&missing_runtime);
|
|
let mut worker_dirs = std::fs::read_dir(missing_store.runtime_dir().join("workers"))
|
|
.unwrap()
|
|
.collect::<Result<Vec<_>, _>>()
|
|
.unwrap();
|
|
worker_dirs.sort_by_key(|entry| entry.path());
|
|
std::fs::remove_file(worker_dirs[0].path().join("worker.json")).unwrap();
|
|
drop(missing_runtime);
|
|
let loaded = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions {
|
|
root: missing_root.clone(),
|
|
runtime_id: "test-runtime".to_string(),
|
|
display_name: None,
|
|
})
|
|
.expect("invalid worker snapshot should not make runtime store unreadable");
|
|
assert!(loaded.list_workers().unwrap().is_empty());
|
|
assert!(
|
|
loaded
|
|
.diagnostics()
|
|
.unwrap()
|
|
.iter()
|
|
.any(|diagnostic| diagnostic.code == "worker_snapshot_ignored")
|
|
);
|
|
let _ = std::fs::remove_dir_all(missing_root);
|
|
}
|
|
}
|