refactor: separate worker workspace identity
This commit is contained in:
parent
74b07bd787
commit
142b60e1b3
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -5848,6 +5848,7 @@ dependencies = [
|
||||||
"thiserror 2.0.18",
|
"thiserror 2.0.18",
|
||||||
"tokio",
|
"tokio",
|
||||||
"tokio-tungstenite 0.29.0",
|
"tokio-tungstenite 0.29.0",
|
||||||
|
"toml",
|
||||||
"tower",
|
"tower",
|
||||||
"worker",
|
"worker",
|
||||||
]
|
]
|
||||||
|
|
|
||||||
|
|
@ -100,8 +100,13 @@ pub struct WorkerMetadata {
|
||||||
pub worker_name: String,
|
pub worker_name: String,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub active: Option<WorkerActiveSegmentRef>,
|
pub active: Option<WorkerActiveSegmentRef>,
|
||||||
|
/// Legacy local path hint retained for host/runtime compatibility. It is not
|
||||||
|
/// Worker workspace identity or authority.
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub workspace_root: Option<PathBuf>,
|
pub workspace_root: Option<PathBuf>,
|
||||||
|
/// Path-free workspace identity supplied by the host/runtime boundary.
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub workspace_id: Option<String>,
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||||
pub spawned_children: Vec<WorkerSpawnedChild>,
|
pub spawned_children: Vec<WorkerSpawnedChild>,
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||||
|
|
@ -119,6 +124,7 @@ impl WorkerMetadata {
|
||||||
worker_name: worker_name.into(),
|
worker_name: worker_name.into(),
|
||||||
active,
|
active,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -130,6 +136,11 @@ impl WorkerMetadata {
|
||||||
self.workspace_root = Some(workspace_root);
|
self.workspace_root = Some(workspace_root);
|
||||||
self
|
self
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn with_workspace_id(mut self, workspace_id: impl Into<String>) -> Self {
|
||||||
|
self.workspace_id = Some(workspace_id.into());
|
||||||
|
self
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Sync persistence backend for Worker metadata.
|
/// Sync persistence backend for Worker metadata.
|
||||||
|
|
@ -180,6 +191,44 @@ pub trait WorkerMetadataStore: Send + Sync {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Set the active pointer and workspace ownership while preserving unrelated fields.
|
/// Set the active pointer and workspace ownership while preserving unrelated fields.
|
||||||
|
fn set_active_with_workspace_context(
|
||||||
|
&self,
|
||||||
|
worker_name: &str,
|
||||||
|
active: Option<WorkerActiveSegmentRef>,
|
||||||
|
resolved_manifest_snapshot: Option<serde_json::Value>,
|
||||||
|
workspace_id: Option<String>,
|
||||||
|
workspace_root: Option<PathBuf>,
|
||||||
|
) -> Result<WorkerMetadata, WorkerStoreError> {
|
||||||
|
self.update_by_name(worker_name, |metadata| {
|
||||||
|
metadata.active = active;
|
||||||
|
metadata.resolved_manifest_snapshot = resolved_manifest_snapshot;
|
||||||
|
if let Some(workspace_id) = workspace_id {
|
||||||
|
metadata.workspace_id = Some(workspace_id);
|
||||||
|
}
|
||||||
|
if let Some(workspace_root) = workspace_root {
|
||||||
|
metadata.workspace_root = Some(workspace_root);
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Set the active pointer and path-free workspace identity while preserving unrelated fields.
|
||||||
|
fn set_active_with_workspace_id(
|
||||||
|
&self,
|
||||||
|
worker_name: &str,
|
||||||
|
active: Option<WorkerActiveSegmentRef>,
|
||||||
|
resolved_manifest_snapshot: Option<serde_json::Value>,
|
||||||
|
workspace_id: Option<String>,
|
||||||
|
) -> Result<WorkerMetadata, WorkerStoreError> {
|
||||||
|
self.update_by_name(worker_name, |metadata| {
|
||||||
|
metadata.active = active;
|
||||||
|
metadata.resolved_manifest_snapshot = resolved_manifest_snapshot;
|
||||||
|
if let Some(workspace_id) = workspace_id {
|
||||||
|
metadata.workspace_id = Some(workspace_id);
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Set the active pointer and legacy local workspace-root hint while preserving unrelated fields.
|
||||||
fn set_active_with_workspace_root(
|
fn set_active_with_workspace_root(
|
||||||
&self,
|
&self,
|
||||||
worker_name: &str,
|
worker_name: &str,
|
||||||
|
|
|
||||||
|
|
@ -32,6 +32,7 @@ reqwest = { version = "0.13", optional = true, default-features = false, feature
|
||||||
tar.workspace = true
|
tar.workspace = true
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
tokio = { workspace = true, features = ["net", "rt", "sync", "time"] }
|
tokio = { workspace = true, features = ["net", "rt", "sync", "time"] }
|
||||||
|
toml.workspace = true
|
||||||
tower = { workspace = true, features = ["util"], optional = true }
|
tower = { workspace = true, features = ["util"], optional = true }
|
||||||
worker.workspace = true
|
worker.workspace = true
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,7 @@
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use std::path::PathBuf;
|
use std::path::{Path, PathBuf};
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::sync::{Arc, Mutex, mpsc};
|
use std::sync::{Arc, Mutex, mpsc};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
@ -35,7 +35,10 @@ use session_store::{CombinedStore, FsWorkerStore};
|
||||||
use tokio::runtime::Runtime;
|
use tokio::runtime::Runtime;
|
||||||
use tokio::sync::broadcast;
|
use tokio::sync::broadcast;
|
||||||
|
|
||||||
use worker::{Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle};
|
use worker::{
|
||||||
|
Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle, WorkerWorkspaceContext,
|
||||||
|
WorkspaceId,
|
||||||
|
};
|
||||||
|
|
||||||
const DEFAULT_BACKEND_ID: &str = "worker-crate";
|
const DEFAULT_BACKEND_ID: &str = "worker-crate";
|
||||||
const RUNTIME_TASK_TIMEOUT: Duration = Duration::from_secs(10);
|
const RUNTIME_TASK_TIMEOUT: Duration = Duration::from_secs(10);
|
||||||
|
|
@ -230,6 +233,41 @@ fn sanitize_worker_name_component(value: &str) -> String {
|
||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
|
enum RuntimeWorkspaceBackendRef {
|
||||||
|
None,
|
||||||
|
LocalFilesystem { root: PathBuf },
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RuntimeWorkspaceBackendRef {
|
||||||
|
fn from_working_directory(binding: Option<&WorkingDirectoryBinding>) -> Self {
|
||||||
|
match binding {
|
||||||
|
Some(binding) => Self::LocalFilesystem {
|
||||||
|
root: binding.root().to_path_buf(),
|
||||||
|
},
|
||||||
|
None => Self::None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn worker_context(&self) -> WorkerWorkspaceContext {
|
||||||
|
match self {
|
||||||
|
Self::None => WorkerWorkspaceContext::no_workspace(),
|
||||||
|
Self::LocalFilesystem { root } => local_workspace_context(root),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn local_workspace_context(root: &Path) -> WorkerWorkspaceContext {
|
||||||
|
WorkerWorkspaceContext::local_filesystem(read_workspace_id_hint(root))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn read_workspace_id_hint(root: &Path) -> Option<WorkspaceId> {
|
||||||
|
let contents = std::fs::read_to_string(root.join(".yoi/workspace.toml")).ok()?;
|
||||||
|
let value = toml::from_str::<toml::Value>(&contents).ok()?;
|
||||||
|
let id = value.get("id")?.as_str()?.to_string();
|
||||||
|
WorkspaceId::new(id).ok()
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "http-server")]
|
#[cfg(feature = "http-server")]
|
||||||
async fn fetch_profile_source_archive_http(
|
async fn fetch_profile_source_archive_http(
|
||||||
location: &ProfileSourceArchiveHttpRef,
|
location: &ProfileSourceArchiveHttpRef,
|
||||||
|
|
@ -308,6 +346,9 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
.unwrap_or(WorkerFilesystemAuthority::None);
|
.unwrap_or(WorkerFilesystemAuthority::None);
|
||||||
|
let workspace_backend_ref =
|
||||||
|
RuntimeWorkspaceBackendRef::from_working_directory(request.working_directory.as_ref());
|
||||||
|
let workspace_context = workspace_backend_ref.worker_context();
|
||||||
let selector = profile.as_deref().unwrap_or("builtin:default");
|
let selector = profile.as_deref().unwrap_or("builtin:default");
|
||||||
let archive = self
|
let archive = self
|
||||||
.resolve_profile_source_archive(&request.request.profile_source)
|
.resolve_profile_source_archive(&request.request.profile_source)
|
||||||
|
|
@ -344,7 +385,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
worker_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -927,12 +968,16 @@ mod tests {
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.map(|binding| binding.root().to_path_buf())
|
.map(|binding| binding.root().to_path_buf())
|
||||||
.unwrap_or_else(|| self.cwd.clone());
|
.unwrap_or_else(|| self.cwd.clone());
|
||||||
|
let workspace_backend_ref = RuntimeWorkspaceBackendRef::from_working_directory(
|
||||||
|
request.working_directory.as_ref(),
|
||||||
|
);
|
||||||
|
let workspace_context = workspace_backend_ref.worker_context();
|
||||||
let scope = Scope::writable(&scope_root).map_err(|err| err.to_string())?;
|
let scope = Scope::writable(&scope_root).map_err(|err| err.to_string())?;
|
||||||
let worker = Worker::new(
|
let worker = Worker::new(
|
||||||
manifest,
|
manifest,
|
||||||
Engine::new(self.client.clone()),
|
Engine::new(self.client.clone()),
|
||||||
store,
|
store,
|
||||||
self.cwd.clone(),
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -522,17 +522,16 @@ fn install_ticket_event_companion_notify_hook<C, St>(
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
let Some(companion_worker_name) = companion_worker_name_for_workspace(worker.workspace_root())
|
let Some(local) = worker.local_working_directory() else {
|
||||||
else {
|
return;
|
||||||
|
};
|
||||||
|
let Some(companion_worker_name) = companion_worker_name_for_workspace(&local.root) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
if companion_worker_name == worker.manifest().worker.name {
|
if companion_worker_name == worker.manifest().worker.name {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
let Some(local) = worker.local_working_directory() else {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
let Ok(ticket_config) = TicketConfig::load_workspace(&local.cwd) else {
|
let Ok(ticket_config) = TicketConfig::load_workspace(&local.cwd) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
|
|
@ -603,7 +602,7 @@ where
|
||||||
// below so the worker borrow doesn't conflict with reads on `worker`.
|
// below so the worker borrow doesn't conflict with reads on `worker`.
|
||||||
let scope_handle = worker.scope().clone();
|
let scope_handle = worker.scope().clone();
|
||||||
let local_filesystem = worker.local_working_directory().cloned();
|
let local_filesystem = worker.local_working_directory().cloned();
|
||||||
let workspace_root = worker.workspace_root().to_path_buf();
|
let local_workspace_root = local_filesystem.as_ref().map(|local| local.root.clone());
|
||||||
let task_feature = worker.task_feature();
|
let task_feature = worker.task_feature();
|
||||||
let session_id_for_usage = worker.segment_id().to_string();
|
let session_id_for_usage = worker.segment_id().to_string();
|
||||||
let memory_config = worker.manifest().memory.clone();
|
let memory_config = worker.manifest().memory.clone();
|
||||||
|
|
@ -654,8 +653,8 @@ where
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
// Ticket tools are typed operations over the currently checked-out work
|
// Ticket tools are typed operations over the currently checked-out work
|
||||||
// tree. They require explicit local filesystem authority; workspace_root
|
// tree. They require explicit local filesystem authority and must not
|
||||||
// is context only and must not be used as a cwd fallback.
|
// use workspace identity as a cwd fallback.
|
||||||
let ticket_cwd = local_filesystem
|
let ticket_cwd = local_filesystem
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.map(|local| &local.cwd)
|
.map(|local| &local.cwd)
|
||||||
|
|
@ -679,10 +678,12 @@ where
|
||||||
) {
|
) {
|
||||||
feature_registry = feature_registry.with_module(module);
|
feature_registry = feature_registry.with_module(module);
|
||||||
}
|
}
|
||||||
if let Some(module) =
|
if let Some(workspace_root) = local_workspace_root.as_ref() {
|
||||||
crate::feature::mcp::discover_stdio_tool_feature(&mcp_config, &workspace_root).await
|
if let Some(module) =
|
||||||
{
|
crate::feature::mcp::discover_stdio_tool_feature(&mcp_config, workspace_root).await
|
||||||
feature_registry = feature_registry.with_module(module);
|
{
|
||||||
|
feature_registry = feature_registry.with_module(module);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
{
|
{
|
||||||
|
|
@ -698,7 +699,13 @@ where
|
||||||
"[feature.memory].enabled = true requires a [memory] configuration section",
|
"[feature.memory].enabled = true requires a [memory] configuration section",
|
||||||
)
|
)
|
||||||
})?;
|
})?;
|
||||||
let layout = memory::WorkspaceLayout::resolve(mem, &workspace_root);
|
let workspace_root = local_workspace_root.as_ref().ok_or_else(|| {
|
||||||
|
std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidInput,
|
||||||
|
"memory tools require local Worker filesystem authority",
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
let layout = memory::WorkspaceLayout::resolve(mem, workspace_root);
|
||||||
let query_cfg = memory::tool::QueryConfig::from(mem);
|
let query_cfg = memory::tool::QueryConfig::from(mem);
|
||||||
worker.register_tool(memory::tool::read_tool_with_usage(
|
worker.register_tool(memory::tool::read_tool_with_usage(
|
||||||
layout.clone(),
|
layout.clone(),
|
||||||
|
|
@ -732,11 +739,17 @@ where
|
||||||
"worker spawn tools require local Worker filesystem authority",
|
"worker spawn tools require local Worker filesystem authority",
|
||||||
)
|
)
|
||||||
})?;
|
})?;
|
||||||
|
let spawner_workspace_root = local_workspace_root.clone().ok_or_else(|| {
|
||||||
|
std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidInput,
|
||||||
|
"worker spawn tools require local Worker filesystem authority",
|
||||||
|
)
|
||||||
|
})?;
|
||||||
worker.register_tool(spawn_worker_tool(
|
worker.register_tool(spawn_worker_tool(
|
||||||
spawner_name.clone(),
|
spawner_name.clone(),
|
||||||
spawner_socket,
|
spawner_socket,
|
||||||
runtime_base.clone(),
|
runtime_base.clone(),
|
||||||
workspace_root.clone(),
|
spawner_workspace_root,
|
||||||
spawner_cwd.clone(),
|
spawner_cwd.clone(),
|
||||||
spawned_registry.clone(),
|
spawned_registry.clone(),
|
||||||
self_parent_socket,
|
self_parent_socket,
|
||||||
|
|
|
||||||
|
|
@ -1145,6 +1145,7 @@ mod tests {
|
||||||
worker_name: "parent".into(),
|
worker_name: "parent".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: vec![
|
spawned_children: vec![
|
||||||
child("child-live", &live_socket),
|
child("child-live", &live_socket),
|
||||||
child("child-stale", &stale_socket),
|
child("child-stale", &stale_socket),
|
||||||
|
|
@ -1165,6 +1166,7 @@ mod tests {
|
||||||
active_child_segment,
|
active_child_segment,
|
||||||
)),
|
)),
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -1179,6 +1181,7 @@ mod tests {
|
||||||
active_child_segment,
|
active_child_segment,
|
||||||
)),
|
)),
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -1190,6 +1193,7 @@ mod tests {
|
||||||
worker_name: "child-pending".into(),
|
worker_name: "child-pending".into(),
|
||||||
active: Some(WorkerActiveSegmentRef::pending_segment(pending_session_id)),
|
active: Some(WorkerActiveSegmentRef::pending_segment(pending_session_id)),
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -1204,6 +1208,7 @@ mod tests {
|
||||||
new_segment_id(),
|
new_segment_id(),
|
||||||
)),
|
)),
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -1215,6 +1220,7 @@ mod tests {
|
||||||
worker_name: "peer".into(),
|
worker_name: "peer".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1421,6 +1427,7 @@ mod tests {
|
||||||
worker_name: "source".into(),
|
worker_name: "source".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1461,6 +1468,7 @@ mod tests {
|
||||||
worker_name: "source".into(),
|
worker_name: "source".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1474,6 +1482,7 @@ mod tests {
|
||||||
worker_name: "target".into(),
|
worker_name: "target".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1579,6 +1588,7 @@ mod tests {
|
||||||
worker_name: "source".into(),
|
worker_name: "source".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1592,6 +1602,7 @@ mod tests {
|
||||||
worker_name: "target".into(),
|
worker_name: "target".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: vec![session_store::WorkerPeer {
|
peers: vec![session_store::WorkerPeer {
|
||||||
|
|
@ -1695,6 +1706,7 @@ mod tests {
|
||||||
worker_name: "source".into(),
|
worker_name: "source".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: vec![child("target", &socket)],
|
spawned_children: vec![child("target", &socket)],
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,10 @@ use std::ffi::OsString;
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
use std::process::ExitCode;
|
use std::process::ExitCode;
|
||||||
|
|
||||||
use crate::{PromptLoader, Worker, WorkerController, WorkerFilesystemAuthority};
|
use crate::{
|
||||||
|
PromptLoader, Worker, WorkerController, WorkerFilesystemAuthority, WorkerWorkspaceContext,
|
||||||
|
WorkspaceId,
|
||||||
|
};
|
||||||
use clap::{CommandFactory, FromArgMatches, Parser};
|
use clap::{CommandFactory, FromArgMatches, Parser};
|
||||||
use manifest::{
|
use manifest::{
|
||||||
Permission, ProfileResolveOptions, ProfileResolver, ProfileSelector, ScopeConfig, ScopeRule,
|
Permission, ProfileResolveOptions, ProfileResolver, ProfileSelector, ScopeConfig, ScopeRule,
|
||||||
|
|
@ -103,6 +106,24 @@ fn runtime_workspace_root(cli: &Cli) -> Result<PathBuf, String> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn runtime_workspace_context(workspace_root: &Path) -> WorkerWorkspaceContext {
|
||||||
|
WorkerWorkspaceContext::local_filesystem(read_workspace_id_hint(workspace_root))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn read_workspace_id_hint(workspace_root: &Path) -> Option<WorkspaceId> {
|
||||||
|
let path = workspace_root.join(".yoi/workspace.toml");
|
||||||
|
let contents = std::fs::read_to_string(path).ok()?;
|
||||||
|
let value = toml::from_str::<toml::Value>(&contents).ok()?;
|
||||||
|
let id = value.get("id")?.as_str()?.to_string();
|
||||||
|
match WorkspaceId::new(id) {
|
||||||
|
Ok(id) => Some(id),
|
||||||
|
Err(err) => {
|
||||||
|
tracing::warn!("ignoring invalid workspace id in .yoi/workspace.toml: {err}");
|
||||||
|
None
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn runtime_worker_name(cli: &Cli, workspace_root: &Path) -> String {
|
fn runtime_worker_name(cli: &Cli, workspace_root: &Path) -> String {
|
||||||
cli.worker
|
cli.worker
|
||||||
.as_deref()
|
.as_deref()
|
||||||
|
|
@ -513,6 +534,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
let store = CombinedStore::new(session_store, worker_metadata_store);
|
let store = CombinedStore::new(session_store, worker_metadata_store);
|
||||||
let filesystem_authority =
|
let filesystem_authority =
|
||||||
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone());
|
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone());
|
||||||
|
let workspace_context = runtime_workspace_context(&workspace_root);
|
||||||
|
|
||||||
let mut worker = if cli.adopt {
|
let mut worker = if cli.adopt {
|
||||||
let callback = match cli.callback.clone() {
|
let callback = match cli.callback.clone() {
|
||||||
|
|
@ -527,7 +549,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
callback,
|
callback,
|
||||||
workspace_root.clone(),
|
workspace_context.clone(),
|
||||||
filesystem_authority.clone(),
|
filesystem_authority.clone(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -558,7 +580,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_root.clone(),
|
workspace_context.clone(),
|
||||||
filesystem_authority.clone(),
|
filesystem_authority.clone(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -578,7 +600,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_root.clone(),
|
workspace_context.clone(),
|
||||||
filesystem_authority.clone(),
|
filesystem_authority.clone(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -599,7 +621,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_root.clone(),
|
workspace_context.clone(),
|
||||||
filesystem_authority.clone(),
|
filesystem_authority.clone(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -621,7 +643,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode {
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_root.clone(),
|
workspace_context.clone(),
|
||||||
filesystem_authority.clone(),
|
filesystem_authority.clone(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
|
||||||
|
|
@ -40,5 +40,5 @@ pub use segment_log_sink::SegmentLogSink;
|
||||||
pub use shared_state::WorkerSharedState;
|
pub use shared_state::WorkerSharedState;
|
||||||
pub use worker::{
|
pub use worker::{
|
||||||
LocalWorkingDirectory, Worker, WorkerError, WorkerFilesystemAuthority, WorkerRunResult,
|
LocalWorkingDirectory, Worker, WorkerError, WorkerFilesystemAuthority, WorkerRunResult,
|
||||||
apply_worker_manifest,
|
WorkerWorkspaceContext, WorkspaceClient, WorkspaceId, WorkspaceIdError, apply_worker_manifest,
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -384,6 +384,7 @@ mod tests {
|
||||||
worker_name: "orchestrator".into(),
|
worker_name: "orchestrator".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
@ -395,6 +396,7 @@ mod tests {
|
||||||
worker_name: "companion".into(),
|
worker_name: "companion".into(),
|
||||||
active: None,
|
active: None,
|
||||||
workspace_root: None,
|
workspace_root: None,
|
||||||
|
workspace_id: None,
|
||||||
spawned_children: Vec::new(),
|
spawned_children: Vec::new(),
|
||||||
reclaimed_children: Vec::new(),
|
reclaimed_children: Vec::new(),
|
||||||
peers: Vec::new(),
|
peers: Vec::new(),
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,6 @@
|
||||||
use std::path::{Path, PathBuf};
|
#[cfg(test)]
|
||||||
|
use std::path::Path;
|
||||||
|
use std::path::PathBuf;
|
||||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
@ -92,6 +94,113 @@ pub struct LocalWorkingDirectory {
|
||||||
pub cwd: PathBuf,
|
pub cwd: PathBuf,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Path-free workspace identity carried by a Worker.
|
||||||
|
///
|
||||||
|
/// The value is intentionally opaque to Worker code: Runtime/host layers own
|
||||||
|
/// backend lookup, endpoint/auth/secret materialisation, and any mapping from a
|
||||||
|
/// local checkout path to an id. Worker code may only compare/log the id and pass
|
||||||
|
/// it through to narrow workspace-aware handles.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||||
|
pub struct WorkspaceId(String);
|
||||||
|
|
||||||
|
impl WorkspaceId {
|
||||||
|
pub fn new(id: impl Into<String>) -> Result<Self, WorkspaceIdError> {
|
||||||
|
let id = id.into();
|
||||||
|
if id.trim().is_empty() {
|
||||||
|
return Err(WorkspaceIdError::Empty);
|
||||||
|
}
|
||||||
|
Ok(Self(id))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn as_str(&self) -> &str {
|
||||||
|
&self.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, thiserror::Error, Clone, PartialEq, Eq)]
|
||||||
|
pub enum WorkspaceIdError {
|
||||||
|
#[error("workspace id must not be empty")]
|
||||||
|
Empty,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Narrow path-free workspace API handle injected by Runtime/host code.
|
||||||
|
///
|
||||||
|
/// This is deliberately not a filesystem authority surface. A Worker may have a
|
||||||
|
/// workspace client without local filesystem authority, or neither. Local
|
||||||
|
/// path-backed implementations are represented only as a capability marker here;
|
||||||
|
/// the actual paths remain under [`WorkerFilesystemAuthority::Local`] or in host
|
||||||
|
/// adapter code.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub enum WorkspaceClient {
|
||||||
|
/// Runtime/host supplied a workspace API handle. The string is an opaque
|
||||||
|
/// diagnostic/backend kind, not an endpoint, path, or secret-bearing value.
|
||||||
|
Available { kind: String },
|
||||||
|
/// Workspace-aware operations must fail closed or stay disabled.
|
||||||
|
Unavailable { reason: String },
|
||||||
|
}
|
||||||
|
|
||||||
|
impl WorkspaceClient {
|
||||||
|
pub fn available(kind: impl Into<String>) -> Self {
|
||||||
|
Self::Available { kind: kind.into() }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn unavailable(reason: impl Into<String>) -> Self {
|
||||||
|
Self::Unavailable {
|
||||||
|
reason: reason.into(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn local_filesystem() -> Self {
|
||||||
|
Self::available("local-filesystem")
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn is_available(&self) -> bool {
|
||||||
|
matches!(self, Self::Available { .. })
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Workspace context supplied to a Worker separately from filesystem authority.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct WorkerWorkspaceContext {
|
||||||
|
workspace_id: Option<WorkspaceId>,
|
||||||
|
client: WorkspaceClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl WorkerWorkspaceContext {
|
||||||
|
pub fn no_workspace() -> Self {
|
||||||
|
Self {
|
||||||
|
workspace_id: None,
|
||||||
|
client: WorkspaceClient::unavailable("no workspace configured"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn unavailable(workspace_id: Option<WorkspaceId>, reason: impl Into<String>) -> Self {
|
||||||
|
Self {
|
||||||
|
workspace_id,
|
||||||
|
client: WorkspaceClient::unavailable(reason),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn with_client(workspace_id: Option<WorkspaceId>, client: WorkspaceClient) -> Self {
|
||||||
|
Self {
|
||||||
|
workspace_id,
|
||||||
|
client,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn local_filesystem(workspace_id: Option<WorkspaceId>) -> Self {
|
||||||
|
Self::with_client(workspace_id, WorkspaceClient::local_filesystem())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn workspace_id(&self) -> Option<&WorkspaceId> {
|
||||||
|
self.workspace_id.as_ref()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn client(&self) -> &WorkspaceClient {
|
||||||
|
&self.client
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// `(SessionId, SegmentId)` pair the Worker is currently writing to.
|
/// `(SessionId, SegmentId)` pair the Worker is currently writing to.
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
pub struct SegmentLocation {
|
pub struct SegmentLocation {
|
||||||
|
|
@ -109,10 +218,11 @@ where
|
||||||
let store = store.clone();
|
let store = store.clone();
|
||||||
Arc::new(move |metadata| {
|
Arc::new(move |metadata| {
|
||||||
store
|
store
|
||||||
.set_active_with_workspace_root(
|
.set_active_with_workspace_context(
|
||||||
&metadata.worker_name,
|
&metadata.worker_name,
|
||||||
metadata.active,
|
metadata.active,
|
||||||
metadata.resolved_manifest_snapshot,
|
metadata.resolved_manifest_snapshot,
|
||||||
|
metadata.workspace_id,
|
||||||
metadata.workspace_root,
|
metadata.workspace_root,
|
||||||
)
|
)
|
||||||
.map(|_| ())
|
.map(|_| ())
|
||||||
|
|
@ -287,9 +397,9 @@ pub struct Worker<C: LlmClient, St: Store> {
|
||||||
/// Explicit local filesystem authority, or `None` for Workers with no
|
/// Explicit local filesystem authority, or `None` for Workers with no
|
||||||
/// local cwd and no filesystem/Bash tool surface.
|
/// local cwd and no filesystem/Bash tool surface.
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
/// Absolute runtime workspace root used for project records, workflow,
|
/// Path-free workspace identity/client context injected by Runtime/host.
|
||||||
/// memory, Ticket config, Profile context, and spawned-child inheritance.
|
/// This never grants local filesystem authority.
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
/// Shared, atomically-swappable view of the Worker's resolved scope.
|
/// Shared, atomically-swappable view of the Worker's resolved scope.
|
||||||
/// Cloned out to `ScopedFs` instances (builtin tools, fs_view,
|
/// Cloned out to `ScopedFs` instances (builtin tools, fs_view,
|
||||||
/// compact worker) so scope updates propagate to every consumer
|
/// compact worker) so scope updates propagate to every consumer
|
||||||
|
|
@ -483,7 +593,7 @@ impl<C: LlmClient + Clone + 'static, St: Store + Clone + 'static> Worker<C, St>
|
||||||
worker_metadata_writer: None,
|
worker_metadata_writer: None,
|
||||||
segment_state: self.segment_state.clone(),
|
segment_state: self.segment_state.clone(),
|
||||||
filesystem_authority: self.filesystem_authority.clone(),
|
filesystem_authority: self.filesystem_authority.clone(),
|
||||||
workspace_root: self.workspace_root.clone(),
|
workspace_context: self.workspace_context.clone(),
|
||||||
scope: self.scope.clone(),
|
scope: self.scope.clone(),
|
||||||
delegation_scope: self.delegation_scope.clone(),
|
delegation_scope: self.delegation_scope.clone(),
|
||||||
hook_builder: HookRegistryBuilder::new(),
|
hook_builder: HookRegistryBuilder::new(),
|
||||||
|
|
@ -643,10 +753,10 @@ impl<C: LlmClient + Clone + 'static, St: Store + Clone + 'static> Worker<C, St>
|
||||||
impl<C: LlmClient, St: Store> Worker<C, St> {
|
impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
/// Create a new Worker from a pre-built Engine and store.
|
/// Create a new Worker from a pre-built Engine and store.
|
||||||
///
|
///
|
||||||
/// Callers must pass explicit filesystem authority and build a [`Scope`]
|
/// Callers must pass path-free workspace context separately from explicit
|
||||||
/// — typically via [`Scope::from_config`] when coming from a
|
/// filesystem authority and build a [`Scope`] — typically via
|
||||||
/// manifest, or [`Scope::writable`] in tests. Use
|
/// [`Scope::from_config`] when coming from a manifest, or [`Scope::writable`]
|
||||||
/// [`WorkerFilesystemAuthority::None`] for no-workdir Workers.
|
/// in tests. Use [`WorkerFilesystemAuthority::None`] for no-workdir Workers.
|
||||||
///
|
///
|
||||||
/// Note: this constructor does **not** parse `manifest.worker.system_prompt`
|
/// Note: this constructor does **not** parse `manifest.worker.system_prompt`
|
||||||
/// as a template. `Worker::from_manifest` is the production path for
|
/// as a template. `Worker::from_manifest` is the production path for
|
||||||
|
|
@ -656,7 +766,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
manifest: WorkerManifest,
|
manifest: WorkerManifest,
|
||||||
worker: Engine<C>,
|
worker: Engine<C>,
|
||||||
store: St,
|
store: St,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
scope: Scope,
|
scope: Scope,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
|
|
@ -675,7 +785,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
worker_metadata_writer: None,
|
worker_metadata_writer: None,
|
||||||
segment_state: SegmentState::new(session_id, segment_id, 0),
|
segment_state: SegmentState::new(session_id, segment_id, 0),
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
scope: SharedScope::new(scope),
|
scope: SharedScope::new(scope),
|
||||||
delegation_scope,
|
delegation_scope,
|
||||||
hook_builder: HookRegistryBuilder::new(),
|
hook_builder: HookRegistryBuilder::new(),
|
||||||
|
|
@ -797,10 +907,16 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
self.filesystem_authority.as_local()
|
self.filesystem_authority.as_local()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The Worker's runtime workspace root. This stays separate from `cwd` for
|
/// Path-free workspace identity, if Runtime/host associated this Worker
|
||||||
/// spawned children whose SpawnWorker `cwd` only changes tool defaults.
|
/// with a workspace.
|
||||||
pub fn workspace_root(&self) -> &Path {
|
pub fn workspace_id(&self) -> Option<&WorkspaceId> {
|
||||||
&self.workspace_root
|
self.workspace_context.workspace_id()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Narrow workspace client/availability handle injected by Runtime/host.
|
||||||
|
/// This never grants local filesystem authority.
|
||||||
|
pub fn workspace_client(&self) -> &WorkspaceClient {
|
||||||
|
self.workspace_context.client()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn worker_metadata_store(&self) -> St
|
pub(crate) fn worker_metadata_store(&self) -> St
|
||||||
|
|
@ -988,7 +1104,14 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn worker_metadata(&self, active: Option<WorkerActiveSegmentRef>) -> WorkerMetadata {
|
fn worker_metadata(&self, active: Option<WorkerActiveSegmentRef>) -> WorkerMetadata {
|
||||||
worker_metadata_for_manifest(&self.manifest, &self.workspace_root, active)
|
worker_metadata_for_manifest(
|
||||||
|
&self.manifest,
|
||||||
|
self.workspace_id(),
|
||||||
|
self.filesystem_authority
|
||||||
|
.as_local()
|
||||||
|
.map(|local| local.root.as_path()),
|
||||||
|
active,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_worker_metadata_pending(&self) -> Result<(), WorkerError> {
|
fn write_worker_metadata_pending(&self) -> Result<(), WorkerError> {
|
||||||
|
|
@ -1363,10 +1486,15 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
.map(|d| d.name)
|
.map(|d| d.name)
|
||||||
.collect()
|
.collect()
|
||||||
};
|
};
|
||||||
let agents_md_read = read_agents_md(&self.workspace_root);
|
let agents_md_read = self
|
||||||
for warning in agents_md_read.warnings {
|
.filesystem_authority
|
||||||
if let Some(n) = alerter.as_ref() {
|
.as_local()
|
||||||
n.alert(AlertLevel::Warn, AlertSource::AgentsMd, warning);
|
.map(|local| read_agents_md(&local.root));
|
||||||
|
if let Some(read) = agents_md_read.as_ref() {
|
||||||
|
for warning in &read.warnings {
|
||||||
|
if let Some(n) = alerter.as_ref() {
|
||||||
|
n.alert(AlertLevel::Warn, AlertSource::AgentsMd, warning.clone());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Resident-injection collection. Each resident section has its own
|
// Resident-injection collection. Each resident section has its own
|
||||||
|
|
@ -1429,7 +1557,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
language: worker_language,
|
language: worker_language,
|
||||||
scope: &scope_snapshot,
|
scope: &scope_snapshot,
|
||||||
tool_names,
|
tool_names,
|
||||||
agents_md: agents_md_read.body,
|
agents_md: agents_md_read.and_then(|read| read.body),
|
||||||
resident_summary: resident_summary.as_deref(),
|
resident_summary: resident_summary.as_deref(),
|
||||||
resident_knowledge: resident_slice,
|
resident_knowledge: resident_slice,
|
||||||
resident_workflows: resident_workflow_slice,
|
resident_workflows: resident_workflow_slice,
|
||||||
|
|
@ -2963,6 +3091,15 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
total_now.saturating_sub(total_at_pointer)
|
total_now.saturating_sub(total_at_pointer)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn local_memory_layout(
|
||||||
|
&self,
|
||||||
|
memory_cfg: &manifest::MemoryConfig,
|
||||||
|
) -> Option<memory::WorkspaceLayout> {
|
||||||
|
self.filesystem_authority
|
||||||
|
.as_local()
|
||||||
|
.map(|local| memory::WorkspaceLayout::resolve(memory_cfg, &local.root))
|
||||||
|
}
|
||||||
|
|
||||||
/// extract (memory.extract) post-run trigger.
|
/// extract (memory.extract) post-run trigger.
|
||||||
///
|
///
|
||||||
/// Called by the Controller before spawning the background memory task so
|
/// Called by the Controller before spawning the background memory task so
|
||||||
|
|
@ -2979,10 +3116,13 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
let Some(memory_cfg) = self.manifest.memory.clone() else {
|
let Some(memory_cfg) = self.manifest.memory.clone() else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
|
let Some(layout) = self.local_memory_layout(&memory_cfg) else {
|
||||||
|
tracing::debug!("workspace memory extract unavailable: no local filesystem authority");
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
// `Some(0)` means disabled, same as `None`. Otherwise the
|
// `Some(0)` means disabled, same as `None`. Otherwise the
|
||||||
// `tokens_since >= 0` comparison would fire on every post-run.
|
// `tokens_since >= 0` comparison would fire on every post-run.
|
||||||
let Some(threshold) = memory_cfg.extract_threshold.filter(|n| *n > 0) else {
|
let Some(threshold) = memory_cfg.extract_threshold.filter(|n| *n > 0) else {
|
||||||
let layout = memory::WorkspaceLayout::resolve(&memory_cfg, &self.workspace_root);
|
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.extract_model
|
.extract_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3012,7 +3152,6 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
|
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
|
||||||
.is_err()
|
.is_err()
|
||||||
{
|
{
|
||||||
let layout = memory::WorkspaceLayout::resolve(&memory_cfg, &self.workspace_root);
|
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.extract_model
|
.extract_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3068,7 +3207,10 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
) -> Result<ExtractDecision, WorkerError> {
|
) -> Result<ExtractDecision, WorkerError> {
|
||||||
use memory::extract;
|
use memory::extract;
|
||||||
|
|
||||||
let layout = memory::WorkspaceLayout::resolve(memory_cfg, &self.workspace_root);
|
let Some(layout) = self.local_memory_layout(memory_cfg) else {
|
||||||
|
tracing::debug!("workspace memory extract unavailable: no local filesystem authority");
|
||||||
|
return Ok(ExtractDecision::Skipped);
|
||||||
|
};
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.extract_model
|
.extract_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3383,6 +3525,12 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
let Some(memory_cfg) = self.manifest.memory.clone() else {
|
let Some(memory_cfg) = self.manifest.memory.clone() else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
|
let Some(layout) = self.local_memory_layout(&memory_cfg) else {
|
||||||
|
tracing::debug!(
|
||||||
|
"workspace memory consolidation unavailable: no local filesystem authority"
|
||||||
|
);
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
// `Some(0)` collapses to `None` — staging count / bytes always
|
// `Some(0)` collapses to `None` — staging count / bytes always
|
||||||
// satisfies `>= 0`, which would fire consolidation on every post-run.
|
// satisfies `>= 0`, which would fire consolidation on every post-run.
|
||||||
// Treating zero as disabled lines up with `extract_threshold` and
|
// Treating zero as disabled lines up with `extract_threshold` and
|
||||||
|
|
@ -3391,7 +3539,6 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
let files_threshold = memory_cfg.consolidation_threshold_files.filter(|n| *n > 0);
|
let files_threshold = memory_cfg.consolidation_threshold_files.filter(|n| *n > 0);
|
||||||
let bytes_threshold = memory_cfg.consolidation_threshold_bytes.filter(|n| *n > 0);
|
let bytes_threshold = memory_cfg.consolidation_threshold_bytes.filter(|n| *n > 0);
|
||||||
if files_threshold.is_none() && bytes_threshold.is_none() {
|
if files_threshold.is_none() && bytes_threshold.is_none() {
|
||||||
let layout = memory::WorkspaceLayout::resolve(&memory_cfg, &self.workspace_root);
|
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.consolidation_model
|
.consolidation_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3419,7 +3566,6 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
|
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
|
||||||
.is_err()
|
.is_err()
|
||||||
{
|
{
|
||||||
let layout = memory::WorkspaceLayout::resolve(&memory_cfg, &self.workspace_root);
|
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.consolidation_model
|
.consolidation_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3472,7 +3618,12 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||||
) -> Result<ConsolidateDecision, WorkerError> {
|
) -> Result<ConsolidateDecision, WorkerError> {
|
||||||
use memory::consolidate;
|
use memory::consolidate;
|
||||||
|
|
||||||
let layout = memory::WorkspaceLayout::resolve(memory_cfg, &self.workspace_root);
|
let Some(layout) = self.local_memory_layout(memory_cfg) else {
|
||||||
|
tracing::debug!(
|
||||||
|
"workspace memory consolidation unavailable: no local filesystem authority"
|
||||||
|
);
|
||||||
|
return Ok(ConsolidateDecision::Skipped);
|
||||||
|
};
|
||||||
let model = memory_cfg
|
let model = memory_cfg
|
||||||
.consolidation_model
|
.consolidation_model
|
||||||
.as_ref()
|
.as_ref()
|
||||||
|
|
@ -3892,21 +4043,23 @@ where
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let cwd = current_cwd()?;
|
let cwd = current_cwd()?;
|
||||||
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
||||||
Self::from_manifest_with_context(manifest, store, loader, cwd, authority).await
|
let workspace_context = WorkerWorkspaceContext::local_filesystem(None);
|
||||||
|
Self::from_manifest_with_context(manifest, store, loader, workspace_context, authority)
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn from_manifest_with_context(
|
pub async fn from_manifest_with_context(
|
||||||
manifest: WorkerManifest,
|
manifest: WorkerManifest,
|
||||||
store: St,
|
store: St,
|
||||||
loader: PromptLoader,
|
loader: PromptLoader,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let mut common = prepare_worker_common_with_context(
|
let mut common = prepare_worker_common_with_context(
|
||||||
&manifest,
|
&manifest,
|
||||||
&loader,
|
&loader,
|
||||||
/* parse_template */ true,
|
/* parse_template */ true,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
manifest.scope.clone(),
|
manifest.scope.clone(),
|
||||||
)?;
|
)?;
|
||||||
|
|
@ -3947,7 +4100,7 @@ where
|
||||||
worker_metadata_writer,
|
worker_metadata_writer,
|
||||||
segment_state: SegmentState::new(session_id, segment_id, 0),
|
segment_state: SegmentState::new(session_id, segment_id, 0),
|
||||||
filesystem_authority: common.filesystem_authority,
|
filesystem_authority: common.filesystem_authority,
|
||||||
workspace_root: common.workspace_root,
|
workspace_context: common.workspace_context,
|
||||||
scope: SharedScope::new(common.scope),
|
scope: SharedScope::new(common.scope),
|
||||||
delegation_scope: common.delegation_scope,
|
delegation_scope: common.delegation_scope,
|
||||||
hook_builder: HookRegistryBuilder::new(),
|
hook_builder: HookRegistryBuilder::new(),
|
||||||
|
|
@ -4007,12 +4160,13 @@ where
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let cwd = current_cwd()?;
|
let cwd = current_cwd()?;
|
||||||
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
||||||
|
let workspace_context = WorkerWorkspaceContext::local_filesystem(None);
|
||||||
Self::from_manifest_spawned_with_context(
|
Self::from_manifest_spawned_with_context(
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
callback_socket,
|
callback_socket,
|
||||||
cwd,
|
workspace_context,
|
||||||
authority,
|
authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -4023,14 +4177,14 @@ where
|
||||||
store: St,
|
store: St,
|
||||||
loader: PromptLoader,
|
loader: PromptLoader,
|
||||||
callback_socket: PathBuf,
|
callback_socket: PathBuf,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let mut common = prepare_worker_common_with_context(
|
let mut common = prepare_worker_common_with_context(
|
||||||
&manifest,
|
&manifest,
|
||||||
&loader,
|
&loader,
|
||||||
/* parse_template */ true,
|
/* parse_template */ true,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
manifest.scope.clone(),
|
manifest.scope.clone(),
|
||||||
)?;
|
)?;
|
||||||
|
|
@ -4058,7 +4212,7 @@ where
|
||||||
worker_metadata_writer,
|
worker_metadata_writer,
|
||||||
segment_state: SegmentState::new(session_id, segment_id, 0),
|
segment_state: SegmentState::new(session_id, segment_id, 0),
|
||||||
filesystem_authority: common.filesystem_authority,
|
filesystem_authority: common.filesystem_authority,
|
||||||
workspace_root: common.workspace_root,
|
workspace_context: common.workspace_context,
|
||||||
scope: SharedScope::new(common.scope),
|
scope: SharedScope::new(common.scope),
|
||||||
delegation_scope: common.delegation_scope,
|
delegation_scope: common.delegation_scope,
|
||||||
hook_builder: HookRegistryBuilder::new(),
|
hook_builder: HookRegistryBuilder::new(),
|
||||||
|
|
@ -4114,12 +4268,13 @@ where
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let cwd = current_cwd()?;
|
let cwd = current_cwd()?;
|
||||||
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
||||||
|
let workspace_context = WorkerWorkspaceContext::local_filesystem(None);
|
||||||
Self::restore_from_worker_metadata_with_context(
|
Self::restore_from_worker_metadata_with_context(
|
||||||
worker_name,
|
worker_name,
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
cwd,
|
workspace_context,
|
||||||
authority,
|
authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -4130,7 +4285,7 @@ where
|
||||||
manifest: WorkerManifest,
|
manifest: WorkerManifest,
|
||||||
store: St,
|
store: St,
|
||||||
loader: PromptLoader,
|
loader: PromptLoader,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let metadata =
|
let metadata =
|
||||||
|
|
@ -4161,7 +4316,7 @@ where
|
||||||
manifest,
|
manifest,
|
||||||
store,
|
store,
|
||||||
loader,
|
loader,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -4193,8 +4348,15 @@ where
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
let cwd = current_cwd()?;
|
let cwd = current_cwd()?;
|
||||||
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone());
|
||||||
|
let workspace_context = WorkerWorkspaceContext::local_filesystem(None);
|
||||||
Self::restore_from_manifest_with_context(
|
Self::restore_from_manifest_with_context(
|
||||||
session_id, segment_id, manifest, store, loader, cwd, authority,
|
session_id,
|
||||||
|
segment_id,
|
||||||
|
manifest,
|
||||||
|
store,
|
||||||
|
loader,
|
||||||
|
workspace_context,
|
||||||
|
authority,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
@ -4205,7 +4367,7 @@ where
|
||||||
manifest: WorkerManifest,
|
manifest: WorkerManifest,
|
||||||
store: St,
|
store: St,
|
||||||
loader: PromptLoader,
|
loader: PromptLoader,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
// Read raw entries once so we can both reconstruct state and
|
// Read raw entries once so we can both reconstruct state and
|
||||||
|
|
@ -4223,7 +4385,7 @@ where
|
||||||
&manifest,
|
&manifest,
|
||||||
&loader,
|
&loader,
|
||||||
/* parse_template */ false,
|
/* parse_template */ false,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
scope_config,
|
scope_config,
|
||||||
)?;
|
)?;
|
||||||
|
|
@ -4289,7 +4451,7 @@ where
|
||||||
worker_metadata_writer,
|
worker_metadata_writer,
|
||||||
segment_state: SegmentState::new(session_id, segment_id, state.entries_count),
|
segment_state: SegmentState::new(session_id, segment_id, state.entries_count),
|
||||||
filesystem_authority: common.filesystem_authority,
|
filesystem_authority: common.filesystem_authority,
|
||||||
workspace_root: common.workspace_root,
|
workspace_context: common.workspace_context,
|
||||||
scope: SharedScope::new(common.scope),
|
scope: SharedScope::new(common.scope),
|
||||||
delegation_scope: common.delegation_scope,
|
delegation_scope: common.delegation_scope,
|
||||||
hook_builder: HookRegistryBuilder::new(),
|
hook_builder: HookRegistryBuilder::new(),
|
||||||
|
|
@ -4441,11 +4603,17 @@ fn request_config_from_engine_manifest(wm: &manifest::EngineManifest) -> Request
|
||||||
|
|
||||||
fn worker_metadata_for_manifest(
|
fn worker_metadata_for_manifest(
|
||||||
manifest: &WorkerManifest,
|
manifest: &WorkerManifest,
|
||||||
workspace_root: &Path,
|
workspace_id: Option<&WorkspaceId>,
|
||||||
|
local_workspace_root: Option<&std::path::Path>,
|
||||||
active: Option<WorkerActiveSegmentRef>,
|
active: Option<WorkerActiveSegmentRef>,
|
||||||
) -> WorkerMetadata {
|
) -> WorkerMetadata {
|
||||||
let mut metadata = WorkerMetadata::new(manifest.worker.name.clone(), active)
|
let mut metadata = WorkerMetadata::new(manifest.worker.name.clone(), active);
|
||||||
.with_workspace_root(workspace_root.to_path_buf());
|
if let Some(workspace_id) = workspace_id {
|
||||||
|
metadata = metadata.with_workspace_id(workspace_id.as_str().to_owned());
|
||||||
|
}
|
||||||
|
if let Some(local_workspace_root) = local_workspace_root {
|
||||||
|
metadata = metadata.with_workspace_root(local_workspace_root.to_path_buf());
|
||||||
|
}
|
||||||
if should_persist_resolved_manifest_snapshot(manifest) {
|
if should_persist_resolved_manifest_snapshot(manifest) {
|
||||||
metadata.resolved_manifest_snapshot = serde_json::to_value(manifest).ok();
|
metadata.resolved_manifest_snapshot = serde_json::to_value(manifest).ok();
|
||||||
}
|
}
|
||||||
|
|
@ -4875,15 +5043,15 @@ pub enum WorkerError {
|
||||||
#[error(transparent)]
|
#[error(transparent)]
|
||||||
Scope(ScopeError),
|
Scope(ScopeError),
|
||||||
|
|
||||||
#[error("workspace root is not readable under the configured scope: {}", .workspace_root.display())]
|
#[error("local filesystem authority root is not readable under the configured scope: {}", .root.display())]
|
||||||
WorkspaceRootOutsideScope { workspace_root: PathBuf },
|
LocalFilesystemRootOutsideScope { root: PathBuf },
|
||||||
|
|
||||||
#[error("cwd is not readable under the configured scope: {}", .cwd.display())]
|
#[error("cwd is not readable under the configured scope: {}", .cwd.display())]
|
||||||
CwdOutsideScope { cwd: PathBuf },
|
CwdOutsideScope { cwd: PathBuf },
|
||||||
|
|
||||||
#[error("failed to resolve workspace root {}: {source}", .workspace_root.display())]
|
#[error("failed to resolve local filesystem authority root {}: {source}", .root.display())]
|
||||||
InvalidWorkspaceRoot {
|
InvalidLocalFilesystemRoot {
|
||||||
workspace_root: PathBuf,
|
root: PathBuf,
|
||||||
#[source]
|
#[source]
|
||||||
source: std::io::Error,
|
source: std::io::Error,
|
||||||
},
|
},
|
||||||
|
|
@ -4974,13 +5142,13 @@ pub enum WorkerError {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Bundle of resources that every high-level Worker constructor needs:
|
/// Bundle of resources that every high-level Worker constructor needs:
|
||||||
/// cwd, runtime workspace root, scope, an LLM client, the prompt catalog,
|
/// filesystem authority, path-free workspace context, scope, an LLM client, the prompt catalog,
|
||||||
/// and (optionally) a parsed system-prompt template. Built once by
|
/// and (optionally) a parsed system-prompt template. Built once by
|
||||||
/// [`prepare_worker_common_with_context`] from the resolved manifest and then split into Worker
|
/// [`prepare_worker_common_with_context`] from the resolved manifest and then split into Worker
|
||||||
/// fields.
|
/// fields.
|
||||||
struct WorkerCommon {
|
struct WorkerCommon {
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
scope: Scope,
|
scope: Scope,
|
||||||
delegation_scope: DelegationScope,
|
delegation_scope: DelegationScope,
|
||||||
client: Box<dyn LlmClient>,
|
client: Box<dyn LlmClient>,
|
||||||
|
|
@ -5067,22 +5235,16 @@ fn prepare_worker_common_with_context(
|
||||||
manifest: &WorkerManifest,
|
manifest: &WorkerManifest,
|
||||||
loader: &PromptLoader,
|
loader: &PromptLoader,
|
||||||
parse_template: bool,
|
parse_template: bool,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
scope_config: ScopeConfig,
|
scope_config: ScopeConfig,
|
||||||
) -> Result<WorkerCommon, WorkerError> {
|
) -> Result<WorkerCommon, WorkerError> {
|
||||||
let workspace_root = std::fs::canonicalize(&workspace_root).map_err(|source| {
|
|
||||||
WorkerError::InvalidWorkspaceRoot {
|
|
||||||
workspace_root: workspace_root.clone(),
|
|
||||||
source,
|
|
||||||
}
|
|
||||||
})?;
|
|
||||||
let filesystem_authority = match filesystem_authority {
|
let filesystem_authority = match filesystem_authority {
|
||||||
WorkerFilesystemAuthority::None => WorkerFilesystemAuthority::None,
|
WorkerFilesystemAuthority::None => WorkerFilesystemAuthority::None,
|
||||||
WorkerFilesystemAuthority::Local(local) => {
|
WorkerFilesystemAuthority::Local(local) => {
|
||||||
let root = std::fs::canonicalize(&local.root).map_err(|source| {
|
let root = std::fs::canonicalize(&local.root).map_err(|source| {
|
||||||
WorkerError::InvalidWorkspaceRoot {
|
WorkerError::InvalidLocalFilesystemRoot {
|
||||||
workspace_root: local.root.clone(),
|
root: local.root.clone(),
|
||||||
source,
|
source,
|
||||||
}
|
}
|
||||||
})?;
|
})?;
|
||||||
|
|
@ -5095,8 +5257,8 @@ fn prepare_worker_common_with_context(
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let mut scope_config = scope_config;
|
let mut scope_config = scope_config;
|
||||||
if let Some(mem) = manifest.memory.as_ref() {
|
if let (Some(mem), Some(local)) = (manifest.memory.as_ref(), filesystem_authority.as_local()) {
|
||||||
let layout = memory::WorkspaceLayout::resolve(mem, &workspace_root);
|
let layout = memory::WorkspaceLayout::resolve(mem, &local.root);
|
||||||
scope_config.deny.extend(memory::deny_write_rules(&layout));
|
scope_config.deny.extend(memory::deny_write_rules(&layout));
|
||||||
scope_config
|
scope_config
|
||||||
.deny
|
.deny
|
||||||
|
|
@ -5108,7 +5270,7 @@ fn prepare_worker_common_with_context(
|
||||||
manifest,
|
manifest,
|
||||||
loader,
|
loader,
|
||||||
parse_template,
|
parse_template,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
@ -5118,14 +5280,16 @@ fn prepare_worker_common_from_scope(
|
||||||
manifest: &WorkerManifest,
|
manifest: &WorkerManifest,
|
||||||
loader: &PromptLoader,
|
loader: &PromptLoader,
|
||||||
parse_template: bool,
|
parse_template: bool,
|
||||||
workspace_root: PathBuf,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
scope: Scope,
|
scope: Scope,
|
||||||
) -> Result<WorkerCommon, WorkerError> {
|
) -> Result<WorkerCommon, WorkerError> {
|
||||||
if !scope.is_readable(&workspace_root) {
|
|
||||||
return Err(WorkerError::WorkspaceRootOutsideScope { workspace_root });
|
|
||||||
}
|
|
||||||
if let Some(local) = filesystem_authority.as_local() {
|
if let Some(local) = filesystem_authority.as_local() {
|
||||||
|
if !scope.is_readable(&local.root) {
|
||||||
|
return Err(WorkerError::LocalFilesystemRootOutsideScope {
|
||||||
|
root: local.root.clone(),
|
||||||
|
});
|
||||||
|
}
|
||||||
if !scope.is_readable(&local.cwd) {
|
if !scope.is_readable(&local.cwd) {
|
||||||
return Err(WorkerError::CwdOutsideScope {
|
return Err(WorkerError::CwdOutsideScope {
|
||||||
cwd: local.cwd.clone(),
|
cwd: local.cwd.clone(),
|
||||||
|
|
@ -5137,10 +5301,11 @@ fn prepare_worker_common_from_scope(
|
||||||
|
|
||||||
let client = provider::build_client(&manifest.model)?;
|
let client = provider::build_client(&manifest.model)?;
|
||||||
let prompts = PromptCatalog::load(loader, manifest.worker.prompt_pack.as_deref())?;
|
let prompts = PromptCatalog::load(loader, manifest.worker.prompt_pack.as_deref())?;
|
||||||
let memory_layout = manifest
|
let memory_layout = manifest.memory.as_ref().and_then(|mem| {
|
||||||
.memory
|
filesystem_authority
|
||||||
.as_ref()
|
.as_local()
|
||||||
.map(|mem| memory::WorkspaceLayout::resolve(mem, &workspace_root));
|
.map(|local| memory::WorkspaceLayout::resolve(mem, &local.root))
|
||||||
|
});
|
||||||
let mut workflow_registry = match memory_layout.as_ref() {
|
let mut workflow_registry = match memory_layout.as_ref() {
|
||||||
Some(layout) => {
|
Some(layout) => {
|
||||||
workflow_crate::load_workflows(layout).map_err(WorkerError::WorkflowLoad)?
|
workflow_crate::load_workflows(layout).map_err(WorkerError::WorkflowLoad)?
|
||||||
|
|
@ -5160,7 +5325,7 @@ fn prepare_worker_common_from_scope(
|
||||||
|
|
||||||
Ok(WorkerCommon {
|
Ok(WorkerCommon {
|
||||||
filesystem_authority,
|
filesystem_authority,
|
||||||
workspace_root,
|
workspace_context,
|
||||||
scope,
|
scope,
|
||||||
delegation_scope,
|
delegation_scope,
|
||||||
client,
|
client,
|
||||||
|
|
@ -5249,7 +5414,7 @@ mod spawned_context_tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn spawn_worker_context_keeps_workspace_root_separate_from_tool_pwd() {
|
fn spawn_worker_context_separates_workspace_identity_from_tool_pwd() {
|
||||||
let tmp = tempfile::tempdir().unwrap();
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
let workspace_root = tmp.path().join("workspace-root");
|
let workspace_root = tmp.path().join("workspace-root");
|
||||||
let cwd = tmp.path().join("child-worktree");
|
let cwd = tmp.path().join("child-worktree");
|
||||||
|
|
@ -5262,14 +5427,21 @@ mod spawned_context_tests {
|
||||||
&manifest,
|
&manifest,
|
||||||
&PromptLoader::builtins_only(),
|
&PromptLoader::builtins_only(),
|
||||||
false,
|
false,
|
||||||
workspace_root.clone(),
|
WorkerWorkspaceContext::local_filesystem(Some(WorkspaceId::new("ws-test").unwrap())),
|
||||||
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
||||||
manifest.scope.clone(),
|
manifest.scope.clone(),
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
common.workspace_root,
|
common
|
||||||
|
.workspace_context
|
||||||
|
.workspace_id()
|
||||||
|
.map(WorkspaceId::as_str),
|
||||||
|
Some("ws-test")
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
common.filesystem_authority.as_local().unwrap().root,
|
||||||
workspace_root.canonicalize().unwrap()
|
workspace_root.canonicalize().unwrap()
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|
@ -5283,7 +5455,42 @@ mod spawned_context_tests {
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn prepare_context_reports_workspace_root_when_workspace_root_is_unreadable() {
|
fn workspace_identity_and_client_do_not_grant_filesystem_authority() {
|
||||||
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
|
let workspace_root = tmp.path().join("workspace-root");
|
||||||
|
let cwd = workspace_root.join("nested");
|
||||||
|
std::fs::create_dir_all(&cwd).unwrap();
|
||||||
|
let mut manifest = minimal_manifest_for_context_test(&workspace_root, &cwd);
|
||||||
|
manifest.memory = Some(manifest::MemoryConfig::default());
|
||||||
|
let loader = PromptLoader::new(None, Some(workspace_root.clone()));
|
||||||
|
let workspace_id = WorkspaceId::new("ws-api-only").unwrap();
|
||||||
|
let common = prepare_worker_common_with_context(
|
||||||
|
&manifest,
|
||||||
|
&loader,
|
||||||
|
false,
|
||||||
|
WorkerWorkspaceContext::with_client(
|
||||||
|
Some(workspace_id.clone()),
|
||||||
|
WorkspaceClient::available("test-api"),
|
||||||
|
),
|
||||||
|
WorkerFilesystemAuthority::None,
|
||||||
|
manifest.scope.clone(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert_eq!(common.filesystem_authority, WorkerFilesystemAuthority::None);
|
||||||
|
assert_eq!(
|
||||||
|
common
|
||||||
|
.workspace_context
|
||||||
|
.workspace_id()
|
||||||
|
.map(WorkspaceId::as_str),
|
||||||
|
Some(workspace_id.as_str())
|
||||||
|
);
|
||||||
|
assert!(common.workspace_context.client().is_available());
|
||||||
|
assert!(common.memory_layout.is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn prepare_context_reports_local_filesystem_root_when_unreadable() {
|
||||||
let tmp = tempfile::tempdir().unwrap();
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
let workspace_root = tmp.path().join("workspace-root");
|
let workspace_root = tmp.path().join("workspace-root");
|
||||||
let cwd = tmp.path().join("child-worktree");
|
let cwd = tmp.path().join("child-worktree");
|
||||||
|
|
@ -5295,7 +5502,7 @@ mod spawned_context_tests {
|
||||||
&manifest,
|
&manifest,
|
||||||
&PromptLoader::builtins_only(),
|
&PromptLoader::builtins_only(),
|
||||||
false,
|
false,
|
||||||
workspace_root.clone(),
|
WorkerWorkspaceContext::local_filesystem(Some(WorkspaceId::new("ws-test").unwrap())),
|
||||||
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
||||||
ScopeConfig {
|
ScopeConfig {
|
||||||
allow: vec![ScopeRule {
|
allow: vec![ScopeRule {
|
||||||
|
|
@ -5306,17 +5513,15 @@ mod spawned_context_tests {
|
||||||
deny: Vec::new(),
|
deny: Vec::new(),
|
||||||
},
|
},
|
||||||
) {
|
) {
|
||||||
Ok(_) => panic!("expected workspace-root scope error"),
|
Ok(_) => panic!("expected local filesystem root scope error"),
|
||||||
Err(err) => err,
|
Err(err) => err,
|
||||||
};
|
};
|
||||||
|
|
||||||
match err {
|
match err {
|
||||||
WorkerError::WorkspaceRootOutsideScope {
|
WorkerError::LocalFilesystemRootOutsideScope { root: got } => {
|
||||||
workspace_root: got,
|
|
||||||
} => {
|
|
||||||
assert_eq!(got, workspace_root.canonicalize().unwrap());
|
assert_eq!(got, workspace_root.canonicalize().unwrap());
|
||||||
}
|
}
|
||||||
other => panic!("expected workspace-root scope error, got {other:?}"),
|
other => panic!("expected local filesystem root scope error, got {other:?}"),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -5333,7 +5538,7 @@ mod spawned_context_tests {
|
||||||
&manifest,
|
&manifest,
|
||||||
&PromptLoader::builtins_only(),
|
&PromptLoader::builtins_only(),
|
||||||
false,
|
false,
|
||||||
workspace_root.clone(),
|
WorkerWorkspaceContext::local_filesystem(Some(WorkspaceId::new("ws-test").unwrap())),
|
||||||
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()),
|
||||||
ScopeConfig {
|
ScopeConfig {
|
||||||
allow: vec![ScopeRule {
|
allow: vec![ScopeRule {
|
||||||
|
|
@ -5390,23 +5595,15 @@ mod worker_metadata_restore_manifest_tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn metadata_writer_persists_workspace_root_through_store_update() {
|
fn metadata_writer_persists_workspace_id_through_store_update() {
|
||||||
let temp = tempfile::tempdir().unwrap();
|
let temp = tempfile::tempdir().unwrap();
|
||||||
let store = session_store::FsWorkerStore::new(temp.path().join("workers")).unwrap();
|
let store = session_store::FsWorkerStore::new(temp.path().join("workers")).unwrap();
|
||||||
let workspace_root = temp.path().join("workspace-root");
|
|
||||||
std::fs::create_dir_all(&workspace_root).unwrap();
|
|
||||||
let writer = worker_metadata_writer_for_store(&store);
|
let writer = worker_metadata_writer_for_store(&store);
|
||||||
|
|
||||||
writer(
|
writer(WorkerMetadata::new("runtime-worker", None).with_workspace_id("ws-test")).unwrap();
|
||||||
WorkerMetadata::new("runtime-worker", None).with_workspace_root(workspace_root.clone()),
|
|
||||||
)
|
|
||||||
.unwrap();
|
|
||||||
|
|
||||||
let stored = store.read_by_name("runtime-worker").unwrap().unwrap();
|
let stored = store.read_by_name("runtime-worker").unwrap().unwrap();
|
||||||
assert_eq!(
|
assert_eq!(stored.workspace_id.as_deref(), Some("ws-test"));
|
||||||
stored.workspace_root.as_deref(),
|
|
||||||
Some(workspace_root.as_path())
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|
@ -5503,7 +5700,7 @@ permission = "read"
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert!(manifest.profile.is_none());
|
assert!(manifest.profile.is_none());
|
||||||
assert!(
|
assert!(
|
||||||
worker_metadata_for_manifest(&manifest, Path::new("/snapshot/workspace"), None)
|
worker_metadata_for_manifest(&manifest, None, None, None)
|
||||||
.resolved_manifest_snapshot
|
.resolved_manifest_snapshot
|
||||||
.is_none()
|
.is_none()
|
||||||
);
|
);
|
||||||
|
|
@ -5540,8 +5737,7 @@ permission = "read"
|
||||||
config: None,
|
config: None,
|
||||||
}];
|
}];
|
||||||
|
|
||||||
let metadata =
|
let metadata = worker_metadata_for_manifest(&manifest, None, None, None);
|
||||||
worker_metadata_for_manifest(&manifest, Path::new("/snapshot/workspace"), None);
|
|
||||||
let snapshot = metadata
|
let snapshot = metadata
|
||||||
.resolved_manifest_snapshot
|
.resolved_manifest_snapshot
|
||||||
.expect("plugin-resolved manifest should be snapshotted");
|
.expect("plugin-resolved manifest should be snapshotted");
|
||||||
|
|
@ -5800,7 +5996,7 @@ mod build_summary_prompt_tests {
|
||||||
manifest,
|
manifest,
|
||||||
Engine::new(NoopClient),
|
Engine::new(NoopClient),
|
||||||
store,
|
store,
|
||||||
cwd.clone(),
|
WorkerWorkspaceContext::local_filesystem(None),
|
||||||
authority,
|
authority,
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
@ -5956,7 +6152,7 @@ mod build_summary_prompt_tests {
|
||||||
manifest,
|
manifest,
|
||||||
Engine::new(NoopClient),
|
Engine::new(NoopClient),
|
||||||
store,
|
store,
|
||||||
cwd.clone(),
|
WorkerWorkspaceContext::local_filesystem(None),
|
||||||
authority,
|
authority,
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
@ -6091,7 +6287,7 @@ mod build_summary_prompt_tests {
|
||||||
manifest,
|
manifest,
|
||||||
Engine::new(NoopClient),
|
Engine::new(NoopClient),
|
||||||
store,
|
store,
|
||||||
cwd.clone(),
|
WorkerWorkspaceContext::local_filesystem(None),
|
||||||
authority,
|
authority,
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -168,7 +168,7 @@ async fn make_worker_with_manifest(
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store,
|
store,
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -174,7 +174,7 @@ async fn make_worker_with(
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store,
|
store,
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,7 @@ use session_store::{FsStore, LogEntry};
|
||||||
|
|
||||||
use worker::{
|
use worker::{
|
||||||
Event, Method, Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle,
|
Event, Method, Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle,
|
||||||
WorkerManifest, WorkerStatus,
|
WorkerManifest, WorkerStatus, WorkerWorkspaceContext,
|
||||||
};
|
};
|
||||||
|
|
||||||
type TestStore = CombinedStore<FsStore, FsWorkerStore>;
|
type TestStore = CombinedStore<FsStore, FsWorkerStore>;
|
||||||
|
|
@ -190,9 +190,16 @@ async fn make_worker_with_pwd_and_manifest(
|
||||||
|
|
||||||
let worker = Engine::new(client);
|
let worker = Engine::new(client);
|
||||||
let authority = WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone());
|
let authority = WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone());
|
||||||
let worker = Worker::new(manifest, worker, store, pwd.clone(), authority, scope)
|
let worker = Worker::new(
|
||||||
.await
|
manifest,
|
||||||
.unwrap();
|
worker,
|
||||||
|
store,
|
||||||
|
WorkerWorkspaceContext::local_filesystem(None),
|
||||||
|
authority,
|
||||||
|
scope,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
(worker, pwd)
|
(worker, pwd)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -194,7 +194,7 @@ async fn make_worker(
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store,
|
store,
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
@ -464,7 +464,7 @@ async fn metric_write_failure_emits_warn_alert_and_does_not_abort_run() {
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store.clone(),
|
store.clone(),
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
@ -540,7 +540,7 @@ permission = "write"
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store.clone(),
|
store.clone(),
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -127,7 +127,7 @@ async fn make_worker_with_body(
|
||||||
manifest,
|
manifest,
|
||||||
worker,
|
worker,
|
||||||
store,
|
store,
|
||||||
pwd.clone(),
|
worker::WorkerWorkspaceContext::local_filesystem(None),
|
||||||
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()),
|
||||||
scope,
|
scope,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -43,7 +43,7 @@ rustPlatform.buildRustPackage rec {
|
||||||
filter = sourceFilter;
|
filter = sourceFilter;
|
||||||
};
|
};
|
||||||
|
|
||||||
cargoHash = "sha256-vM2ea6W3hqnmzqOXfzQo+ZYrQBoy2Kh69vz1fkU+DIs=";
|
cargoHash = "sha256-av+Iix50MMLcrzxw8UDTT0tUqBFIPkSN1i23WYegNcM=";
|
||||||
|
|
||||||
depsExtraArgs = {
|
depsExtraArgs = {
|
||||||
# Older fetchCargoVendor utilities used crates.io's API download endpoint,
|
# Older fetchCargoVendor utilities used crates.io's API download endpoint,
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user