use boatramp_core::function::FunctionSummary;
use serde::Deserialize;
use crate::client;
use crate::config::ProjectConfig;
mod build;
mod scaffold;
use build::build_project;
use scaffold::init_project;
#[cfg(test)]
use scaffold::sanitize_crate_name;
#[derive(Debug, thiserror::Error)]
pub enum FunctionError {
#[error(transparent)]
Client(#[from] crate::client::ClientError),
#[error("control-plane request: {0}")]
Http(#[from] reqwest::Error),
#[error("reading request body: {0}")]
Io(#[from] std::io::Error),
#[error("specify exactly one of --cron, --queue, or --blob")]
BadTrigger,
#[error("unknown template language {0:?} (supported: rust, js, python)")]
UnknownLang(String),
#[error("{0} already exists")]
AlreadyExists(std::path::PathBuf),
#[error("build failed (Rust: wasm32-wasip2 target; JS: node/npx; Python: uv/componentize-py)")]
BuildFailed,
#[error("no component produced under target/wasm32-wasip2/release")]
NoComponent,
#[cfg(feature = "handlers")]
#[error("running the component: {0}")]
Harness(String),
#[cfg(feature = "handlers")]
#[error("function test assertion failed")]
HarnessFailed,
}
type Result<T> = std::result::Result<T, FunctionError>;
#[derive(Debug, clap::Args)]
pub struct FunctionArgs {
#[command(subcommand)]
command: FunctionCommand,
}
#[derive(Debug, clap::Subcommand)]
enum FunctionCommand {
Ls {
#[arg(long)]
site: Option<String>,
#[arg(long)]
server: Option<String>,
},
Get {
name: String,
#[arg(long)]
server: Option<String>,
},
Deploy {
name: String,
#[arg(long)]
component: std::path::PathBuf,
#[arg(long)]
runtime: Option<String>,
#[arg(long)]
webhook_secret_env: Option<String>,
#[arg(long)]
server: Option<String>,
},
Rollback {
name: String,
#[arg(long)]
to: String,
#[arg(long)]
server: Option<String>,
},
Alias {
name: String,
label: String,
version: String,
#[arg(long)]
server: Option<String>,
},
Rm {
name: String,
#[arg(long)]
server: Option<String>,
},
Invoke {
name: String,
#[arg(long, conflicts_with = "data_file")]
data: Option<String>,
#[arg(long)]
data_file: Option<std::path::PathBuf>,
#[arg(long)]
content_type: Option<String>,
#[arg(long)]
r#async: bool,
#[arg(long)]
idempotency_key: Option<String>,
#[arg(long)]
version: Option<String>,
#[arg(long)]
server: Option<String>,
},
Invocation {
name: String,
id: String,
#[arg(long)]
server: Option<String>,
},
Usage {
name: String,
#[arg(long)]
server: Option<String>,
},
Trigger(TriggerArgs),
Init {
name: String,
#[arg(long, default_value = "rust")]
lang: String,
#[arg(long)]
dir: Option<std::path::PathBuf>,
},
Build {
#[arg(long)]
dir: Option<std::path::PathBuf>,
},
#[cfg(feature = "handlers")]
Test {
#[arg(long)]
component: std::path::PathBuf,
#[arg(long, default_value = "/")]
path: String,
#[arg(long, default_value = "GET")]
method: String,
#[arg(long)]
data: Option<String>,
#[arg(long)]
content_type: Option<String>,
#[arg(long)]
expect_status: Option<u16>,
#[arg(long)]
expect_body: Option<String>,
},
#[cfg(feature = "handlers")]
Dev {
#[arg(long)]
component: std::path::PathBuf,
#[arg(long, default_value = "8787")]
port: u16,
},
}
#[derive(Debug, clap::Args)]
struct TriggerArgs {
#[command(subcommand)]
command: TriggerCommand,
}
#[derive(Debug, clap::Subcommand)]
enum TriggerCommand {
Add {
name: String,
id: String,
#[arg(long, conflicts_with_all = ["queue", "blob"])]
cron: Option<String>,
#[arg(long, conflicts_with = "blob")]
queue: Option<String>,
#[arg(long)]
blob: Option<String>,
#[arg(long)]
server: Option<String>,
},
Ls {
name: String,
#[arg(long)]
server: Option<String>,
},
Rm {
name: String,
id: String,
#[arg(long)]
server: Option<String>,
},
}
#[derive(Debug, Deserialize)]
struct TriggerView {
id: String,
kind: serde_json::Value,
}
#[derive(Debug, Deserialize)]
struct StoredFunction {
name: String,
active: String,
#[serde(default)]
config: StoredConfig,
}
#[derive(Debug, Default, Deserialize)]
struct StoredConfig {
#[serde(default)]
runtime: String,
}
#[derive(Debug, Deserialize)]
struct InvocationRecord {
id: String,
status: String,
#[serde(default)]
attempts: u32,
#[serde(default)]
result: Option<InvocationResultView>,
}
#[derive(Debug, Deserialize)]
struct InvocationResultView {
status: u16,
}
#[derive(Debug, Default, Deserialize)]
struct UsageView {
function: String,
#[serde(default)]
invocations: u64,
#[serde(default)]
successes: u64,
#[serde(default)]
failures: u64,
#[serde(default)]
duration_ms_total: u64,
#[serde(default)]
bytes_in_total: u64,
#[serde(default)]
bytes_out_total: u64,
}
pub async fn run(args: FunctionArgs, config: &ProjectConfig) -> Result<()> {
let seg = client::project_seg(&client::resolve_project(config), "functions");
match args.command {
FunctionCommand::Ls { site, server } => {
let funcs = fetch(server, site, config).await?;
if funcs.is_empty() {
println!("no functions");
return Ok(());
}
for f in funcs {
println!(
"{} [{}] {} {}",
f.name,
f.runtime,
short(&f.version),
f.triggers.join(", ")
);
}
}
FunctionCommand::Get { name, server } => {
let funcs = fetch(server, None, config).await?;
match funcs.into_iter().find(|f| f.name == name) {
Some(f) => {
println!("{}", f.name);
println!(" runtime: {}", f.runtime);
println!(" version: {}", f.version);
for t in &f.triggers {
println!(" trigger: {t}");
}
}
None => println!("no function {name:?}"),
}
}
FunctionCommand::Deploy {
name,
component,
runtime,
webhook_secret_env,
server,
} => {
let (server, http) = client::connect(server, config)?;
let cp = client::ControlPlane::new(
server.clone(),
http.clone(),
client::resolve_project(config),
);
let hash = cp.put_file_blob(&component).await?;
let mut cfg = serde_json::Map::new();
if let Some(r) = &runtime {
cfg.insert("runtime".to_string(), serde_json::json!(r));
}
if let Some(secret_env) = &webhook_secret_env {
cfg.insert(
"webhook".to_string(),
serde_json::json!({ "secret_env": secret_env }),
);
}
let body = serde_json::json!({
"component": hash,
"config": serde_json::Value::Object(cfg),
"lifecycle": "independent",
});
let f: StoredFunction = http
.put(format!("{server}/api/{seg}/{name}"))
.json(&body)
.send()
.await?
.error_for_status()?
.json()
.await?;
println!(
"deployed {} [{}] {}",
f.name,
f.config.runtime,
short(&f.active)
);
}
FunctionCommand::Rollback { name, to, server } => {
let (server, http) = client::connect(server, config)?;
let f: StoredFunction = http
.post(format!("{server}/api/{seg}/{name}/rollback"))
.json(&serde_json::json!({ "to": to }))
.send()
.await?
.error_for_status()?
.json()
.await?;
println!("rolled {} back to {}", f.name, short(&f.active));
}
FunctionCommand::Alias {
name,
label,
version,
server,
} => {
let (server, http) = client::connect(server, config)?;
http.put(format!("{server}/api/{seg}/{name}/aliases/{label}"))
.json(&serde_json::json!({ "version": version }))
.send()
.await?
.error_for_status()?;
println!("aliased {name}:{label} -> {}", short(&version));
}
FunctionCommand::Rm { name, server } => {
let (server, http) = client::connect(server, config)?;
http.delete(format!("{server}/api/{seg}/{name}"))
.send()
.await?
.error_for_status()?;
println!("removed {name}");
}
FunctionCommand::Invoke {
name,
data,
data_file,
content_type,
r#async,
idempotency_key,
version,
server,
} => {
let (server, http) = client::connect(server, config)?;
let body = read_invoke_body(data, data_file).await?;
let mut qs: Vec<String> = Vec::new();
if r#async {
qs.push("mode=async".to_string());
}
if let Some(v) = &version {
qs.push(format!("version={v}"));
}
let url = if qs.is_empty() {
format!("{server}/api/{seg}/{name}/invoke")
} else {
format!("{server}/api/{seg}/{name}/invoke?{}", qs.join("&"))
};
let mut req = http.post(url).body(body);
if let Some(ct) = &content_type {
req = req.header("content-type", ct.as_str());
}
if let Some(key) = &idempotency_key {
req = req.header("idempotency-key", key.as_str());
}
let resp = req.send().await?;
let status = resp.status();
let bytes = resp.bytes().await?;
if r#async {
match serde_json::from_slice::<InvocationRecord>(&bytes) {
Ok(inv) => println!("queued {} [{}]", inv.id, inv.status),
Err(_) => eprint!("{}", String::from_utf8_lossy(&bytes)),
}
} else {
use std::io::Write;
let _ = std::io::stdout().write_all(&bytes);
if !status.is_success() {
eprintln!("invoke returned HTTP {}", status.as_u16());
}
}
}
FunctionCommand::Invocation { name, id, server } => {
let (server, http) = client::connect(server, config)?;
let inv: InvocationRecord = http
.get(format!("{server}/api/{seg}/{name}/invocations/{id}"))
.send()
.await?
.error_for_status()?
.json()
.await?;
println!("{} [{}] attempts={}", inv.id, inv.status, inv.attempts);
if let Some(result) = &inv.result {
println!(" result: HTTP {}", result.status);
}
}
FunctionCommand::Usage { name, server } => {
let (server, http) = client::connect(server, config)?;
let usage: UsageView = http
.get(format!("{server}/api/{seg}/{name}/usage"))
.send()
.await?
.error_for_status()?
.json()
.await?;
println!("{}", usage.function);
println!(
" invocations: {} ({} ok, {} failed)",
usage.invocations, usage.successes, usage.failures
);
println!(" duration: {} ms total", usage.duration_ms_total);
println!(
" bytes: {} in / {} out",
usage.bytes_in_total, usage.bytes_out_total
);
}
FunctionCommand::Trigger(args) => run_trigger(args, config).await?,
FunctionCommand::Init { name, lang, dir } => init_project(&name, &lang, dir)?,
FunctionCommand::Build { dir } => build_project(dir).await?,
#[cfg(feature = "handlers")]
FunctionCommand::Test {
component,
path,
method,
data,
content_type,
expect_status,
expect_body,
} => {
harness::test_component(
component,
&path,
&method,
data,
content_type,
expect_status,
expect_body,
)
.await?;
}
#[cfg(feature = "handlers")]
FunctionCommand::Dev { component, port } => harness::dev_serve(component, port).await?,
}
Ok(())
}
async fn run_trigger(args: TriggerArgs, config: &ProjectConfig) -> Result<()> {
let seg = client::project_seg(&client::resolve_project(config), "functions");
match args.command {
TriggerCommand::Add {
name,
id,
cron,
queue,
blob,
server,
} => {
let (server, http) = client::connect(server, config)?;
let kind = match (cron, queue, blob) {
(Some(schedule), None, None) => {
serde_json::json!({ "type": "cron", "schedule": schedule })
}
(None, Some(topic), None) => {
serde_json::json!({ "type": "queue", "topic": topic })
}
(None, None, Some(prefix)) => {
serde_json::json!({ "type": "blob", "prefix": prefix })
}
_ => return Err(FunctionError::BadTrigger),
};
http.put(format!("{server}/api/{seg}/{name}/triggers/{id}"))
.json(&kind)
.send()
.await?
.error_for_status()?;
println!("added trigger {name}/{id}");
}
TriggerCommand::Ls { name, server } => {
let (server, http) = client::connect(server, config)?;
let list: Vec<TriggerView> = http
.get(format!("{server}/api/{seg}/{name}/triggers"))
.send()
.await?
.error_for_status()?
.json()
.await?;
if list.is_empty() {
println!("no triggers");
return Ok(());
}
for t in list {
let kind = t.kind.get("type").and_then(|v| v.as_str()).unwrap_or("?");
println!("{} [{}]", t.id, kind);
}
}
TriggerCommand::Rm { name, id, server } => {
let (server, http) = client::connect(server, config)?;
http.delete(format!("{server}/api/{seg}/{name}/triggers/{id}"))
.send()
.await?
.error_for_status()?;
println!("removed trigger {name}/{id}");
}
}
Ok(())
}
async fn read_invoke_body(
data: Option<String>,
data_file: Option<std::path::PathBuf>,
) -> Result<Vec<u8>> {
if let Some(inline) = data {
return Ok(inline.into_bytes());
}
let Some(path) = data_file else {
return Ok(Vec::new());
};
if path.as_os_str() == "-" {
use tokio::io::AsyncReadExt;
let mut buf = Vec::new();
tokio::io::stdin().read_to_end(&mut buf).await?;
Ok(buf)
} else {
Ok(tokio::fs::read(&path).await?)
}
}
#[cfg(feature = "handlers")]
mod harness {
use super::{FunctionError, Result};
use boatramp_handlers::{Bindings, HandlerEngine, Limits};
use http_body_util::{BodyExt, Full};
const LOCAL_HASH: &str = "function-local";
#[allow(clippy::too_many_arguments)]
pub(super) async fn test_component(
component: std::path::PathBuf,
path: &str,
method: &str,
data: Option<String>,
content_type: Option<String>,
expect_status: Option<u16>,
expect_body: Option<String>,
) -> Result<()> {
let wasm = tokio::fs::read(&component).await?;
let engine = HandlerEngine::new(Limits::default(), 4).map_err(harness_err)?;
let (status, body) = run_once(
&engine,
&wasm,
method,
path,
content_type.as_deref(),
data.unwrap_or_default().into_bytes(),
)
.await?;
let text = String::from_utf8_lossy(&body);
println!("HTTP {status}");
print!("{text}");
if !text.ends_with('\n') {
println!();
}
let mut ok = true;
if let Some(exp) = expect_status {
if status != exp {
eprintln!("FAIL: expected status {exp}, got {status}");
ok = false;
}
}
if let Some(sub) = &expect_body {
if !text.contains(sub.as_str()) {
eprintln!("FAIL: body does not contain {sub:?}");
ok = false;
}
}
if ok {
println!("ok");
Ok(())
} else {
Err(FunctionError::HarnessFailed)
}
}
pub(super) async fn dev_serve(component: std::path::PathBuf, port: u16) -> Result<()> {
use hyper::server::conn::http1;
use hyper::service::service_fn;
use hyper_util::rt::TokioIo;
use std::sync::Arc;
let wasm = Arc::new(tokio::fs::read(&component).await?);
let engine = Arc::new(HandlerEngine::new(Limits::default(), 16).map_err(harness_err)?);
let listener = tokio::net::TcpListener::bind(("127.0.0.1", port)).await?;
println!(
"serving {} on http://127.0.0.1:{port} (Ctrl-C to stop)",
component.display()
);
loop {
let (stream, _) = listener.accept().await?;
let io = TokioIo::new(stream);
let engine = engine.clone();
let wasm = wasm.clone();
tokio::spawn(async move {
let service = service_fn(move |req: http::Request<hyper::body::Incoming>| {
let engine = engine.clone();
let wasm = wasm.clone();
async move {
let response = match engine
.serve(LOCAL_HASH, &wasm, req, Bindings::new("fn/local"))
.await
{
Ok(resp) => {
let (parts, body) = resp.into_parts();
let bytes = body
.collect()
.await
.map(http_body_util::Collected::to_bytes)
.unwrap_or_default();
http::Response::from_parts(parts, Full::new(bytes))
}
Err(err) => http::Response::builder()
.status(500)
.body(Full::new(bytes::Bytes::from(format!(
"function error: {err:?}\n"
))))
.unwrap(),
};
Ok::<_, std::convert::Infallible>(response)
}
});
let _ = http1::Builder::new().serve_connection(io, service).await;
});
}
}
async fn run_once(
engine: &HandlerEngine,
wasm: &[u8],
method: &str,
path: &str,
content_type: Option<&str>,
body: Vec<u8>,
) -> Result<(u16, Vec<u8>)> {
let mut builder = http::Request::builder()
.method(method)
.uri(format!("http://localhost{path}"));
if let Some(ct) = content_type {
builder = builder.header("content-type", ct);
}
let request = builder
.body(Full::new(bytes::Bytes::from(body)))
.map_err(harness_err)?;
let response = engine
.serve(LOCAL_HASH, wasm, request, Bindings::new("fn/local"))
.await
.map_err(|e| FunctionError::Harness(format!("{e:?}")))?;
let status = response.status().as_u16();
let bytes = response
.into_body()
.collect()
.await
.map_err(|e| FunctionError::Harness(format!("{e:?}")))?
.to_bytes();
Ok((status, bytes.to_vec()))
}
fn harness_err<E: std::fmt::Display>(err: E) -> FunctionError {
FunctionError::Harness(err.to_string())
}
}
async fn fetch(
server: Option<String>,
site: Option<String>,
config: &ProjectConfig,
) -> Result<Vec<FunctionSummary>> {
let server = client::resolve_server(server, config)?;
let http = client::http_client(client::token(config).as_deref());
let seg = client::project_seg(&client::resolve_project(config), "functions");
let url = match &site {
Some(s) => format!("{server}/api/{seg}?site={s}"),
None => format!("{server}/api/{seg}"),
};
Ok(http
.get(url)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
fn short(id: &str) -> &str {
let id = id.strip_prefix("sha256:").unwrap_or(id);
&id[..id.len().min(12)]
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn sanitizes_crate_names() {
assert_eq!(sanitize_crate_name("My Cool Fn"), "my-cool-fn");
assert_eq!(sanitize_crate_name("resize_images"), "resize-images");
assert_eq!(sanitize_crate_name(" --Foo.Bar-- "), "foo-bar");
assert_eq!(sanitize_crate_name("!!!"), "function");
}
#[test]
fn init_scaffolds_a_buildable_tree() {
let root = std::env::temp_dir().join(format!("boatramp-init-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("My Cool Fn", "rust", Some(root.clone())).unwrap();
let proj = root.join("my-cool-fn");
let cargo = std::fs::read_to_string(proj.join("Cargo.toml")).unwrap();
assert!(cargo.contains("name = \"my-cool-fn\""));
assert!(!cargo.contains("BOATRAMP_FUNCTION_NAME"));
assert!(!proj.join("Cargo.toml.tmpl").exists());
assert!(proj.join("src/lib.rs").exists());
assert!(proj.join("wit/handler.wit").exists());
assert!(matches!(
init_project("x", "cobol", Some(root.clone())),
Err(FunctionError::UnknownLang(_))
));
assert!(matches!(
init_project("My Cool Fn", "rust", Some(root.clone())),
Err(FunctionError::AlreadyExists(_))
));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn init_scaffolds_a_js_tree() {
let root = std::env::temp_dir().join(format!("boatramp-initjs-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("My JS Fn", "js", Some(root.clone())).unwrap();
let proj = root.join("my-js-fn");
let pkg = std::fs::read_to_string(proj.join("package.json")).unwrap();
assert!(pkg.contains("\"name\": \"my-js-fn\""));
assert!(!pkg.contains("BOATRAMP_FUNCTION_NAME"));
assert!(!proj.join("package.json.tmpl").exists());
assert!(proj.join("handler.js").exists());
assert!(proj.join("wit/handler.wit").exists());
assert!(init_project("Other", "javascript", Some(root.clone())).is_ok());
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn init_scaffolds_a_python_tree() {
let root = std::env::temp_dir().join(format!("boatramp-initpy-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("My Py Fn", "python", Some(root.clone())).unwrap();
let proj = root.join("my-py-fn");
let pyproject = std::fs::read_to_string(proj.join("pyproject.toml")).unwrap();
assert!(pyproject.contains("name = \"my-py-fn\""));
assert!(!pyproject.contains("BOATRAMP_FUNCTION_NAME"));
assert!(!proj.join("pyproject.toml.tmpl").exists());
assert!(proj.join("app.py").exists());
assert!(proj.join("wit/handler.wit").exists());
assert!(init_project("Other", "py", Some(root.clone())).is_ok());
assert!(matches!(
init_project("x", "ruby", Some(root.clone())),
Err(FunctionError::UnknownLang(_))
));
let _ = std::fs::remove_dir_all(&root);
}
#[cfg(feature = "handlers")]
#[tokio::test]
async fn harness_runs_a_component_and_asserts() {
const HTTP_200: &[u8] =
include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
let tmp =
std::env::temp_dir().join(format!("boatramp-harness-{}.wasm", std::process::id()));
std::fs::write(&tmp, HTTP_200).unwrap();
harness::test_component(
tmp.clone(),
"/",
"GET",
None,
None,
Some(200),
Some("hello from boatramp".into()),
)
.await
.unwrap();
assert!(matches!(
harness::test_component(tmp.clone(), "/", "GET", None, None, Some(404), None).await,
Err(FunctionError::HarnessFailed)
));
let _ = std::fs::remove_file(&tmp);
}
#[tokio::test]
#[ignore = "compiles a wasm component; run with --ignored / in the flake check"]
async fn init_then_build_produces_a_component() {
let root = std::env::temp_dir().join(format!("boatramp-roundtrip-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("roundtrip demo", "rust", Some(root.clone())).unwrap();
let proj = root.join("roundtrip-demo");
build_project(Some(proj.clone())).await.unwrap();
let wasm = proj.join("target/wasm32-wasip2/release/roundtrip_demo.wasm");
assert!(wasm.exists(), "expected a built component at {wasm:?}");
let bytes = std::fs::read(&wasm).unwrap();
assert_eq!(&bytes[..4], b"\0asm");
let _ = std::fs::remove_dir_all(&root);
}
#[tokio::test]
#[ignore = "runs jco via npx (network + slow); run with --ignored / in the flake check"]
async fn init_then_build_js_produces_a_component() {
let root =
std::env::temp_dir().join(format!("boatramp-jsroundtrip-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("js roundtrip", "js", Some(root.clone())).unwrap();
let proj = root.join("js-roundtrip");
build_project(Some(proj.clone())).await.unwrap();
let wasm = proj.join("js-roundtrip.wasm");
assert!(wasm.exists(), "expected a built component at {wasm:?}");
let bytes = std::fs::read(&wasm).unwrap();
assert_eq!(&bytes[..4], b"\0asm");
let _ = std::fs::remove_dir_all(&root);
}
#[tokio::test]
#[ignore = "runs componentize-py via uvx (network + slow); run with --ignored"]
async fn init_then_build_python_produces_a_component() {
let root =
std::env::temp_dir().join(format!("boatramp-pyroundtrip-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
init_project("py roundtrip", "python", Some(root.clone())).unwrap();
let proj = root.join("py-roundtrip");
build_project(Some(proj.clone())).await.unwrap();
let wasm = proj.join("py-roundtrip.wasm");
assert!(wasm.exists(), "expected a built component at {wasm:?}");
let bytes = std::fs::read(&wasm).unwrap();
assert_eq!(&bytes[..4], b"\0asm");
let _ = std::fs::remove_dir_all(&root);
}
}