From c0532fda4e08cbba65b188d705a86fe92db9b44f Mon Sep 17 00:00:00 2001 From: Hare Date: Fri, 7 Aug 2026 16:31:30 +0900 Subject: [PATCH] runtime: project granted worker sessions --- crates/worker-runtime/src/catalog.rs | 12 + crates/worker-runtime/src/http_server.rs | 4 + crates/worker-runtime/src/runtime.rs | 2 + crates/worker-runtime/src/worker_backend.rs | 283 +++++++++++++++++- crates/workspace-server/src/hosts.rs | 18 ++ .../src/runtime_subscription_tests.rs | 2 + crates/workspace-server/src/server.rs | 218 ++++++++++++++ 7 files changed, 535 insertions(+), 4 deletions(-) diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index b299ee6e..3adbfb45 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -4,6 +4,10 @@ use crate::profile_archive::{ProfileSourceArchive, ProfileSourceArchiveRef}; use serde::{Deserialize, Serialize}; use std::path::PathBuf; +fn is_false(value: &bool) -> bool { + !*value +} + /// Profile selector boundary. This is a selector, not a resolved runtime config. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "kind", content = "value", rename_all = "snake_case")] @@ -217,6 +221,14 @@ pub struct CreateWorkerRequest { pub working_directory_request: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub working_directory: Option, + /// Backend-only feature enablement. Grants still define local Runtime peers; + /// the Workspace provider reauthorizes its dynamic set per operation. + #[serde(default, skip_serializing_if = "is_false")] + pub worker_observation_enabled: bool, + /// Backend-authored, bounded peer session grants. Runtime revalidates each + /// requested capture against this exact canonical `(runtime_id, worker_id)` set. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub worker_observation_grants: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] pub workspace_api: Option, } diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 4e925a47..ad586010 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -2061,6 +2061,8 @@ mod tests { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } @@ -2588,6 +2590,8 @@ mod ws_tests { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index d6f6f23b..d087ccca 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -2479,6 +2479,8 @@ mod tests { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index ae302a48..ed8946a0 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -32,18 +32,23 @@ use crate::working_directory::{ use async_trait::async_trait; use manifest::paths; use protocol::{Method, Segment, WorkerStatus}; -use session_store::FsStore; -use session_store::{CombinedStore, FsWorkerStore}; +use session_store::{CombinedStore, FsStore, FsWorkerStore, collect_state}; use tokio::runtime::Runtime; #[cfg(feature = "ws-server")] use tokio::sync::broadcast; use workdir::{LocalWorkdirSession, Workdir, WorkdirSessionCapabilities, WorkdirSessionHandle}; +use worker::feature::builtin::{ + CompositeWorkerObservationProvider, WorkerObservationError, WorkerObservationProvider, + WorkerObservationSubject, WorkerObservationSubjectRef, WorkerSessionCapture, + WorkspaceClientWorkerObservationProvider, +}; #[cfg(feature = "ws-server")] use worker::ipc::protocol_session::{live_log_entry_event, subscribe_worker_protocol_session}; use worker::{ - PromptLoader, RuntimeWorkspaceHttpClient, Worker, WorkerController, WorkerError, - WorkerFilesystemAuthority, WorkerHandle, WorkerWorkspaceContext, WorkspaceId, + PromptLoader, RuntimeWorkspaceHttpClient, SegmentLogSink, Worker, WorkerController, + WorkerError, WorkerFilesystemAuthority, WorkerHandle, WorkerSharedState, + WorkerWorkspaceContext, WorkspaceId, }; const DEFAULT_BACKEND_ID: &str = "worker-crate"; @@ -102,8 +107,123 @@ pub trait RuntimeWorkerFactory: Send + Sync + 'static { /// Production factory that resolves a normal Worker profile and spawns it under /// `WorkerController`. +#[derive(Default)] +struct RuntimeWorkerObservationHub { + workers: Mutex>, +} + +#[derive(Clone)] +struct RuntimeObservedWorker { + workspace_id: Option, + shared_state: std::sync::Weak, + sink: SegmentLogSink, +} + +impl RuntimeWorkerObservationHub { + fn register(&self, worker_ref: WorkerRef, workspace_id: Option, handle: &WorkerHandle) { + if let Ok(mut workers) = self.workers.lock() { + workers.insert( + worker_ref, + RuntimeObservedWorker { + workspace_id, + shared_state: Arc::downgrade(&handle.shared_state), + sink: handle.sink.clone(), + }, + ); + } + } + + fn get( + &self, + worker_ref: &WorkerRef, + ) -> Option<(Option, Arc, SegmentLogSink)> { + let mut workers = self.workers.lock().ok()?; + let entry = workers.get(worker_ref)?.clone(); + let Some(shared_state) = entry.shared_state.upgrade() else { + workers.remove(worker_ref); + return None; + }; + Some((entry.workspace_id, shared_state, entry.sink)) + } +} + +struct RuntimeGrantedWorkerObservationProvider { + runtime_id: String, + workspace_id: String, + grants: std::collections::HashSet, + hub: Arc, +} + +#[async_trait] +impl WorkerObservationProvider for RuntimeGrantedWorkerObservationProvider { + async fn list_worker_sessions( + &self, + ) -> Result, WorkerObservationError> { + let mut subjects = Vec::new(); + for grant in &self.grants { + if grant.runtime_id != self.runtime_id { + continue; + } + let Ok(worker_ref) = grant.local_worker_ref() else { + continue; + }; + let Some((workspace_id, state, _)) = self.hub.get(&worker_ref) else { + continue; + }; + if workspace_id.as_deref() != Some(self.workspace_id.as_str()) { + continue; + } + subjects.push(WorkerObservationSubject { + subject: WorkerObservationSubjectRef::RuntimeWorker { + runtime_id: grant.runtime_id.clone(), + worker_id: grant.worker_id.clone(), + }, + display_name: grant.worker_id.clone(), + relation: "granted_peer".to_string(), + status: format!("{:?}", state.get_status()).to_lowercase(), + }); + } + subjects.sort_by(|left, right| left.subject.cmp(&right.subject)); + Ok(subjects) + } + + async fn capture_worker_session( + &self, + subject: &WorkerObservationSubjectRef, + ) -> Result { + let WorkerObservationSubjectRef::RuntimeWorker { + runtime_id, + worker_id, + } = subject + else { + return Err(WorkerObservationError::NotFound); + }; + let grant = crate::identity::RuntimeWorkerRef::new(runtime_id.clone(), worker_id.clone()); + if !self.grants.contains(&grant) || runtime_id != &self.runtime_id { + return Err(WorkerObservationError::NotFound); + } + let worker_ref = grant + .local_worker_ref() + .map_err(|_| WorkerObservationError::NotFound)?; + let (workspace_id, _, sink) = self + .hub + .get(&worker_ref) + .ok_or(WorkerObservationError::NotFound)?; + if workspace_id.as_deref() != Some(self.workspace_id.as_str()) { + return Err(WorkerObservationError::NotFound); + } + let entries = sink.subscribe_with_snapshot().0; + let state = collect_state(&entries); + Ok(WorkerSessionCapture { + segment_id: format!("runtime:{runtime_id}:worker:{worker_id}"), + items: state.history, + }) + } +} + #[derive(Clone)] pub struct ProfileRuntimeWorkerFactory { + observation_hub: Arc, profile_base_dir: PathBuf, store_dir: Option, worker_metadata_dir: Option, @@ -116,6 +236,7 @@ impl ProfileRuntimeWorkerFactory { pub fn new(profile_base_dir: impl Into) -> Self { let profile_base_dir = profile_base_dir.into(); Self { + observation_hub: Arc::new(RuntimeWorkerObservationHub::default()), profile_base_dir, store_dir: None, worker_metadata_dir: None, @@ -398,6 +519,18 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { .unwrap_or(WorkerFilesystemAuthority::None); let workspace_backend_ref = RuntimeWorkspaceBackendRef::from_worker_request(&request.request); + let observation_runtime_id = request + .request + .workspace_api + .as_ref() + .and_then(|api| api.runtime_id.clone()); + let observation_workspace_id = request + .request + .workspace_api + .as_ref() + .map(|api| api.workspace_id.clone()); + let observation_grants = request.request.worker_observation_grants.clone(); + let observation_enabled = request.request.worker_observation_enabled; let workspace_context = workspace_backend_ref.worker_context(&request.worker_ref); let selector = profile.as_ref(); let archive = self @@ -457,11 +590,35 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { } else { worker.bind_workdir_session(None); } + if let (Some(runtime_id), Some(workspace_id)) = + (observation_runtime_id, observation_workspace_id.clone()) + && observation_enabled + { + let mut providers: Vec> = vec![Arc::new( + WorkspaceClientWorkerObservationProvider::new(worker.workspace_client_handle()), + )]; + if !observation_grants.is_empty() { + providers.push(Arc::new(RuntimeGrantedWorkerObservationProvider { + runtime_id, + workspace_id, + grants: observation_grants.into_iter().take(100).collect(), + hub: self.observation_hub.clone(), + })); + } + worker.bind_worker_observation_provider(Some(Arc::new( + CompositeWorkerObservationProvider::new(providers), + ))); + } let runtime_base = self.runtime_base_dir()?; let (handle, _shutdown_rx) = WorkerController::spawn_runtime_managed(worker, &runtime_base) .await .map_err(|err| format!("failed to spawn Worker controller: {err}"))?; + self.observation_hub.register( + request.worker_ref.clone(), + observation_workspace_id, + &handle, + ); Ok(handle) } @@ -482,6 +639,18 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { .unwrap_or(WorkerFilesystemAuthority::None); let workspace_backend_ref = RuntimeWorkspaceBackendRef::from_worker_request(&request.request); + let observation_runtime_id = request + .request + .workspace_api + .as_ref() + .and_then(|api| api.runtime_id.clone()); + let observation_workspace_id = request + .request + .workspace_api + .as_ref() + .map(|api| api.workspace_id.clone()); + let observation_grants = request.request.worker_observation_grants.clone(); + let observation_enabled = request.request.worker_observation_enabled; let workspace_context = workspace_backend_ref.worker_context(&request.worker_ref); let (manifest, loader) = Self::restore_fallback_manifest(&worker_name)?; @@ -552,11 +721,35 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { } else { worker.bind_workdir_session(None); } + if let (Some(runtime_id), Some(workspace_id)) = + (observation_runtime_id, observation_workspace_id.clone()) + && observation_enabled + { + let mut providers: Vec> = vec![Arc::new( + WorkspaceClientWorkerObservationProvider::new(worker.workspace_client_handle()), + )]; + if !observation_grants.is_empty() { + providers.push(Arc::new(RuntimeGrantedWorkerObservationProvider { + runtime_id, + workspace_id, + grants: observation_grants.into_iter().take(100).collect(), + hub: self.observation_hub.clone(), + })); + } + worker.bind_worker_observation_provider(Some(Arc::new( + CompositeWorkerObservationProvider::new(providers), + ))); + } let runtime_base = self.runtime_base_dir()?; let (handle, _shutdown_rx) = WorkerController::spawn_runtime_managed(worker, &runtime_base) .await .map_err(|err| format!("failed to spawn restored Worker controller: {err}"))?; + self.observation_hub.register( + request.worker_ref.clone(), + observation_workspace_id, + &handle, + ); Ok(handle) } } @@ -1603,6 +1796,8 @@ mod tests { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } @@ -1688,6 +1883,86 @@ mod tests { runtime_base.join(working_directory_id).join("checkout") } + #[tokio::test] + async fn runtime_provider_projects_only_explicit_live_canonical_grants() { + let hub = Arc::new(RuntimeWorkerObservationHub::default()); + let worker_ref = WorkerRef::new(crate::identity::WorkerId::new(7)); + let shared_state = Arc::new(WorkerSharedState::new( + "peer-worker".to_string(), + session_store::new_segment_id(), + "[worker]\nname = \"peer-worker\"".to_string(), + protocol::Greeting { + worker_name: "peer-worker".to_string(), + cwd: "/tmp".to_string(), + provider: "test".to_string(), + model: "test".to_string(), + scope_summary: String::new(), + tools: Vec::new(), + context_window: 1_000, + context_tokens: 0, + }, + )); + hub.workers.lock().unwrap().insert( + worker_ref, + RuntimeObservedWorker { + workspace_id: Some("workspace-1".to_string()), + shared_state: Arc::downgrade(&shared_state), + sink: SegmentLogSink::new(), + }, + ); + let grant = crate::identity::RuntimeWorkerRef::new("runtime-1", "7"); + let provider = RuntimeGrantedWorkerObservationProvider { + runtime_id: "runtime-1".to_string(), + workspace_id: "workspace-1".to_string(), + grants: std::collections::HashSet::from([grant.clone()]), + hub: hub.clone(), + }; + + let listed = provider.list_worker_sessions().await.unwrap(); + assert_eq!(listed.len(), 1); + assert_eq!( + listed[0].subject, + WorkerObservationSubjectRef::RuntimeWorker { + runtime_id: "runtime-1".to_string(), + worker_id: "7".to_string(), + } + ); + provider + .capture_worker_session(&listed[0].subject) + .await + .expect("granted live peer should be capturable"); + let hidden = provider + .capture_worker_session(&WorkerObservationSubjectRef::RuntimeWorker { + runtime_id: "runtime-1".to_string(), + worker_id: "8".to_string(), + }) + .await + .unwrap_err(); + assert!(matches!(hidden, WorkerObservationError::NotFound)); + + let cross_workspace = RuntimeGrantedWorkerObservationProvider { + runtime_id: "runtime-1".to_string(), + workspace_id: "workspace-2".to_string(), + grants: std::collections::HashSet::from([grant]), + hub: hub.clone(), + }; + assert!( + cross_workspace + .list_worker_sessions() + .await + .unwrap() + .is_empty() + ); + let hidden = cross_workspace + .capture_worker_session(&listed[0].subject) + .await + .unwrap_err(); + assert!(matches!(hidden, WorkerObservationError::NotFound)); + + drop(shared_state); + assert!(provider.list_worker_sessions().await.unwrap().is_empty()); + } + #[test] fn runtime_worker_name_is_runtime_local() { let worker_ref = crate::identity::WorkerRef::new(crate::identity::WorkerId::new(1)); diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 86669e35..20b28e66 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -372,6 +372,12 @@ pub struct WorkerSpawnRequest { pub resolved_config_bundle: Option, #[serde(skip, default)] pub resolved_workspace_api: Option, + /// Backend-owned feature enablement; client input cannot set it. + #[serde(skip, default)] + pub resolved_worker_observation_enabled: bool, + /// Backend-authored peer-session grants. Browser/model input cannot set this field. + #[serde(skip, default)] + pub resolved_worker_observation_grants: Vec, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -1858,6 +1864,8 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { initial_input: request.initial_input.clone(), working_directory_request: request.resolved_working_directory_request.clone(), working_directory: request.resolved_working_directory.clone(), + worker_observation_enabled: request.resolved_worker_observation_enabled, + worker_observation_grants: request.resolved_worker_observation_grants.clone(), workspace_api: Some(workspace_api), }; let workspace_scope = RuntimeWorkspaceScope::new(workspace_id, "embedded-backend"); @@ -2959,6 +2967,8 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { initial_input: request.initial_input.clone(), working_directory_request: request.resolved_working_directory_request.clone(), working_directory: request.resolved_working_directory.clone(), + worker_observation_enabled: request.resolved_worker_observation_enabled, + worker_observation_grants: request.resolved_worker_observation_grants.clone(), workspace_api: Some(workspace_api), }; match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) { @@ -4538,6 +4548,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_workspace_api()), } } @@ -4687,6 +4699,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_workspace_api()), }, ) @@ -4782,6 +4796,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_workspace_api()), }, ) @@ -4816,6 +4832,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_workspace_api()), }, ) diff --git a/crates/workspace-server/src/runtime_subscription_tests.rs b/crates/workspace-server/src/runtime_subscription_tests.rs index 921af63f..26b37d16 100644 --- a/crates/workspace-server/src/runtime_subscription_tests.rs +++ b/crates/workspace-server/src/runtime_subscription_tests.rs @@ -74,6 +74,8 @@ fn create_request(name: &str) -> CreateWorkerRequest { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index c03ff144..f6645cf7 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -42,6 +42,7 @@ use webauthn_rs::prelude::{ }; use workdir::WorkdirSessionHandle; use workdir::http::{WorkdirSessionOperation, WorkdirSessionOperationResult}; +use worker::feature::builtin::{WorkerObservationSubject, WorkerObservationSubjectRef}; use worker_runtime::resource::{BackendResourceError, BackendResourceFetchRequest}; use worker_runtime::worker_backend::{ProfileRuntimeWorkerFactory, WorkerRuntimeExecutionBackend}; @@ -921,6 +922,14 @@ pub fn build_router(api: WorkspaceApi) -> Router { get(scoped_workspace_orchestrator_status) .post(scoped_start_workspace_orchestrator), ) + .route( + "/api/w/{workspace_id}/worker-observation/sessions", + get(scoped_list_worker_observation_sessions), + ) + .route( + "/api/w/{workspace_id}/worker-observation/session", + post(scoped_capture_worker_observation_session), + ) .route( "/api/w/{workspace_id}/workers", get(scoped_list_workers).post(scoped_create_workspace_worker), @@ -3828,6 +3837,8 @@ fn start_memory_staging_consolidation( resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, }, )?; @@ -4200,6 +4211,120 @@ async fn scoped_workspace_orchestrator_status( Ok(Json(workspace_orchestrator_response(&api, "observed"))) } +async fn scoped_list_worker_observation_sessions( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + let source = authenticate_worker_mutation_source(&api, &path.workspace_id, &headers)?; + authorize_workspace_orchestrator_observation(&api, &source)?; + let sessions = workers_response(api.clone())? + .items + .into_iter() + .filter(|worker| { + !matches!( + worker.state.as_str(), + "stopped" | "failed" | "rejected" | "disconnected" + ) && (worker.worker.runtime_id != source.runtime_id + || worker.worker.worker_id != source.worker_id) + }) + .take(100) + .map(|worker| WorkerObservationSubject { + subject: WorkerObservationSubjectRef::RuntimeWorker { + runtime_id: worker.worker.runtime_id, + worker_id: worker.worker.worker_id, + }, + display_name: worker.display_name, + relation: "workspace_orchestrator_grant".to_string(), + status: worker.state, + }) + .collect::>(); + Ok(Json(serde_json::json!({ "sessions": sessions }))) +} + +async fn scoped_capture_worker_observation_session( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(subject): Json, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + let source = authenticate_worker_mutation_source(&api, &path.workspace_id, &headers)?; + authorize_workspace_orchestrator_observation(&api, &source)?; + let WorkerObservationSubjectRef::RuntimeWorker { + runtime_id, + worker_id, + } = subject + else { + return Err(ApiError::from(Error::UnknownWorker { + worker: RuntimeWorkerRef::new("subworker", "inaccessible"), + })); + }; + let target = RuntimeWorkerRef::new(runtime_id, worker_id); + let granted = workers_response(api.clone())? + .items + .into_iter() + .any(|worker| { + worker.worker == target + && !matches!( + worker.state.as_str(), + "stopped" | "failed" | "rejected" | "disconnected" + ) + }); + if !granted { + return Err(ApiError::from(Error::UnknownWorker { worker: target })); + } + + let mut connection = connect_workspace_worker_protocol(&api, &target).await?; + let event = tokio::time::timeout(std::time::Duration::from_secs(10), connection.events.recv()) + .await + .map_err(|_| { + ApiError::from(Error::RuntimeOperationFailed { + runtime_id: target.runtime_id.clone(), + code: "worker_observation_timeout".to_string(), + message: "timed out waiting for the committed session snapshot".to_string(), + }) + })? + .ok_or_else(|| { + ApiError::from(Error::RuntimeOperationFailed { + runtime_id: target.runtime_id.clone(), + code: "worker_observation_closed".to_string(), + message: "worker protocol closed before the session snapshot".to_string(), + }) + })?; + let protocol::Event::Snapshot { entries, .. } = event else { + return Err(ApiError::from(Error::RuntimeOperationFailed { + runtime_id: target.runtime_id.clone(), + code: "worker_observation_missing_snapshot".to_string(), + message: "worker protocol did not begin with a committed session snapshot".to_string(), + })); + }; + Ok(Json(serde_json::json!({ + "segment_id": format!("runtime:{}:worker:{}", target.runtime_id, target.worker_id), + "entries": entries, + }))) +} + +fn authorize_workspace_orchestrator_observation( + api: &WorkspaceApi, + source: &WorkerMutationSource, +) -> ApiResult<()> { + let Some(orchestrator) = find_workspace_orchestrator(api) else { + return Err(ApiError::from(Error::UnknownWorker { + worker: RuntimeWorkerRef::new(source.runtime_id.clone(), source.worker_id.clone()), + })); + }; + if orchestrator.worker.runtime_id != source.runtime_id + || orchestrator.worker.worker_id != source.worker_id + { + return Err(ApiError::from(Error::UnknownWorker { + worker: RuntimeWorkerRef::new(source.runtime_id.clone(), source.worker_id.clone()), + })); + } + Ok(()) +} + async fn scoped_start_workspace_orchestrator( State(api): State, AxumPath(path): AxumPath, @@ -4258,6 +4383,17 @@ async fn scoped_start_workspace_orchestrator( resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: true, + resolved_worker_observation_grants: workers_response(api.clone()) + .map(|response| { + response + .items + .into_iter() + .take(100) + .map(|worker| worker.worker) + .collect() + }) + .unwrap_or_default(), resolved_workspace_api: None, }, )?; @@ -6655,6 +6791,8 @@ async fn create_workspace_worker( resolved_working_directory_request: None, resolved_working_directory, resolved_config_bundle, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, }, )?; @@ -9921,6 +10059,70 @@ mod tests { ); assert_ne!(dedicated.worker.worker_id, generic.worker_ref.worker_id); + let mut observation_headers = HeaderMap::new(); + observation_headers.insert( + "x-yoi-runtime-id", + axum::http::HeaderValue::from_str(&dedicated.worker.runtime_id).unwrap(), + ); + observation_headers.insert( + "x-yoi-worker-id", + axum::http::HeaderValue::from_str(&dedicated.worker.worker_id).unwrap(), + ); + let Json(sessions) = scoped_list_worker_observation_sessions( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + observation_headers.clone(), + ) + .await + .unwrap(); + assert!( + sessions["sessions"] + .as_array() + .unwrap() + .iter() + .any(|session| { + session["subject"]["kind"] == "runtime_worker" + && session["subject"]["runtime_id"] == generic.worker_ref.runtime_id + && session["subject"]["worker_id"] == generic.worker_ref.worker_id + }) + ); + let Json(capture) = scoped_capture_worker_observation_session( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + observation_headers, + Json(WorkerObservationSubjectRef::RuntimeWorker { + runtime_id: generic.worker_ref.runtime_id.clone(), + worker_id: generic.worker_ref.worker_id.clone(), + }), + ) + .await + .unwrap(); + assert!(capture["entries"].is_array()); + + let mut unauthorized_headers = HeaderMap::new(); + unauthorized_headers.insert( + "x-yoi-runtime-id", + axum::http::HeaderValue::from_str(&generic.worker_ref.runtime_id).unwrap(), + ); + unauthorized_headers.insert( + "x-yoi-worker-id", + axum::http::HeaderValue::from_str(&generic.worker_ref.worker_id).unwrap(), + ); + let error = scoped_list_worker_observation_sessions( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + unauthorized_headers, + ) + .await + .unwrap_err(); + assert_eq!(error.into_response().status(), StatusCode::NOT_FOUND); + let Json(existing) = scoped_start_workspace_orchestrator( State(api.clone()), AxumPath(ScopedWorkspacePath { @@ -10691,6 +10893,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_worker_workspace_api( EMBEDDED_WORKER_RUNTIME_ID, )), @@ -10871,6 +11075,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)), }; let source_worker = api @@ -11085,6 +11291,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_worker_workspace_api( EMBEDDED_WORKER_RUNTIME_ID, )), @@ -11225,6 +11433,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, }; let Json(first) = scoped_create_runtime_worker( @@ -11455,6 +11665,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, }; let Json(created) = scoped_create_runtime_worker( @@ -12303,6 +12515,8 @@ mod tests { initial_input: None, working_directory_request: None, working_directory: None, + worker_observation_enabled: false, + worker_observation_grants: Vec::new(), workspace_api: None, } } @@ -13454,6 +13668,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_worker_workspace_api( "embedded-worker-runtime", )), @@ -13972,6 +14188,8 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: Some(runtime_test_bundle()), + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, }; let spawned = api