diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml new file mode 100644 index 0000000..864e5d8 --- /dev/null +++ b/.github/workflows/lint.yml @@ -0,0 +1,52 @@ +name: Lint & Format + +on: + pull_request: + branches: + - main + push: + branches: + - main + +jobs: + rust-fmt: + name: Rust Format Check + runs-on: ubuntu-latest + + steps: + - name: Checkout repo + uses: actions/checkout@v4 + + - name: Install Rust + uses: dtolnay/rust-toolchain@stable + with: + components: rustfmt + + - name: Check formatting (capsule-cli) + working-directory: crates/capsule-cli + run: cargo fmt --all -- --check + + - name: Check formatting (capsule-core) + working-directory: crates/capsule-core + run: cargo fmt --all -- --check + + rust-clippy: + name: Rust Clippy Check + runs-on: ubuntu-latest + + steps: + - name: Checkout repo + uses: actions/checkout@v4 + + - name: Install Rust + uses: dtolnay/rust-toolchain@stable + with: + components: clippy + + - name: Run clippy (capsule-cli) + working-directory: crates/capsule-cli + run: cargo clippy --all-targets --all-features -- -D warnings + + - name: Run clippy (capsule-core) + working-directory: crates/capsule-core + run: cargo clippy --all-targets --all-features -- -D warnings diff --git a/.gitignore b/.gitignore index 7acd56b..e048071 100644 --- a/.gitignore +++ b/.gitignore @@ -11,4 +11,3 @@ capsule.egg-info **/*.rs.bk *.pdb **/mutants.out*/ - diff --git a/ROADMAP.md b/ROADMAP.md index af5b839..3d49f41 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -13,7 +13,7 @@ This document tracks the development status of Capsule. We follow a "Release Ear - [x] **Core:** Rust Host capable of loading Wasm Components. - [x] **SDK (Python):** Basic `@task` decorator with JSON serialization. - [x] **CLI:** `capsule run main.py` with JIT compilation. -- [ ] **Limits:** Basic Fuel metering for CPU protection. +- [x] **Limits:** Basic Fuel metering for CPU protection. --- diff --git a/crates/capsule-cli/Cargo.lock b/crates/capsule-cli/Cargo.lock index 725ec89..869af59 100644 --- a/crates/capsule-cli/Cargo.lock +++ b/crates/capsule-cli/Cargo.lock @@ -252,6 +252,7 @@ version = "0.1.0" dependencies = [ "capsule-core", "clap", + "indicatif", "serde_json", "tokio", ] @@ -261,6 +262,7 @@ name = "capsule-core" version = "0.1.0" dependencies = [ "anyhow", + "humantime", "nanoid", "rusqlite", "serde", @@ -343,6 +345,19 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75" +[[package]] +name = "console" +version = "0.15.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "054ccb5b10f9f2cbf51eb355ca1d05c2d279ce1804688d0db74b4733a5aeafd8" +dependencies = [ + "encode_unicode", + "libc", + "once_cell", + "unicode-width", + "windows-sys 0.59.0", +] + [[package]] name = "core-foundation-sys" version = "0.8.7" @@ -607,6 +622,12 @@ version = "0.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d" +[[package]] +name = "encode_unicode" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34aa73646ffb006b8f5147f3dc182bd4bcb190227ce861fc4a4844bf8e3cb2c0" + [[package]] name = "encoding_rs" version = "0.8.35" @@ -854,6 +875,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "humantime" +version = "2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "135b12329e5e3ce057a9f972339ea52bc954fe1e9358ef27f95e89716fbc5424" + [[package]] name = "iana-time-zone" version = "0.1.64" @@ -998,6 +1025,19 @@ dependencies = [ "serde_core", ] +[[package]] +name = "indicatif" +version = "0.17.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "183b3088984b400f4cfac3620d5e076c84da5364016b4f49473de574b2586235" +dependencies = [ + "console", + "number_prefix", + "portable-atomic", + "unicode-width", + "web-time", +] + [[package]] name = "io-extras" version = "0.18.4" @@ -1199,6 +1239,12 @@ dependencies = [ "rand", ] +[[package]] +name = "number_prefix" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "830b246a0e5f20af87141b25c173cd1b609bd7779a4617d6ec582abaf90870f3" + [[package]] name = "object" version = "0.32.2" @@ -1262,6 +1308,12 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" +[[package]] +name = "portable-atomic" +version = "1.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f84267b20a16ea918e43c6a88433c2d54fa145c92a811b5b047ccbe153674483" + [[package]] name = "postcard" version = "1.1.3" @@ -2324,6 +2376,16 @@ dependencies = [ "wast 243.0.0", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "wiggle" version = "29.0.1" diff --git a/crates/capsule-cli/Cargo.toml b/crates/capsule-cli/Cargo.toml index 0d8b88a..ad27063 100644 --- a/crates/capsule-cli/Cargo.toml +++ b/crates/capsule-cli/Cargo.toml @@ -3,11 +3,13 @@ name = "capsule-cli" version = "0.1.0" edition = "2024" +[[bin]] +name = "capsule" +path = "src/main.rs" [dependencies] clap = { version = "4.5.53", features = ["derive"] } capsule-core = { path = "../capsule-core" } tokio = { version = "1.48.0", features = ["rt", "rt-multi-thread", "macros"] } serde_json = "1" - - +indicatif = "0.17" diff --git a/crates/capsule-cli/src/commands/run.rs b/crates/capsule-cli/src/commands/run.rs index f6b9bba..423ecc3 100644 --- a/crates/capsule-cli/src/commands/run.rs +++ b/crates/capsule-cli/src/commands/run.rs @@ -1,23 +1,26 @@ use std::fmt; use std::path::Path; +use std::time::Instant; + +use indicatif::{ProgressBar, ProgressStyle}; use capsule_core::wasm::commands::create::CreateInstance; use capsule_core::wasm::commands::run::RunInstance; -use capsule_core::wasm::execution_policy::ExecutionPolicy; use capsule_core::wasm::compiler::python::{PythonWasmCompiler, PythonWasmCompilerError}; -use capsule_core::wasm::runtime::WasmRuntimeError; +use capsule_core::wasm::execution_policy::{Compute, ExecutionPolicy}; use capsule_core::wasm::runtime::Runtime; +use capsule_core::wasm::runtime::WasmRuntimeError; pub enum RunError { - IoError(String), - CompileFailed(String), + IoError(String), + CompileFailed(String), } impl fmt::Display for RunError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { RunError::IoError(msg) => write!(f, "{}", msg), - RunError::CompileFailed(msg) => write!(f, "{}", msg), + RunError::CompileFailed(msg) => write!(f, "🛠️ Mission aborted: {}", msg), } } } @@ -42,28 +45,81 @@ impl From for RunError { pub async fn execute(file_path: &Path, args: Vec) -> Result { let compiler = PythonWasmCompiler::new(file_path)?; + + let spinner = ProgressBar::new_spinner(); + spinner.set_style( + ProgressStyle::default_spinner() + .tick_strings(&["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]) + .template("{spinner:.cyan} {msg}") + .unwrap(), + ); + spinner.enable_steady_tick(std::time::Duration::from_millis(80)); + + spinner.set_message("Preparing capsule environment"); + let wasm_path = compiler.compile_wasm()?; + spinner.finish_and_clear(); + + let spinner = ProgressBar::new_spinner(); + spinner.set_style( + ProgressStyle::default_spinner() + .tick_strings(&["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]) + .template("{spinner:.cyan} {msg}") + .unwrap(), + ); + spinner.enable_steady_tick(std::time::Duration::from_millis(80)); + + spinner.set_message("Launching capsule runtime"); let runtime_config = capsule_core::wasm::runtime::RuntimeConfig { cache_dir: compiler.cache_dir.clone(), }; let runtime = Runtime::with_config(runtime_config)?; - let execution_policy = ExecutionPolicy::default(); - let create_instance_command = CreateInstance::new(execution_policy.clone(), args.clone()) - .wasm_path(wasm_path); + let execution_policy = ExecutionPolicy::default().compute(Some(Compute::High)); + let create_instance_command = + CreateInstance::new(execution_policy.clone(), args.clone()).wasm_path(wasm_path); let (store, instance, task_id) = runtime.execute(create_instance_command).await?; + spinner.finish_and_clear(); + + println!("📡 Capsule in orbit. Systems nominal."); + + let start_time = Instant::now(); let args_json = serde_json::json!({ "task_name": "main", "args": args, "kwargs": {} - }).to_string(); + }) + .to_string(); - let run_instance_command = RunInstance::new(task_id, execution_policy, store, instance, args_json); + let run_instance_command = + RunInstance::new(task_id, execution_policy, store, instance, args_json); let result = runtime.execute(run_instance_command).await?; + let elapsed = start_time.elapsed(); + let time_str = format_duration(elapsed); + println!("✓ Complete ({})", time_str); + Ok(result) } + +fn format_duration(duration: std::time::Duration) -> String { + let total_secs = duration.as_secs_f64(); + + if total_secs < 60.0 { + format!("{:.2}s", total_secs) + } else if total_secs < 3600.0 { + let minutes = (total_secs / 60.0).floor() as u64; + let seconds = total_secs % 60.0; + format!("{}m {:.0}s", minutes, seconds) + } else { + let hours = (total_secs / 3600.0).floor() as u64; + let remaining_secs = total_secs % 3600.0; + let minutes = (remaining_secs / 60.0).floor() as u64; + let seconds = remaining_secs % 60.0; + format!("{}h {}m {:.0}s", hours, minutes, seconds) + } +} diff --git a/crates/capsule-cli/src/main.rs b/crates/capsule-cli/src/main.rs index d14b723..15ffb0a 100644 --- a/crates/capsule-cli/src/main.rs +++ b/crates/capsule-cli/src/main.rs @@ -1,16 +1,16 @@ pub mod cli; pub mod commands; +use clap::Parser; use std::fmt; use std::path::Path; -use clap::Parser; use cli::{Cli, Commands}; -use commands::{run, RunError}; +use commands::{RunError, run}; #[derive(Debug)] pub enum CliError { - RunError(String), + RunError(String), } impl fmt::Display for CliError { @@ -27,8 +27,6 @@ impl From for CliError { } } - - #[tokio::main] async fn main() -> Result<(), CliError> { let cli = Cli::parse(); diff --git a/crates/capsule-core/Cargo.lock b/crates/capsule-core/Cargo.lock index d160849..acb92fe 100644 --- a/crates/capsule-core/Cargo.lock +++ b/crates/capsule-core/Cargo.lock @@ -237,6 +237,7 @@ name = "capsule-core" version = "0.1.0" dependencies = [ "anyhow", + "humantime", "nanoid", "rusqlite", "rustfmt", @@ -838,6 +839,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "humantime" +version = "2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "135b12329e5e3ce057a9f972339ea52bc954fe1e9358ef27f95e89716fbc5424" + [[package]] name = "iana-time-zone" version = "0.1.64" diff --git a/crates/capsule-core/Cargo.toml b/crates/capsule-core/Cargo.toml index f87b9a5..c93d23a 100644 --- a/crates/capsule-core/Cargo.toml +++ b/crates/capsule-core/Cargo.toml @@ -8,6 +8,7 @@ crate-type = ["rlib"] [dependencies] anyhow = "1" +humantime = "2" nanoid = "0.4.0" rusqlite = "0.37.0" serde = { version = "1.0.228", features = ["derive"] } @@ -18,6 +19,3 @@ wasmtime-wasi = "29.0.0" [dev-dependencies] rustfmt = "0.10.0" - - - diff --git a/crates/capsule-core/src/config/database.rs b/crates/capsule-core/src/config/database.rs index fc99f6b..6701230 100644 --- a/crates/capsule-core/src/config/database.rs +++ b/crates/capsule-core/src/config/database.rs @@ -58,9 +58,9 @@ impl Database { Some(path) => { let database_path = &format!("{}/{}", path, database_name); - std::fs::create_dir_all(&path)?; + std::fs::create_dir_all(path)?; - Connection::open(&database_path)? + Connection::open(database_path)? } None => Connection::open(":memory:")?, }; @@ -309,7 +309,7 @@ mod tests { }) .expect("Failed to query test"); - assert!(result.len() > 0, "Failed to execute test"); + assert!(!result.is_empty(), "Failed to execute test"); } #[test] diff --git a/crates/capsule-core/src/config/log.rs b/crates/capsule-core/src/config/log.rs index 288f824..b91c835 100644 --- a/crates/capsule-core/src/config/log.rs +++ b/crates/capsule-core/src/config/log.rs @@ -50,6 +50,7 @@ pub enum InstanceState { Completed, Failed, Interrupted, + TimedOut, } impl fmt::Display for InstanceState { @@ -60,6 +61,7 @@ impl fmt::Display for InstanceState { InstanceState::Completed => "completed", InstanceState::Failed => "failed", InstanceState::Interrupted => "interrupted", + InstanceState::TimedOut => "timed_out", }; write!(f, "{}", state_str) } @@ -378,11 +380,11 @@ mod tests { conn.execute("INSERT INTO instance_log (id, agent_name, agent_version, task_id, task_name, state, fuel_limit, fuel_consumed) VALUES (?, ?, ?, ?, ?, ?, ?, ?)", [ &nanoid!(10), - &"test_agent".to_string(), - &"1.0.0".to_string(), - &"test_task_123".to_string(), - &"Test Task".to_string(), - &"created".to_string(), + "test_agent", + "1.0.0", + "test_task_123", + "Test Task", + "created", "15000000", "0", ]).expect("Failed to insert test data"); diff --git a/crates/capsule-core/src/wasm/commands/run.rs b/crates/capsule-core/src/wasm/commands/run.rs index 973ee3b..6abb55c 100644 --- a/crates/capsule-core/src/wasm/commands/run.rs +++ b/crates/capsule-core/src/wasm/commands/run.rs @@ -46,11 +46,29 @@ impl RuntimeCommand for RunInstance { }) .await?; - let result = self + let wasm_future = self .instance .capsule_host_task_runner() - .call_run(&mut self.store, &self.args_json) - .await; + .call_run(&mut self.store, &self.args_json); + + let result = match self.policy.timeout_duration() { + Some(duration) => match tokio::time::timeout(duration, wasm_future).await { + Ok(inner_result) => inner_result, + Err(_elapsed) => { + runtime + .log + .update_log(UpdateInstanceLog { + task_id: self.task_id.clone(), + state: InstanceState::TimedOut, + fuel_consumed: self.policy.compute.as_fuel() + - self.store.get_fuel().unwrap_or(0), + }) + .await?; + return Err(WasmRuntimeError::Timeout(self.policy.name)); + } + }, + None => wasm_future.await, + }; match result { Ok(Ok(value)) => { diff --git a/crates/capsule-core/src/wasm/compiler/python.rs b/crates/capsule-core/src/wasm/compiler/python.rs index 82dc0e2..1dfa5ec 100644 --- a/crates/capsule-core/src/wasm/compiler/python.rs +++ b/crates/capsule-core/src/wasm/compiler/python.rs @@ -2,6 +2,7 @@ use std::fmt; use std::fs; use std::path::{Path, PathBuf}; use std::process::Command; +use std::process::Stdio; use super::CAPSULE_WIT; @@ -13,7 +14,9 @@ pub enum PythonWasmCompilerError { impl fmt::Display for PythonWasmCompilerError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { - PythonWasmCompilerError::CompileFailed(msg) => write!(f, "Compilation failed > {}", msg), + PythonWasmCompilerError::CompileFailed(msg) => { + write!(f, "Compilation failed > {}", msg) + } PythonWasmCompilerError::FsError(msg) => write!(f, "File system error > {}", msg), } } @@ -39,11 +42,15 @@ pub struct PythonWasmCompiler { impl PythonWasmCompiler { pub fn new(source_path: &Path) -> Result { - let source_path = source_path.canonicalize() - .map_err(|e| PythonWasmCompilerError::FsError(format!("Cannot resolve source path: {}", e)))?; + let source_path = source_path.canonicalize().map_err(|e| { + PythonWasmCompilerError::FsError(format!("Cannot resolve source path: {}", e)) + })?; - let source_dir = source_path.parent() - .ok_or(PythonWasmCompilerError::FsError("Cannot determine source directory".to_string()))?; + let source_dir = source_path + .parent() + .ok_or(PythonWasmCompilerError::FsError( + "Cannot determine source directory".to_string(), + ))?; let cache_dir = source_dir.join(".capsule"); let output_wasm = cache_dir.join("capsule.wasm"); @@ -61,30 +68,40 @@ impl PythonWasmCompiler { pub fn compile_wasm(&self) -> Result { if self.needs_rebuild(&self.source_path, &self.output_wasm)? { - let module_name = self.source_path + let module_name = self + .source_path .file_stem() - .ok_or(PythonWasmCompilerError::FsError("Invalid source file name".to_string()))? + .ok_or(PythonWasmCompilerError::FsError( + "Invalid source file name".to_string(), + ))? .to_str() - .ok_or(PythonWasmCompilerError::FsError("Invalid UTF-8 in file name".to_string()))?; + .ok_or(PythonWasmCompilerError::FsError( + "Invalid UTF-8 in file name".to_string(), + ))?; - let python_path = self.source_path + let python_path = self + .source_path .parent() - .ok_or(PythonWasmCompilerError::FsError("Cannot determine parent directory".to_string()))?; + .ok_or(PythonWasmCompilerError::FsError( + "Cannot determine parent directory".to_string(), + ))?; let wit_path = self.get_wit_path()?; let sdk_path = self.get_sdk_path()?; if !sdk_path.exists() { - return Err(PythonWasmCompilerError::FsError( - format!("SDK directory not found: {}", sdk_path.display()) - )); + return Err(PythonWasmCompilerError::FsError(format!( + "SDK directory not found: {}", + sdk_path.display() + ))); } if !sdk_path.exists() { - return Err(PythonWasmCompilerError::FsError( - format!("SDK directory not found: {}", sdk_path.display()) - )); + return Err(PythonWasmCompilerError::FsError(format!( + "SDK directory not found: {}", + sdk_path.display() + ))); } let bootloader_path = self.cache_dir.join("_capsule_boot.py"); @@ -121,6 +138,8 @@ from capsule.app import TaskRunner, exports .arg(&sdk_path) .arg("-o") .arg(&self.output_wasm) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) .status()?; if !status.success() { @@ -133,7 +152,11 @@ from capsule.app import TaskRunner, exports Ok(self.output_wasm.clone()) } - fn needs_rebuild(&self, source: &Path, wasm_path: &Path) -> Result { + fn needs_rebuild( + &self, + source: &Path, + wasm_path: &Path, + ) -> Result { if !wasm_path.exists() { return Ok(true); } @@ -145,19 +168,18 @@ from capsule.app import TaskRunner, exports return Ok(true); } - if let Some(source_dir) = source.parent() { - if let Ok(entries) = fs::read_dir(source_dir) { - for entry in entries.flatten() { - let path = entry.path(); - if path.extension().map_or(false, |ext| ext == "py") && path != source { - if let Ok(metadata) = fs::metadata(&path) { - if let Ok(modified) = metadata.modified() { - if modified > wasm_time { - return Ok(true); - } - } - } - } + if let Some(source_dir) = source.parent() + && let Ok(entries) = fs::read_dir(source_dir) + { + for entry in entries.flatten() { + let path = entry.path(); + if path.extension().is_some_and(|ext| ext == "py") + && path != source + && let Ok(metadata) = fs::metadata(&path) + && let Ok(modified) = metadata.modified() + && modified > wasm_time + { + return Ok(true); } } } @@ -189,23 +211,22 @@ from capsule.app import TaskRunner, exports return Ok(PathBuf::from(path)); } - if let Ok(exe_path) = std::env::current_exe() { - if let Some(project_root) = exe_path + if let Ok(exe_path) = std::env::current_exe() + && let Some(project_root) = exe_path .parent() .and_then(|p| p.parent()) .and_then(|p| p.parent()) .and_then(|p| p.parent()) .and_then(|p| p.parent()) - { - let sdk_path = project_root.join("crates/capsule-sdk/python/src"); - if sdk_path.exists() { - return Ok(sdk_path); - } + { + let sdk_path = project_root.join("crates/capsule-sdk/python/src"); + if sdk_path.exists() { + return Ok(sdk_path); } } Err(PythonWasmCompilerError::FsError( - "Cannot find SDK. Set CAPSULE_SDK_PATH environment variable.".to_string() + "Cannot find SDK. Set CAPSULE_SDK_PATH environment variable.".to_string(), )) } } diff --git a/crates/capsule-core/src/wasm/execution_policy.rs b/crates/capsule-core/src/wasm/execution_policy.rs index b1687b7..a2bc5b0 100644 --- a/crates/capsule-core/src/wasm/execution_policy.rs +++ b/crates/capsule-core/src/wasm/execution_policy.rs @@ -1,4 +1,5 @@ use serde::{Deserialize, Serialize}; +use std::time::Duration; #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] @@ -41,7 +42,7 @@ impl Default for ExecutionPolicy { compute: Compute::Medium, ram: None, timeout: None, - max_retries: 1, + max_retries: 0, env_vars: None, } } @@ -87,6 +88,12 @@ impl ExecutionPolicy { self.env_vars = env_vars; self } + + pub fn timeout_duration(&self) -> Option { + self.timeout + .as_ref() + .and_then(|s| humantime::parse_duration(s).ok()) + } } #[cfg(test)] diff --git a/crates/capsule-core/src/wasm/runtime.rs b/crates/capsule-core/src/wasm/runtime.rs index 5a5d7b4..8cccf8c 100644 --- a/crates/capsule-core/src/wasm/runtime.rs +++ b/crates/capsule-core/src/wasm/runtime.rs @@ -11,6 +11,7 @@ pub enum WasmRuntimeError { WasmtimeError(wasmtime::Error), LogError(LogError), ConfigError(String), + Timeout(String), } impl fmt::Display for WasmRuntimeError { @@ -21,6 +22,7 @@ impl fmt::Display for WasmRuntimeError { } WasmRuntimeError::LogError(msg) => write!(f, "Runtime error > {}", msg), WasmRuntimeError::ConfigError(msg) => write!(f, "Runtime error > Config > {}", msg), + WasmRuntimeError::Timeout(task_id) => write!(f, "Task '{}' timed out", task_id), } } } @@ -75,7 +77,10 @@ impl Runtime { pub fn with_config(config: RuntimeConfig) -> Result, WasmRuntimeError> { let mut engine_config = Config::new(); let db_path = config.cache_dir.join("state.db"); - let log = Log::new(Some(db_path.parent().unwrap().to_str().unwrap()), db_path.file_name().unwrap().to_str().unwrap())?; + let log = Log::new( + Some(db_path.parent().unwrap().to_str().unwrap()), + db_path.file_name().unwrap().to_str().unwrap(), + )?; engine_config.wasm_component_model(true); engine_config.async_support(true); diff --git a/crates/capsule-core/src/wasm/state.rs b/crates/capsule-core/src/wasm/state.rs index cacf659..4848e67 100644 --- a/crates/capsule-core/src/wasm/state.rs +++ b/crates/capsule-core/src/wasm/state.rs @@ -52,27 +52,42 @@ impl Host for State { let task_config: TaskConfig = serde_json::from_str(&config).unwrap_or_default(); let policy = task_config.to_execution_policy(); + let max_retries = policy.max_retries; + + let mut last_error: Option = None; + + for attempt in 0..=max_retries { + let create_cmd = CreateInstance::new(policy.clone(), vec![]).task_name(&name); + + let (store, instance, task_id) = match runtime.execute(create_cmd).await { + Ok(result) => result, + Err(e) => { + last_error = Some(format!("Failed to create instance: {}", e)); + continue; + } + }; + + let args_json = format!( + r#"{{"task_name": "{}", "args": {}, "kwargs": {{}}}}"#, + name, args + ); + + let run_cmd = RunInstance::new(task_id, policy.clone(), store, instance, args_json); + + match runtime.execute(run_cmd).await { + Ok(result) => return Ok(result), + Err(e) => { + last_error = Some(format!("Failed to run instance: {}", e)); + if attempt < max_retries { + continue; + } + } + } + } - let create_cmd = CreateInstance::new(policy.clone(), vec![]).task_name(&name); - - let (store, instance, task_id) = runtime - .execute(create_cmd) - .await - .map_err(|e| TaskError::InternalError(format!("Failed to create instance: {}", e)))?; - - let args_json = format!( - r#"{{"task_name": "{}", "args": {}, "kwargs": {{}}}}"#, - name, args - ); - - let run_cmd = RunInstance::new(task_id, policy, store, instance, args_json); - - let result = runtime - .execute(run_cmd) - .await - .map_err(|e| TaskError::InternalError(format!("Failed to run instance: {}", e)))?; - - Ok(result) + Err(TaskError::InternalError(last_error.unwrap_or_else(|| { + "Unknown error after retries".to_string() + }))) } } diff --git a/crates/capsule-core/src/wasm/utilities/task_config.rs b/crates/capsule-core/src/wasm/utilities/task_config.rs index 212e7ea..a96cebd 100644 --- a/crates/capsule-core/src/wasm/utilities/task_config.rs +++ b/crates/capsule-core/src/wasm/utilities/task_config.rs @@ -21,7 +21,10 @@ impl TaskConfig { "LOW" => Compute::Low, "MEDIUM" => Compute::Medium, "HIGH" => Compute::High, - _ => Compute::Medium, + _ => c + .parse::() + .map(Compute::Custom) + .unwrap_or(Compute::Medium), }); let ram = self.ram.as_ref().and_then(|r| Self::parse_ram_string(r)); @@ -67,13 +70,31 @@ mod tests { #[test] fn test_parse_ram_string() { - assert_eq!(TaskConfig::parse_ram_string("2GB"), Some(2 * 1024 * 1024 * 1024)); - assert_eq!(TaskConfig::parse_ram_string("1 GB"), Some(1024 * 1024 * 1024)); - assert_eq!(TaskConfig::parse_ram_string("4gb"), Some(4 * 1024 * 1024 * 1024)); - - assert_eq!(TaskConfig::parse_ram_string("512MB"), Some(512 * 1024 * 1024)); - assert_eq!(TaskConfig::parse_ram_string("256 MB"), Some(256 * 1024 * 1024)); - assert_eq!(TaskConfig::parse_ram_string("128mb"), Some(128 * 1024 * 1024)); + assert_eq!( + TaskConfig::parse_ram_string("2GB"), + Some(2 * 1024 * 1024 * 1024) + ); + assert_eq!( + TaskConfig::parse_ram_string("1 GB"), + Some(1024 * 1024 * 1024) + ); + assert_eq!( + TaskConfig::parse_ram_string("4gb"), + Some(4 * 1024 * 1024 * 1024) + ); + + assert_eq!( + TaskConfig::parse_ram_string("512MB"), + Some(512 * 1024 * 1024) + ); + assert_eq!( + TaskConfig::parse_ram_string("256 MB"), + Some(256 * 1024 * 1024) + ); + assert_eq!( + TaskConfig::parse_ram_string("128mb"), + Some(128 * 1024 * 1024) + ); assert_eq!(TaskConfig::parse_ram_string("1024KB"), Some(1024 * 1024)); assert_eq!(TaskConfig::parse_ram_string("512 KB"), Some(512 * 1024)); @@ -99,7 +120,7 @@ mod tests { assert_eq!(policy.compute, Compute::Medium); assert_eq!(policy.ram, None); assert_eq!(policy.timeout, None); - assert_eq!(policy.max_retries, 1); + assert_eq!(policy.max_retries, 0); assert_eq!(policy.env_vars, None); } @@ -121,7 +142,10 @@ mod tests { assert_eq!(policy.ram, Some(2 * 1024 * 1024 * 1024)); assert_eq!(policy.timeout, Some("30s".to_string())); assert_eq!(policy.max_retries, 3); - assert_eq!(policy.env_vars, Some(vec![("KEY".to_string(), "VALUE".to_string())])); + assert_eq!( + policy.env_vars, + Some(vec![("KEY".to_string(), "VALUE".to_string())]) + ); } #[test] diff --git a/crates/capsule-sdk/python/src/capsule/__init__.py b/crates/capsule-sdk/python/src/capsule/__init__.py index d28fad1..66cbb20 100644 --- a/crates/capsule-sdk/python/src/capsule/__init__.py +++ b/crates/capsule-sdk/python/src/capsule/__init__.py @@ -3,4 +3,3 @@ from .app import TaskRunner exports = app.exports - diff --git a/crates/capsule-wit/deps/sockets/tcp.wit b/crates/capsule-wit/deps/sockets/tcp.wit index 5902b9e..63627e4 100644 --- a/crates/capsule-wit/deps/sockets/tcp.wit +++ b/crates/capsule-wit/deps/sockets/tcp.wit @@ -15,7 +15,7 @@ interface tcp { /// Similar to `SHUT_RDWR` in POSIX. both, } - + /// A TCP socket resource. /// /// The socket can be in one of the following states: @@ -59,10 +59,10 @@ interface tcp { /// - `address-not-bindable`: `local-address` is not an address that the `network` can bind to. (EADDRNOTAVAIL) /// - `not-in-progress`: A `bind` operation is not in progress. /// - `would-block`: Can't finish the operation, it is still in progress. (EWOULDBLOCK, EAGAIN) - /// + /// /// # Implementors note /// When binding to a non-zero port, this bind operation shouldn't be affected by the TIME_WAIT - /// state of a recently closed socket on the same local address. In practice this means that the SO_REUSEADDR + /// state of a recently closed socket on the same local address. In practice this means that the SO_REUSEADDR /// socket option should be set implicitly on all platforms, except on Windows where this is the default behavior /// and SO_REUSEADDR performs something different entirely. /// diff --git a/crates/capsule-wit/deps/sockets/udp.wit b/crates/capsule-wit/deps/sockets/udp.wit index d987a0a..48722fa 100644 --- a/crates/capsule-wit/deps/sockets/udp.wit +++ b/crates/capsule-wit/deps/sockets/udp.wit @@ -6,7 +6,7 @@ interface udp { /// A received datagram. record incoming-datagram { /// The payload. - /// + /// /// Theoretical max size: ~64 KiB. In practice, typically less than 1500 bytes. data: list, @@ -79,7 +79,7 @@ interface udp { /// This method may be called multiple times on the same socket to change its association, but /// only the most recently returned pair of streams will be operational. Implementations may trap if /// the streams returned by a previous invocation haven't been dropped yet before calling `stream` again. - /// + /// /// The POSIX equivalent in pseudo-code is: /// ```text /// if (was previously connected) { @@ -91,7 +91,7 @@ interface udp { /// ``` /// /// Unlike in POSIX, the socket must already be explicitly bound. - /// + /// /// # Typical errors /// - `invalid-argument`: The `remote-address` has the wrong address family. (EAFNOSUPPORT) /// - `invalid-argument`: The IP address in `remote-address` is set to INADDR_ANY (`0.0.0.0` / `::`). (EDESTADDRREQ, EADDRNOTAVAIL) @@ -115,7 +115,7 @@ interface udp { /// > stored in the object pointed to by `address` is unspecified. /// /// WASI is stricter and requires `local-address` to return `invalid-state` when the socket hasn't been bound yet. - /// + /// /// # Typical errors /// - `invalid-state`: The socket is not bound to any local address. /// @@ -217,7 +217,7 @@ interface udp { /// When this function returns ok(0), the `subscribe` pollable will /// become ready when this function will report at least ok(1), or an /// error. - /// + /// /// Never returns `would-block`. check-send: func() -> result; @@ -256,7 +256,7 @@ interface udp { /// - /// - send: func(datagrams: list) -> result; - + /// Create a `pollable` which will resolve once the stream is ready to send again. /// /// Note: this function is here for WASI Preview2 only.