use std::sync::Arc;
use serde::{Deserialize, Serialize};
use crate::fs::VirtualFileSystem;
#[derive(Default)]
pub struct AgentOsConfig {
pub software: Vec<SoftwareInput>,
pub packages: Vec<PackageRef>,
pub packages_mount_at: Option<String>,
pub loopback_exempt_ports: Vec<u16>,
pub allowed_node_builtins: Option<Vec<String>>,
pub root_filesystem: RootFilesystemConfig,
pub mounts: Vec<MountConfig>,
pub additional_instructions: Option<String>,
pub schedule_driver: Option<Arc<dyn ScheduleDriver>>,
pub tool_kits: Vec<ToolKit>,
pub sidecar_js_bridge_callback: Option<SidecarJsBridgeCallback>,
pub permissions: Option<Permissions>,
pub limits: Option<AgentOsLimits>,
pub sidecar: Option<AgentOsSidecarConfig>,
pub sidecar_binary_path: Option<String>,
}
#[derive(Default)]
pub struct AgentOsConfigBuilder {
config: AgentOsConfig,
}
impl AgentOsConfigBuilder {
pub fn new() -> Self {
Self::default()
}
pub fn packages(mut self, packages: Vec<PackageRef>) -> Self {
self.config.packages = packages;
self
}
pub fn packages_mount_at(mut self, mount_at: impl Into<String>) -> Self {
self.config.packages_mount_at = Some(mount_at.into());
self
}
pub fn loopback_exempt_ports(mut self, ports: Vec<u16>) -> Self {
self.config.loopback_exempt_ports = ports;
self
}
pub fn allowed_node_builtins(mut self, builtins: Vec<String>) -> Self {
self.config.allowed_node_builtins = Some(builtins);
self
}
pub fn root_filesystem(mut self, root: RootFilesystemConfig) -> Self {
self.config.root_filesystem = root;
self
}
pub fn mounts(mut self, mounts: Vec<MountConfig>) -> Self {
self.config.mounts = mounts;
self
}
pub fn additional_instructions(mut self, instructions: impl Into<String>) -> Self {
self.config.additional_instructions = Some(instructions.into());
self
}
pub fn schedule_driver(mut self, driver: Arc<dyn ScheduleDriver>) -> Self {
self.config.schedule_driver = Some(driver);
self
}
pub fn tool_kits(mut self, tool_kits: Vec<ToolKit>) -> Self {
self.config.tool_kits = tool_kits;
self
}
pub fn sidecar_js_bridge_callback(mut self, callback: SidecarJsBridgeCallback) -> Self {
self.config.sidecar_js_bridge_callback = Some(callback);
self
}
pub fn permissions(mut self, permissions: Permissions) -> Self {
self.config.permissions = Some(permissions);
self
}
pub fn limits(mut self, limits: AgentOsLimits) -> Self {
self.config.limits = Some(limits);
self
}
pub fn sidecar(mut self, sidecar: AgentOsSidecarConfig) -> Self {
self.config.sidecar = Some(sidecar);
self
}
pub fn sidecar_binary_path(mut self, path: impl Into<String>) -> Self {
self.config.sidecar_binary_path = Some(path.into());
self
}
pub fn build(self) -> AgentOsConfig {
self.config
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum SoftwareKind {
#[default]
WasmCommands,
Agent,
Tool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SoftwareInput {
pub package: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
#[serde(default)]
pub kind: SoftwareKind,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PackageRef {
#[serde(rename = "packagePath")]
pub path: String,
}
pub type ToolCallback = Arc<
dyn Fn(
serde_json::Value,
) -> futures::future::BoxFuture<'static, Result<serde_json::Value, String>>
+ Send
+ Sync,
>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SidecarJsBridgeCall {
pub call_id: String,
pub mount_id: String,
pub operation: String,
pub args: serde_json::Value,
}
pub type SidecarJsBridgeCallback = Arc<
dyn Fn(
SidecarJsBridgeCall,
)
-> futures::future::BoxFuture<'static, Result<Option<serde_json::Value>, String>>
+ Send
+ Sync,
>;
#[derive(Clone)]
pub struct HostTool {
pub name: String,
pub description: String,
pub input_schema: serde_json::Value,
pub timeout_ms: Option<u64>,
pub execute: ToolCallback,
}
#[derive(Clone)]
pub struct ToolKit {
pub name: String,
pub description: String,
pub tools: Vec<HostTool>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct AgentOsLimits {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resources: Option<ResourceLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub http: Option<HttpLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tools: Option<ToolLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plugins: Option<PluginLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub acp: Option<AcpLimits>,
#[serde(default, rename = "jsRuntime", skip_serializing_if = "Option::is_none")]
pub js_runtime: Option<JsRuntimeLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub python: Option<PythonLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub wasm: Option<WasmLimits>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ResourceLimits {
#[serde(default, rename = "cpuCount", skip_serializing_if = "Option::is_none")]
pub cpu_count: Option<u64>,
#[serde(
default,
rename = "maxProcesses",
skip_serializing_if = "Option::is_none"
)]
pub max_processes: Option<u64>,
#[serde(
default,
rename = "maxOpenFds",
skip_serializing_if = "Option::is_none"
)]
pub max_open_fds: Option<u64>,
#[serde(default, rename = "maxPipes", skip_serializing_if = "Option::is_none")]
pub max_pipes: Option<u64>,
#[serde(default, rename = "maxPtys", skip_serializing_if = "Option::is_none")]
pub max_ptys: Option<u64>,
#[serde(
default,
rename = "maxSockets",
skip_serializing_if = "Option::is_none"
)]
pub max_sockets: Option<u64>,
#[serde(
default,
rename = "maxConnections",
skip_serializing_if = "Option::is_none"
)]
pub max_connections: Option<u64>,
#[serde(
default,
rename = "maxSocketBufferedBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_socket_buffered_bytes: Option<u64>,
#[serde(
default,
rename = "maxSocketDatagramQueueLen",
skip_serializing_if = "Option::is_none"
)]
pub max_socket_datagram_queue_len: Option<u64>,
#[serde(
default,
rename = "maxFilesystemBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_filesystem_bytes: Option<u64>,
#[serde(
default,
rename = "maxInodeCount",
skip_serializing_if = "Option::is_none"
)]
pub max_inode_count: Option<u64>,
#[serde(
default,
rename = "maxBlockingReadMs",
skip_serializing_if = "Option::is_none"
)]
pub max_blocking_read_ms: Option<u64>,
#[serde(
default,
rename = "maxPreadBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_pread_bytes: Option<u64>,
#[serde(
default,
rename = "maxFdWriteBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_fd_write_bytes: Option<u64>,
#[serde(
default,
rename = "maxProcessArgvBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_process_argv_bytes: Option<u64>,
#[serde(
default,
rename = "maxProcessEnvBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_process_env_bytes: Option<u64>,
#[serde(
default,
rename = "maxReaddirEntries",
skip_serializing_if = "Option::is_none"
)]
pub max_readdir_entries: Option<u64>,
#[serde(
default,
rename = "maxWasmFuel",
skip_serializing_if = "Option::is_none"
)]
pub max_wasm_fuel: Option<u64>,
#[serde(
default,
rename = "maxWasmMemoryBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_wasm_memory_bytes: Option<u64>,
#[serde(
default,
rename = "maxWasmStackBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_wasm_stack_bytes: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct HttpLimits {
#[serde(
default,
rename = "maxFetchResponseBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_fetch_response_bytes: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ToolLimits {
#[serde(
default,
rename = "defaultToolTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub default_tool_timeout_ms: Option<u64>,
#[serde(
default,
rename = "maxToolTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub max_tool_timeout_ms: Option<u64>,
#[serde(
default,
rename = "maxRegisteredToolkits",
skip_serializing_if = "Option::is_none"
)]
pub max_registered_toolkits: Option<u64>,
#[serde(
default,
rename = "maxRegisteredToolsPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_registered_tools_per_vm: Option<u64>,
#[serde(
default,
rename = "maxToolsPerToolkit",
skip_serializing_if = "Option::is_none"
)]
pub max_tools_per_toolkit: Option<u64>,
#[serde(
default,
rename = "maxToolSchemaBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_tool_schema_bytes: Option<u64>,
#[serde(
default,
rename = "maxToolExamplesPerTool",
skip_serializing_if = "Option::is_none"
)]
pub max_tool_examples_per_tool: Option<u64>,
#[serde(
default,
rename = "maxToolExampleInputBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_tool_example_input_bytes: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct PluginLimits {
#[serde(
default,
rename = "maxPersistedManifestBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_persisted_manifest_bytes: Option<u64>,
#[serde(
default,
rename = "maxPersistedManifestFileBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_persisted_manifest_file_bytes: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct AcpLimits {
#[serde(
default,
rename = "maxReadLineBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_read_line_bytes: Option<u64>,
#[serde(
default,
rename = "stdoutBufferByteLimit",
skip_serializing_if = "Option::is_none"
)]
pub stdout_buffer_byte_limit: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JsRuntimeLimits {
#[serde(
default,
rename = "v8HeapLimitMb",
skip_serializing_if = "Option::is_none"
)]
pub v8_heap_limit_mb: Option<u64>,
#[serde(
default,
rename = "syncRpcWaitTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub sync_rpc_wait_timeout_ms: Option<u64>,
#[serde(
default,
rename = "cpuTimeLimitMs",
skip_serializing_if = "Option::is_none"
)]
pub cpu_time_limit_ms: Option<u64>,
#[serde(
default,
rename = "wallClockLimitMs",
skip_serializing_if = "Option::is_none"
)]
pub wall_clock_limit_ms: Option<u64>,
#[serde(
default,
rename = "importCacheMaterializeTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub import_cache_materialize_timeout_ms: Option<u64>,
#[serde(
default,
rename = "capturedOutputLimitBytes",
skip_serializing_if = "Option::is_none"
)]
pub captured_output_limit_bytes: Option<u64>,
#[serde(
default,
rename = "stdinBufferLimitBytes",
skip_serializing_if = "Option::is_none"
)]
pub stdin_buffer_limit_bytes: Option<u64>,
#[serde(
default,
rename = "eventPayloadLimitBytes",
skip_serializing_if = "Option::is_none"
)]
pub event_payload_limit_bytes: Option<u64>,
#[serde(
default,
rename = "v8IpcMaxFrameBytes",
skip_serializing_if = "Option::is_none"
)]
pub v8_ipc_max_frame_bytes: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct PythonLimits {
#[serde(
default,
rename = "outputBufferMaxBytes",
skip_serializing_if = "Option::is_none"
)]
pub output_buffer_max_bytes: Option<u64>,
#[serde(
default,
rename = "executionTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub execution_timeout_ms: Option<u64>,
#[serde(
default,
rename = "maxOldSpaceMb",
skip_serializing_if = "Option::is_none"
)]
pub max_old_space_mb: Option<u64>,
#[serde(
default,
rename = "vfsRpcTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub vfs_rpc_timeout_ms: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct WasmLimits {
#[serde(
default,
rename = "maxModuleFileBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_module_file_bytes: Option<u64>,
#[serde(
default,
rename = "capturedOutputLimitBytes",
skip_serializing_if = "Option::is_none"
)]
pub captured_output_limit_bytes: Option<u64>,
#[serde(
default,
rename = "syncReadLimitBytes",
skip_serializing_if = "Option::is_none"
)]
pub sync_read_limit_bytes: Option<u64>,
#[serde(
default,
rename = "prewarmTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub prewarm_timeout_ms: Option<u64>,
#[serde(
default,
rename = "runnerHeapLimitMb",
skip_serializing_if = "Option::is_none"
)]
pub runner_heap_limit_mb: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Permissions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub fs: Option<FsPermissions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub network: Option<PatternPermissions>,
#[serde(
default,
rename = "childProcess",
skip_serializing_if = "Option::is_none"
)]
pub child_process: Option<PatternPermissions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub process: Option<PatternPermissions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub env: Option<PatternPermissions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub binding: Option<PatternPermissions>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum PermissionMode {
Allow,
Deny,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum FsPermissions {
Mode(PermissionMode),
Rules(RulePermissions<FsPermissionRule>),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum PatternPermissions {
Mode(PermissionMode),
Rules(RulePermissions<PatternPermissionRule>),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RulePermissions<T> {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default: Option<PermissionMode>,
pub rules: Vec<T>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FsPermissionRule {
pub mode: PermissionMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operations: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub paths: Option<Vec<String>>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PatternPermissionRule {
pub mode: PermissionMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operations: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub patterns: Option<Vec<String>>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RootFilesystemConfig {
#[serde(default, rename = "type")]
pub kind: RootFilesystemKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub mode: Option<RootFilesystemMode>,
#[serde(
default,
rename = "nativePlugin",
skip_serializing_if = "Option::is_none"
)]
pub native_plugin: Option<MountPlugin>,
#[serde(default, rename = "disableDefaultBaseLayer")]
pub disable_default_base_layer: bool,
#[serde(default)]
pub lowers: Vec<RootLowerInput>,
}
impl Default for RootFilesystemConfig {
fn default() -> Self {
Self {
kind: RootFilesystemKind::Overlay,
mode: None,
native_plugin: None,
disable_default_base_layer: false,
lowers: Vec::new(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RootFilesystemKind {
#[default]
Overlay,
Native,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum RootFilesystemMode {
Ephemeral,
ReadOnly,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "kebab-case")]
pub enum RootLowerInput {
BundledBaseFilesystem,
#[serde(untagged)]
SnapshotExport(crate::fs::RootSnapshotExport),
}
pub enum MountConfig {
Plain {
path: String,
driver: Arc<dyn VirtualFileSystem>,
read_only: bool,
},
Native {
path: String,
plugin: MountPlugin,
read_only: bool,
},
Overlay {
path: String,
filesystem: OverlayMountConfig,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MountPlugin {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub config: Option<serde_json::Value>,
}
pub fn node_modules_mount(host_node_modules_dir: impl Into<String>) -> MountConfig {
MountConfig::Native {
path: "/root/node_modules".to_string(),
plugin: MountPlugin {
id: "host_dir".to_string(),
config: Some(serde_json::json!({
"hostPath": host_node_modules_dir.into(),
"readOnly": true,
})),
},
read_only: true,
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OverlayMountConfig {
#[serde(rename = "type")]
pub kind: String,
pub store: serde_json::Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub mode: Option<RootFilesystemMode>,
pub lowers: Vec<RootLowerInput>,
}
pub enum AgentOsSidecarConfig {
Shared { pool: Option<String> },
Explicit {
handle: Arc<crate::sidecar::AgentOsSidecar>,
},
}
pub type ScheduleCallback = Arc<dyn Fn() -> futures::future::BoxFuture<'static, ()> + Send + Sync>;
#[derive(Clone)]
pub struct ScheduleEntry {
pub id: String,
pub schedule: String,
pub callback: ScheduleCallback,
}
pub trait ScheduleDriver: Send + Sync {
fn schedule(&self, entry: ScheduleEntry) -> ScheduleHandle;
fn cancel(&self, handle: &ScheduleHandle);
fn dispose(&self);
}
#[derive(Clone)]
pub struct ScheduleHandle {
pub id: String,
}
#[derive(Default)]
pub struct TimerScheduleDriver {
timers: Arc<scc::HashMap<String, tokio_util::sync::CancellationToken>>,
}
impl TimerScheduleDriver {
pub fn new() -> Self {
Self {
timers: Arc::new(scc::HashMap::new()),
}
}
fn schedule_next(
timers: Arc<scc::HashMap<String, tokio_util::sync::CancellationToken>>,
entry: ScheduleEntry,
cancel: tokio_util::sync::CancellationToken,
) {
let now = chrono::Utc::now();
let parsed = match crate::cron::parse_schedule(&entry.schedule) {
Ok(parsed) => parsed,
Err(_) => {
let _ = timers.remove(&entry.id);
return;
}
};
let is_cron = parsed.is_cron();
let next = match crate::cron::resolve_next_run(&parsed, now) {
Some(next) => next,
None => {
let _ = timers.remove(&entry.id);
return;
}
};
let delay = (next - now).to_std().unwrap_or(std::time::Duration::ZERO);
tokio::spawn(async move {
tokio::select! {
_ = cancel.cancelled() => {
return;
}
_ = tokio::time::sleep(delay) => {}
}
if cancel.is_cancelled() {
return;
}
(entry.callback)().await;
if is_cron && timers.contains(&entry.id) {
Self::schedule_next(Arc::clone(&timers), entry, cancel);
} else {
let _ = timers.remove(&entry.id);
}
});
}
}
impl ScheduleDriver for TimerScheduleDriver {
fn schedule(&self, entry: ScheduleEntry) -> ScheduleHandle {
let id = entry.id.clone();
let cancel = tokio_util::sync::CancellationToken::new();
if let Some((_, old)) = self.timers.remove(&id) {
old.cancel();
}
let _ = self.timers.insert(id.clone(), cancel.clone());
Self::schedule_next(Arc::clone(&self.timers), entry, cancel);
ScheduleHandle { id }
}
fn cancel(&self, handle: &ScheduleHandle) {
if let Some((_, cancel)) = self.timers.remove(&handle.id) {
cancel.cancel();
}
}
fn dispose(&self) {
self.timers.scan(|_, cancel| cancel.cancel());
self.timers.clear();
}
}