use boatramp_core::time::now_unix;
use std::collections::BTreeMap;
use std::sync::Arc;
use boatramp_core::config::SiteConfig;
use boatramp_core::cose::{self, LocalSigner, PopClaims};
use boatramp_core::deploy::{DeploymentList, Manifest};
use boatramp_core::domain_verify::DomainVerification;
use serde::{Deserialize, Serialize};
use crate::config::ProjectConfig;
#[derive(Debug, thiserror::Error)]
pub enum ClientError {
#[error("no server configured; pass --server or set publish.server")]
NoServer,
#[error("no site configured; pass --site or set publish.site")]
NoSite,
#[error("control-plane request: {0}")]
Http(#[from] reqwest::Error),
#[error(transparent)]
Io(#[from] std::io::Error),
}
type Result<T> = std::result::Result<T, ClientError>;
pub fn token(config: &ProjectConfig) -> Option<String> {
std::env::var("BOATRAMP_TOKEN")
.ok()
.filter(|token| !token.is_empty())
.or_else(|| config.publish.token.clone())
}
pub fn http_client(token: Option<&str>) -> ApiClient {
let holder = std::env::var("BOATRAMP_TOKEN_HOLDER_KEY")
.ok()
.filter(|v| !v.is_empty());
let origin = std::env::var("BOATRAMP_POP_ORIGIN")
.ok()
.filter(|v| !v.is_empty());
let server_pubkey = std::env::var("BOATRAMP_SERVER_PUBKEY")
.ok()
.filter(|v| !v.is_empty());
build_client(
token,
holder.as_deref(),
origin.as_deref(),
server_pubkey.as_deref(),
)
}
pub fn build_client(
token: Option<&str>,
holder_key: Option<&str>,
origin: Option<&str>,
server_pubkey: Option<&str>,
) -> ApiClient {
let mut builder = reqwest::Client::builder();
if let Some(token) = token {
if let Ok(value) = reqwest::header::HeaderValue::from_str(&format!("Bearer {token}")) {
let mut headers = reqwest::header::HeaderMap::new();
headers.insert(reqwest::header::AUTHORIZATION, value);
builder = builder.default_headers(headers);
}
}
if let Some(hex) = server_pubkey {
if let Ok(spki) = boatramp_rpktls::parse_public_key(hex.trim()) {
let trust = boatramp_rpktls::TrustSet::from_map(std::collections::BTreeMap::from([(
0u64, spki,
)]));
if let Ok(config) = boatramp_rpktls::client_config_server_auth(trust, 0) {
builder = builder.use_preconfigured_tls(config);
}
}
}
let inner = builder.build().unwrap_or_default();
let pop = match (token, holder_key, origin) {
(Some(token), Some(holder), Some(origin)) => LocalSigner::from_private_hex(holder.trim())
.ok()
.map(|holder| {
Arc::new(PopSigner {
holder,
token: token.to_string(),
origin: origin.to_string(),
})
}),
_ => None,
};
ApiClient { inner, pop }
}
#[derive(Clone)]
pub struct ApiClient {
inner: reqwest::Client,
pop: Option<Arc<PopSigner>>,
}
impl ApiClient {
fn request<U: reqwest::IntoUrl>(&self, method: reqwest::Method, url: U) -> ApiRequestBuilder {
ApiRequestBuilder {
inner: self.inner.request(method, url),
client: self.inner.clone(),
pop: self.pop.clone(),
}
}
pub fn get<U: reqwest::IntoUrl>(&self, url: U) -> ApiRequestBuilder {
self.request(reqwest::Method::GET, url)
}
pub fn post<U: reqwest::IntoUrl>(&self, url: U) -> ApiRequestBuilder {
self.request(reqwest::Method::POST, url)
}
pub fn put<U: reqwest::IntoUrl>(&self, url: U) -> ApiRequestBuilder {
self.request(reqwest::Method::PUT, url)
}
pub fn delete<U: reqwest::IntoUrl>(&self, url: U) -> ApiRequestBuilder {
self.request(reqwest::Method::DELETE, url)
}
}
pub struct ApiRequestBuilder {
inner: reqwest::RequestBuilder,
client: reqwest::Client,
pop: Option<Arc<PopSigner>>,
}
impl ApiRequestBuilder {
pub fn json<T: Serialize + ?Sized>(mut self, json: &T) -> Self {
self.inner = self.inner.json(json);
self
}
pub fn body<T: Into<reqwest::Body>>(mut self, body: T) -> Self {
self.inner = self.inner.body(body);
self
}
pub fn query<T: Serialize + ?Sized>(mut self, query: &T) -> Self {
self.inner = self.inner.query(query);
self
}
pub fn header(mut self, key: &str, value: &str) -> Self {
self.inner = self.inner.header(key, value);
self
}
pub async fn send(self) -> reqwest::Result<reqwest::Response> {
let mut request = self.inner.build()?;
if let Some(pop) = &self.pop {
pop.sign(&mut request).await;
}
self.client.execute(request).await
}
}
struct PopSigner {
holder: LocalSigner,
token: String,
origin: String,
}
impl PopSigner {
async fn sign(&self, request: &mut reqwest::Request) {
let bh = request
.body()
.and_then(reqwest::Body::as_bytes)
.filter(|b| !b.is_empty() && b.len() <= cose::POP_MAX_BODY_HASH_BYTES)
.map(cose::pop_sha256_hex);
let claims = PopClaims {
htm: request.method().as_str().to_string(),
htp: cose::canon_pop_path(request.url().path()),
aud: self.origin.clone(),
ath: cose::pop_sha256_hex(self.token.as_bytes()),
bh,
};
let Ok(proof) = cose::mint_pop(&claims, &self.holder, now_unix()).await else {
return;
};
if let Ok(value) = reqwest::header::HeaderValue::from_str(&proof) {
request.headers_mut().insert(
reqwest::header::HeaderName::from_static("boatramp-pop"),
value,
);
}
}
}
pub fn resolve_project(config: &ProjectConfig) -> String {
config
.publish
.project
.clone()
.filter(|s| !s.is_empty())
.unwrap_or_else(|| boatramp_core::project::DEFAULT_PROJECT.to_string())
}
pub fn project_seg(project: &str, collection: &str) -> String {
if project == boatramp_core::project::DEFAULT_PROJECT {
collection.to_string()
} else {
format!("projects/{project}/{collection}")
}
}
pub fn resolve_server(server: Option<String>, config: &ProjectConfig) -> Result<String> {
let server = server
.or_else(|| config.publish.server.clone())
.ok_or(ClientError::NoServer)?;
Ok(server.trim_end_matches('/').to_string())
}
pub fn connect(server: Option<String>, config: &ProjectConfig) -> Result<(String, ApiClient)> {
let server = resolve_server(server, config)?;
let http = http_client(token(config).as_deref());
Ok((server, http))
}
pub fn resolve_target(
server: Option<String>,
site: Option<String>,
config: &ProjectConfig,
) -> Result<(String, String)> {
let server = resolve_server(server, config)?;
let site = site
.or_else(|| config.publish.site.clone())
.ok_or(ClientError::NoSite)?;
Ok((server, site))
}
fn host_segment(host: &str) -> String {
host.replace('*', "%2A")
}
pub use boatramp_core::domain_verify::CheckResult;
pub use boatramp_core::logs::{LogEntry, LogsResponse};
pub fn is_blob_hash(s: &str) -> bool {
s.len() == 64
&& s.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
}
pub async fn hash_file(path: &std::path::Path) -> Result<String> {
use sha2::{Digest, Sha256};
use tokio::io::AsyncReadExt;
let mut file = tokio::fs::File::open(path).await?;
let mut hasher = Sha256::new();
let mut buf = vec![0u8; 64 * 1024];
loop {
let n = file.read(&mut buf).await?;
if n == 0 {
break;
}
hasher.update(&buf[..n]);
}
Ok(hex::encode(hasher.finalize()))
}
fn sanitize(url: &str) -> String {
url.rsplit('/')
.find(|s| !s.is_empty())
.unwrap_or("download")
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '.' || c == '-' {
c
} else {
'_'
}
})
.take(64)
.collect()
}
#[derive(Debug, Deserialize)]
pub struct CreateDeploymentResponse {
pub id: String,
pub missing: Vec<String>,
}
pub struct ControlPlane {
http: ApiClient,
base: String,
project: String,
}
impl ControlPlane {
pub fn new(base: String, http: ApiClient, project: String) -> Self {
Self {
http,
base,
project,
}
}
fn sites_seg(&self) -> String {
project_seg(&self.project, "sites")
}
fn functions_seg(&self) -> String {
project_seg(&self.project, "functions")
}
fn compute_seg(&self) -> String {
project_seg(&self.project, "compute")
}
pub async fn fetch_manifest(&self, site: &str, id: &str) -> Result<Manifest> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/deployments/{id}"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn fetch_deployments(&self, site: &str) -> Result<DeploymentList> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/deployments"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn fetch_site_config(&self, site: &str) -> Result<SiteConfig> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/config"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn put_site_config(&self, site: &str, config: &SiteConfig) -> Result<()> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
client
.put(format!("{server}/api/{seg}/{site}/config"))
.json(config)
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn start_domain_verification(
&self,
site: &str,
host: &str,
method: Option<&str>,
) -> Result<DomainVerification> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
let mut url = format!(
"{server}/api/{seg}/{site}/domains/{}/verification",
host_segment(host)
);
if let Some(method) = method {
url.push_str(&format!("?method={method}"));
}
Ok(client
.post(url)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn check_domain_verification(&self, site: &str, host: &str) -> Result<CheckResult> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.post(format!(
"{server}/api/{seg}/{site}/domains/{}/verification/check",
host_segment(host)
))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn attach_domain_unverified(&self, site: &str, host: &str) -> Result<String> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.post(format!(
"{server}/api/{seg}/{site}/domains/{}/attach-unverified",
host_segment(host)
))
.send()
.await?
.error_for_status()?
.text()
.await?)
}
pub async fn remove_domain_verification(&self, site: &str, host: &str) -> Result<()> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
client
.delete(format!(
"{server}/api/{seg}/{site}/domains/{}/verification",
host_segment(host)
))
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn list_domain_verifications(&self, site: &str) -> Result<Vec<DomainVerification>> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/domain-verifications"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn activate(&self, site: &str, id: &str) -> Result<()> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
client
.post(format!(
"{server}/api/{seg}/{site}/deployments/{id}/activate"
))
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn create_deployment(
&self,
site: &str,
manifest: &Manifest,
query: &[(&str, String)],
) -> Result<CreateDeploymentResponse> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.post(format!("{server}/api/{seg}/{site}/deployments"))
.query(query)
.json(manifest)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn upload_blob_source(
&self,
hash: &str,
source: &crate::sync::BlobSource,
) -> Result<()> {
use crate::sync::BlobSource;
let Self {
http, base: server, ..
} = self;
let body = match source {
BlobSource::File(path) => {
let file = tokio::fs::File::open(path).await?;
reqwest::Body::wrap_stream(tokio_util::io::ReaderStream::new(file))
}
BlobSource::Memory(bytes) => reqwest::Body::from(bytes.clone()),
};
http.put(format!("{server}/api/blobs/{hash}"))
.body(body)
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn deploy_function(
&self,
name: &str,
body: &serde_json::Value,
) -> Result<serde_json::Value> {
let seg = self.functions_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.put(format!("{server}/api/{seg}/{name}"))
.json(body)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn put_compute(
&self,
name: &str,
body: &serde_json::Value,
) -> Result<serde_json::Value> {
let seg = self.compute_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.put(format!("{server}/api/{seg}/{name}"))
.json(body)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn set_alias(&self, site: &str, name: &str, id: &str) -> Result<()> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
#[derive(Serialize)]
struct SetAlias<'a> {
id: &'a str,
}
client
.put(format!("{server}/api/{seg}/{site}/aliases/{name}"))
.json(&SetAlias { id })
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn list_aliases(&self, site: &str) -> Result<BTreeMap<String, String>> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/aliases"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn remove_alias(&self, site: &str, name: &str) -> Result<()> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
client
.delete(format!("{server}/api/{seg}/{site}/aliases/{name}"))
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn fetch_logs(
&self,
site: &str,
limit: usize,
after: u64,
stream: Option<&str>,
) -> Result<LogsResponse> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
let mut url =
format!("{server}/api/{seg}/{site}/_boatramp/logs?limit={limit}&after={after}");
if let Some(stream) = stream {
url.push_str("&stream=");
url.push_str(stream);
}
Ok(client
.get(url)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn fetch_handler_stats(&self, site: &str) -> Result<serde_json::Value> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/{seg}/{site}/_boatramp/handlers"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn operate_dlq(
&self,
site: &str,
topic: &str,
alias: Option<&str>,
action: &str,
) -> Result<usize> {
let seg = self.sites_seg();
let Self {
http: client,
base: server,
..
} = self;
#[derive(Serialize)]
struct Request<'a> {
topic: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
alias: Option<&'a str>,
action: &'a str,
}
#[derive(Deserialize)]
struct DlqResponse {
affected: usize,
}
let resp: DlqResponse = client
.post(format!("{server}/api/{seg}/{site}/_boatramp/dlq"))
.json(&Request {
topic,
alias,
action,
})
.send()
.await?
.error_for_status()?
.json()
.await?;
Ok(resp.affected)
}
pub async fn upload_blob(&self, hash: &str, path: &std::path::Path) -> Result<()> {
let Self {
http, base: server, ..
} = self;
let file = tokio::fs::File::open(path).await?;
let body = reqwest::Body::wrap_stream(tokio_util::io::ReaderStream::new(file));
http.put(format!("{server}/api/blobs/{hash}"))
.body(body)
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn put_file_blob(&self, path: &std::path::Path) -> Result<String> {
let hash = hash_file(path).await?;
self.upload_blob(&hash, path).await?;
Ok(hash)
}
pub async fn resolve_artifact(&self, value: &str) -> Result<String> {
if is_blob_hash(value) {
return Ok(value.to_string());
}
if value.starts_with("http://") || value.starts_with("https://") {
use tokio::io::AsyncWriteExt;
let mut resp = self.http.get(value).send().await?.error_for_status()?;
let tmp = std::env::temp_dir().join(format!("boatramp-artifact-{}", sanitize(value)));
let mut out = tokio::fs::File::create(&tmp).await?;
while let Some(chunk) = resp.chunk().await? {
out.write_all(&chunk).await?;
}
out.flush().await?;
drop(out);
let hash = self.put_file_blob(&tmp).await?;
let _ = tokio::fs::remove_file(&tmp).await;
return Ok(hash);
}
self.put_file_blob(std::path::Path::new(value)).await
}
pub async fn list_projects(&self) -> Result<Vec<serde_json::Value>> {
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/projects"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn create_project(&self, body: &serde_json::Value) -> Result<serde_json::Value> {
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.post(format!("{server}/api/projects"))
.json(body)
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn get_project(&self, name: &str) -> Result<serde_json::Value> {
let Self {
http: client,
base: server,
..
} = self;
Ok(client
.get(format!("{server}/api/projects/{name}"))
.send()
.await?
.error_for_status()?
.json()
.await?)
}
pub async fn delete_project(&self, name: &str) -> Result<()> {
let Self {
http: client,
base: server,
..
} = self;
client
.delete(format!("{server}/api/projects/{name}"))
.send()
.await?
.error_for_status()?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn blob_hash_detection_is_exact() {
let hash = "a".repeat(64);
assert!(is_blob_hash(&hash), "64 lowercase hex is a blob hash");
assert!(is_blob_hash(
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
));
assert!(!is_blob_hash(&"a".repeat(63)));
assert!(!is_blob_hash(&"a".repeat(65)));
assert!(
!is_blob_hash(&"A".repeat(64)),
"uppercase is treated as a path, not a hash"
);
assert!(!is_blob_hash("./vmlinux"));
assert!(!is_blob_hash("https://example.com/vmlinux"));
assert!(!is_blob_hash(&"g".repeat(64)), "g is not a hex digit");
}
#[test]
fn sites_seg_is_legacy_for_default_project() {
let cp = |project: &str| {
ControlPlane::new(
"https://cp.example".into(),
build_client(None, None, None, None),
project.into(),
)
};
assert_eq!(
cp(boatramp_core::project::DEFAULT_PROJECT).sites_seg(),
"sites"
);
assert_eq!(cp("acme").sites_seg(), "projects/acme/sites");
}
#[test]
fn resolve_project_falls_back_to_default() {
use crate::config::ProjectConfig;
let mut config = ProjectConfig::default();
assert_eq!(
resolve_project(&config),
boatramp_core::project::DEFAULT_PROJECT
);
config.publish.project = Some(String::new());
assert_eq!(
resolve_project(&config),
boatramp_core::project::DEFAULT_PROJECT
);
config.publish.project = Some("acme".into());
assert_eq!(resolve_project(&config), "acme");
}
#[test]
fn sanitize_url_to_temp_fragment() {
assert_eq!(
sanitize("https://example.com/path/vmlinux-6.1.bin"),
"vmlinux-6.1.bin"
);
assert_eq!(sanitize("https://example.com/a b?c=d"), "a_b_c_d");
assert_eq!(sanitize("https://example.com/"), "example.com");
}
use boatramp_core::authz::GrantedRole;
use boatramp_core::cose::{Claims, Signer, TokenAlg};
use boatramp_core::kv::{KvStore, MemoryKv};
use boatramp_server::{require_auth, Auth};
const POP_ORIGIN: &str = "https://cp.example.test";
async fn spawn_guarded(auth: Auth) -> std::net::SocketAddr {
let app = axum::Router::new()
.route("/api/sites", axum::routing::get(|| async { "ok" }))
.layer(axum::middleware::from_fn_with_state(auth, require_auth));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
addr
}
#[tokio::test]
async fn pop_signing_client_round_trips_against_require_auth() {
let root = LocalSigner::generate(TokenAlg::Es256);
let holder = LocalSigner::generate(TokenAlg::Es256);
let now = now_unix();
let claims = Claims {
roles: vec![GrantedRole::global("admin")],
kind: "role".into(),
ttl_secs: Some(3600),
now_unix: now,
};
let token = cose::mint_delegatable(&claims, &holder.public_key(), &root)
.await
.unwrap();
let holder_priv = holder.private_hex();
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let auth = Auth::with_key(root.public_key(), kv).with_pop(Some(POP_ORIGIN.into()), false);
let addr = spawn_guarded(auth).await;
let url = format!("http://127.0.0.1:{}/api/sites", addr.port());
let signed = build_client(Some(&token), Some(&holder_priv), Some(POP_ORIGIN), None);
let resp = signed.get(&url).send().await.unwrap();
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"signed request → 200"
);
let plain = build_client(Some(&token), None, None, None);
let resp = plain.get(&url).send().await.unwrap();
assert_eq!(
resp.status(),
reqwest::StatusCode::UNAUTHORIZED,
"missing proof → 401"
);
let wrong_origin = build_client(
Some(&token),
Some(&holder_priv),
Some("https://evil.example.test"),
None,
);
let resp = wrong_origin.get(&url).send().await.unwrap();
assert_eq!(
resp.status(),
reqwest::StatusCode::UNAUTHORIZED,
"wrong-origin proof → 401"
);
}
}