use std::sync::Arc;
use serde::{Deserialize, Serialize};
use crate::fs::VirtualFileSystem;
pub use agentos_vm_config::{VmGroupConfig, VmUserAccountConfig, VmUserConfig};
#[derive(Default)]
pub struct AgentOsConfig {
pub database: Option<agentos_vm_config::VmSqliteDescriptor>,
pub user: Option<VmUserConfig>,
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 bindings: Vec<Bindings>,
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 database(mut self, database: agentos_vm_config::VmSqliteDescriptor) -> Self {
self.config.database = Some(database);
self
}
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 user(mut self, user: VmUserConfig) -> Self {
self.config.user = Some(user);
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 bindings(mut self, bindings: Vec<Bindings>) -> Self {
self.config.bindings = bindings;
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,
Binding,
}
#[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 BindingCallback = 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 Binding {
pub name: String,
pub description: String,
pub input_schema: serde_json::Value,
pub timeout_ms: Option<u64>,
pub execute: BindingCallback,
}
#[derive(Clone)]
pub struct Bindings {
pub name: String,
pub description: String,
pub bindings: Vec<Binding>,
}
#[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 bindings: Option<BindingLimits>,
#[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, skip_serializing_if = "Option::is_none")]
pub sqlite: Option<SqliteLimits>,
#[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>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub process: Option<ProcessLimits>,
}
#[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 BindingLimits {
#[serde(
default,
rename = "defaultBindingTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub default_binding_timeout_ms: Option<u64>,
#[serde(
default,
rename = "maxBindingTimeoutMs",
skip_serializing_if = "Option::is_none"
)]
pub max_binding_timeout_ms: Option<u64>,
#[serde(
default,
rename = "maxRegisteredCollections",
skip_serializing_if = "Option::is_none"
)]
pub max_registered_collections: Option<u64>,
#[serde(
default,
rename = "maxRegisteredBindingsPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_registered_bindings_per_vm: Option<u64>,
#[serde(
default,
rename = "maxBindingsPerCollection",
skip_serializing_if = "Option::is_none"
)]
pub max_bindings_per_collection: Option<u64>,
#[serde(
default,
rename = "maxBindingSchemaBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_binding_schema_bytes: Option<u64>,
#[serde(
default,
rename = "maxExamplesPerBinding",
skip_serializing_if = "Option::is_none"
)]
pub max_examples_per_binding: Option<u64>,
#[serde(
default,
rename = "maxBindingExampleInputBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_binding_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>,
#[serde(
default,
rename = "maxCompletedMessageBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_completed_message_bytes: Option<u64>,
#[serde(
default,
rename = "maxTurnOutputBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_turn_output_bytes: Option<u64>,
#[serde(
default,
rename = "maxPromptBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_prompt_bytes: Option<u64>,
#[serde(
default,
rename = "maxPromptBlocks",
skip_serializing_if = "Option::is_none"
)]
pub max_prompt_blocks: Option<u64>,
#[serde(
default,
rename = "maxFallbackContinuationBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_fallback_continuation_bytes: Option<u64>,
#[serde(
default,
rename = "maxSessionHistoryBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_session_history_bytes: Option<u64>,
#[serde(
default,
rename = "maxSessionHistoryEvents",
skip_serializing_if = "Option::is_none"
)]
pub max_session_history_events: Option<u64>,
#[serde(
default,
rename = "maxHistoryPageEntries",
skip_serializing_if = "Option::is_none"
)]
pub max_history_page_entries: Option<u64>,
#[serde(
default,
rename = "maxSessionListEntries",
skip_serializing_if = "Option::is_none"
)]
pub max_session_list_entries: Option<u64>,
#[serde(
default,
rename = "maxSessionsPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_sessions_per_vm: Option<u64>,
#[serde(
default,
rename = "maxPromptsPerSession",
skip_serializing_if = "Option::is_none"
)]
pub max_prompts_per_session: Option<u64>,
#[serde(
default,
rename = "maxPromptsPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_prompts_per_vm: Option<u64>,
#[serde(
default,
rename = "maxPendingPermissionsPerSession",
skip_serializing_if = "Option::is_none"
)]
pub max_pending_permissions_per_session: Option<u64>,
#[serde(
default,
rename = "maxPendingPermissionsPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_pending_permissions_per_vm: Option<u64>,
#[serde(
default,
rename = "maxPermissionOutcomesPerSession",
skip_serializing_if = "Option::is_none"
)]
pub max_permission_outcomes_per_session: Option<u64>,
#[serde(
default,
rename = "maxPermissionOutcomesPerVm",
skip_serializing_if = "Option::is_none"
)]
pub max_permission_outcomes_per_vm: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SqliteLimits {
#[serde(
default,
rename = "maxResultBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_result_bytes: 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>,
#[serde(
default,
rename = "runnerCpuTimeLimitMs",
skip_serializing_if = "Option::is_none"
)]
pub runner_cpu_time_limit_ms: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ProcessLimits {
#[serde(
default,
rename = "maxSpawnFileActions",
skip_serializing_if = "Option::is_none"
)]
pub max_spawn_file_actions: Option<u64>,
#[serde(
default,
rename = "maxSpawnFileActionBytes",
skip_serializing_if = "Option::is_none"
)]
pub max_spawn_file_action_bytes: Option<u64>,
#[serde(
default,
rename = "pendingStdinBytes",
skip_serializing_if = "Option::is_none"
)]
pub pending_stdin_bytes: Option<u64>,
#[serde(
default,
rename = "pendingEventCount",
skip_serializing_if = "Option::is_none"
)]
pub pending_event_count: Option<u64>,
#[serde(
default,
rename = "pendingEventBytes",
skip_serializing_if = "Option::is_none"
)]
pub pending_event_bytes: 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>,
guest_source: Option<String>,
guest_fstype: Option<String>,
read_only: bool,
},
Native {
path: String,
plugin: MountPlugin,
guest_source: Option<String>,
guest_fstype: Option<String>,
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,
})),
},
guest_source: Some(String::from("host_dir")),
guest_fstype: Some(String::from("host_dir")),
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();
}
}