use clap::Subcommand;
use crate::client;
use crate::config::ProjectConfig;
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error(transparent)]
Client(#[from] crate::client::ClientError),
#[error("{0}")]
Usage(String),
#[error(transparent)]
Config(#[from] boatramp_node::config::ConfigError),
#[error("building the {side} blob backend: {source}")]
BackendBuild {
side: &'static str,
#[source]
source: boatramp_node::Error,
},
#[error("resolving the {side} [serve.s3_credential]: {source}")]
Credential {
side: &'static str,
#[source]
source: boatramp_node::s3_credential::S3CredentialError,
},
#[error("building the {side} [secrets] envelope: {reason}")]
Envelope {
side: &'static str,
reason: String,
},
#[error(transparent)]
Migrate(#[from] boatramp_node::blob_migrate::MigrateError),
#[error("blob drain did not complete — see the drain report above (re-run to resume)")]
DrainIncomplete,
#[error("blob purge did not complete — see the purge report above (re-run to resume)")]
PurgeIncomplete,
#[error(transparent)]
Json(#[from] serde_json::Error),
#[error(
"no migration SOURCE: pass --from <config>, or configure [serve.blob_fallback] in \
{node_config} (the drain one-liner) — {node_config} has no fallback secondary to drain"
)]
NoSource {
node_config: String,
},
#[error(
"refusing to migrate: source and destination resolve to the SAME backend ({identity}); \
pick distinct --from/--to backends"
)]
SameSourceDest {
identity: String,
},
}
type Result<T> = std::result::Result<T, Error>;
#[derive(Debug, clap::Args)]
pub struct BlobArgs {
#[arg(long, env = "BOATRAMP_SERVER", global = true)]
server: Option<String>,
#[command(subcommand)]
command: BlobCommand,
}
#[derive(Debug, Subcommand)]
enum BlobCommand {
Put {
file: std::path::PathBuf,
},
#[cfg(feature = "blob-upload")]
MintUpload(MintUploadArgs),
Migrate(MigrateArgs),
Drain(DrainArgs),
Purge(PurgeArgs),
Status(StatusArgs),
}
#[derive(Debug, clap::Args)]
#[command(group(
clap::ArgGroup::new("purge_mode")
.required(true)
.args(["unreferenced", "drained_source"]),
))]
pub struct PurgeArgs {
#[arg(long)]
unreferenced: bool,
#[arg(long)]
drained_source: bool,
#[arg(long)]
apply: bool,
#[arg(long)]
prefix: Option<String>,
#[arg(long)]
json: bool,
}
#[derive(Debug, clap::Args)]
pub struct StatusArgs {
#[arg(long)]
json: bool,
}
#[derive(Debug, clap::Args)]
pub struct DrainArgs {
#[arg(long)]
dry_run: bool,
#[arg(long)]
concurrency: Option<usize>,
#[arg(long)]
prefix: Option<String>,
#[arg(long)]
json: bool,
}
#[derive(Debug, clap::Args)]
pub struct MigrateArgs {
#[arg(long)]
from: Option<std::path::PathBuf>,
#[arg(long)]
to: Option<std::path::PathBuf>,
#[arg(long, default_value = "boatramp.cfg")]
node_config: std::path::PathBuf,
#[arg(long, default_value_t = 8)]
concurrency: usize,
#[arg(long)]
no_verify: bool,
#[arg(long)]
dry_run: bool,
#[arg(long, default_value = "")]
prefix: String,
#[arg(long)]
json: bool,
}
#[cfg(feature = "blob-upload")]
#[derive(Debug, clap::Args)]
pub struct MintUploadArgs {
#[arg(long)]
site: String,
#[arg(long, short = 'c')]
container: String,
#[arg(long, conflicts_with = "prefix")]
key: Option<String>,
#[arg(long)]
prefix: Option<String>,
#[arg(long = "perms", value_delimiter = ',')]
perms: Vec<String>,
#[arg(long, default_value_t = 900)]
ttl: u64,
#[arg(long)]
content_type: Option<String>,
#[arg(long)]
max_bytes: Option<u64>,
#[arg(long)]
sha256: bool,
#[arg(long, value_enum, default_value_t = Emit::Human)]
emit: Emit,
}
#[cfg(feature = "blob-upload")]
#[derive(Debug, Clone, Copy, PartialEq, Eq, clap::ValueEnum)]
enum Emit {
Human,
Env,
Aws,
Rclone,
Json,
}
pub async fn run(args: BlobArgs, config: &ProjectConfig) -> Result<()> {
if let BlobCommand::Migrate(a) = args.command {
return migrate(a).await;
}
let server = client::resolve_server(args.server, config)?;
let cp = client::ControlPlane::new(
server,
client::http_client(client::token(config).as_deref()),
client::resolve_project(config),
);
match args.command {
BlobCommand::Put { file } => {
let hash = cp.put_file_blob(&file).await?;
println!("{hash}");
}
#[cfg(feature = "blob-upload")]
BlobCommand::MintUpload(a) => mint_upload(&cp, a).await?,
BlobCommand::Drain(a) => drain(&cp, a).await?,
BlobCommand::Purge(a) => purge(&cp, a).await?,
BlobCommand::Status(a) => status(&cp, a).await?,
BlobCommand::Migrate(_) => unreachable!("handled before the control-plane client"),
}
Ok(())
}
async fn drain(cp: &client::ControlPlane, a: DrainArgs) -> Result<()> {
let final_report = cp
.blob_drain(a.dry_run, a.concurrency, a.prefix.as_deref())
.await?;
let verify_failed = final_report.get("missing").is_some();
let errored = final_report.get("error").is_some();
if a.json {
println!("{}", serde_json::to_string_pretty(&final_report)?);
} else {
let message = final_report
.get("message")
.and_then(serde_json::Value::as_str)
.unwrap_or("(no message)");
println!("blob drain: {message}");
}
if verify_failed || errored {
return Err(Error::DrainIncomplete);
}
Ok(())
}
async fn purge(cp: &client::ControlPlane, a: PurgeArgs) -> Result<()> {
let mode = if a.unreferenced {
"unreferenced"
} else {
"drained_source"
};
let final_report = cp.blob_purge(mode, a.apply, a.prefix.as_deref()).await?;
let errored = final_report.get("error").is_some();
if a.json {
println!("{}", serde_json::to_string_pretty(&final_report)?);
} else {
let message = final_report
.get("message")
.and_then(serde_json::Value::as_str)
.unwrap_or("(no message)");
println!("blob purge: {message}");
}
if errored {
return Err(Error::PurgeIncomplete);
}
Ok(())
}
async fn status(cp: &client::ControlPlane, a: StatusArgs) -> Result<()> {
let state = cp.blob_status().await?;
if a.json {
println!("{}", serde_json::to_string_pretty(&state)?);
} else {
let active = state
.get("blob_fallback_active")
.and_then(serde_json::Value::as_bool)
.unwrap_or(false);
if active {
println!(
"blob status: a read-fallback secondary is ATTACHED ([serve.blob_fallback]) — the \
node is mid-migration (TRANSITION mode). Drain it (blob drain / blob purge \
--drained-source), drop [serve.blob_fallback], and restart to finish."
);
} else {
println!(
"blob status: no read-fallback secondary attached — the node is not mid-migration."
);
}
}
Ok(())
}
async fn migrate(a: MigrateArgs) -> Result<()> {
let (source, dest) = resolve_sides(&a).await?;
if !a.json {
println!("blob migrate: source = {}", source.identity);
println!("blob migrate: dest = {}", dest.identity);
}
if source.identity == dest.identity {
return Err(Error::SameSourceDest {
identity: source.identity,
});
}
let is_configured_drain = source.is_configured_fallback && dest.is_configured_primary;
let opts = boatramp_node::blob_migrate::MigrateOptions {
concurrency: a.concurrency,
verify: !a.no_verify,
dry_run: a.dry_run,
prefix: a.prefix.clone(),
on_progress: None,
};
let report = boatramp_node::blob_migrate::migrate(source.storage, dest.storage, &opts).await?;
let drained = report.verified && is_configured_drain;
if a.json {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"source": source.identity,
"dest": dest.identity,
"total_objects": report.total_objects,
"copied_objects": report.copied_objects,
"skipped_objects": report.skipped_objects,
"copied_bytes": report.copied_bytes,
"verified": report.verified,
"dry_run": a.dry_run,
"secondary_drained": drained,
}))?
);
} else {
let verb = if a.dry_run { "would copy" } else { "copied" };
println!(
"blob migrate {}: {} {} object(s), skipped {} present, {} byte(s){}",
if a.dry_run { "(dry-run)" } else { "complete" },
verb,
report.copied_objects,
report.skipped_objects,
report.copied_bytes,
if report.verified {
format!(
"; VERIFY OK: {} object(s) present in destination",
report.total_objects
)
} else {
String::new()
},
);
if drained {
println!(
"SECONDARY FULLY DRAINED — safe to remove [serve].blob_fallback and restart the node."
);
}
}
Ok(())
}
struct ResolvedSide {
storage: std::sync::Arc<dyn boatramp_core::Storage>,
identity: String,
is_configured_fallback: bool,
is_configured_primary: bool,
}
async fn resolve_sides(a: &MigrateArgs) -> Result<(ResolvedSide, ResolvedSide)> {
let source = match a.from.as_ref() {
Some(from) => build_primary_side("source", from).await?,
None => {
let config = boatramp_node::config::ServerConfig::load(&a.node_config, None)?;
let serve = config.serve.clone().unwrap_or_default();
let Some(fb) = serve.blob_fallback.clone() else {
return Err(Error::NoSource {
node_config: a.node_config.display().to_string(),
});
};
build_fallback_side("source", &config, &fb).await?
}
};
let dest_path = a.to.clone().unwrap_or_else(|| a.node_config.clone());
let dest_is_configured_primary = a.to.is_none();
let mut dest = build_primary_side("destination", &dest_path).await?;
dest.is_configured_primary = dest_is_configured_primary;
Ok((source, dest))
}
async fn build_primary_side(
side: &'static str,
config_path: &std::path::Path,
) -> Result<ResolvedSide> {
let config = boatramp_node::config::ServerConfig::load(config_path, None)?;
let serve = config.serve.clone().unwrap_or_default();
let data_dir = serve
.data_dir
.clone()
.unwrap_or_else(|| std::path::PathBuf::from("./data"));
let mut blob_args = boatramp_node::blobs::BlobArgs {
blobs: serve
.blobs
.unwrap_or(boatramp_node::backends::BlobBackend::Fs),
s3_bucket: serve.s3_bucket.clone(),
s3_endpoint: serve.s3_endpoint.clone(),
s3_region: serve.s3_region.clone(),
s3_path_style: serve.s3_path_style,
s3_credential: None,
gcs_bucket: serve.gcs_bucket.clone(),
gcs_endpoint: serve.gcs_endpoint.clone(),
gcs_anonymous: serve.gcs_anonymous,
azure_account: serve.azure_account.clone(),
azure_container: serve.azure_container.clone(),
azure_access_key: serve.azure_access_key.clone(),
azure_emulator: serve.azure_emulator,
};
if let Some(cred_cfg) = serve.s3_credential.clone() {
blob_args.s3_credential =
Some(resolve_side_credential(side, &config, &cred_cfg, &data_dir).await?);
}
let identity = backend_identity(&blob_args, &data_dir);
let built = boatramp_node::blobs::build_blobs(&blob_args, &data_dir, None, None)
.await
.map_err(|source| Error::BackendBuild { side, source })?;
Ok(ResolvedSide {
storage: built.storage,
identity,
is_configured_fallback: false,
is_configured_primary: false,
})
}
async fn build_fallback_side(
side: &'static str,
config: &boatramp_node::config::ServerConfig,
fb: &boatramp_node::config::BlobFallbackConfig,
) -> Result<ResolvedSide> {
let serve = config.serve.clone().unwrap_or_default();
let data_dir = serve
.data_dir
.clone()
.unwrap_or_else(|| std::path::PathBuf::from("./data"));
let mut blob_args = boatramp_node::blobs::BlobArgs {
blobs: fb.blobs.unwrap_or(boatramp_node::backends::BlobBackend::Fs),
s3_bucket: fb.s3_bucket.clone(),
s3_endpoint: fb.s3_endpoint.clone(),
s3_region: fb.s3_region.clone(),
s3_path_style: fb.s3_path_style,
s3_credential: None,
gcs_bucket: fb.gcs_bucket.clone(),
gcs_endpoint: fb.gcs_endpoint.clone(),
gcs_anonymous: fb.gcs_anonymous,
azure_account: fb.azure_account.clone(),
azure_container: fb.azure_container.clone(),
azure_access_key: fb.azure_access_key.clone(),
azure_emulator: fb.azure_emulator,
};
if let Some(cred_cfg) = fb.s3_credential.clone() {
blob_args.s3_credential =
Some(resolve_side_credential(side, config, &cred_cfg, &data_dir).await?);
}
let identity = backend_identity(&blob_args, &data_dir);
let built = boatramp_node::blobs::build_blobs(&blob_args, &data_dir, None, None)
.await
.map_err(|source| Error::BackendBuild { side, source })?;
Ok(ResolvedSide {
storage: built.storage,
identity,
is_configured_fallback: true,
is_configured_primary: false,
})
}
async fn resolve_side_credential(
side: &'static str,
config: &boatramp_node::config::ServerConfig,
cred_cfg: &boatramp_node::config::S3CredentialConfig,
data_dir: &std::path::Path,
) -> Result<boatramp_node::s3_credential::SealedS3Credential> {
use std::sync::Arc;
let posture = config
.security
.clone()
.unwrap_or_default()
.resolve()
.map_err(|e| Error::Envelope {
side,
reason: format!("resolving [security] posture: {e}"),
})?;
let envelope = build_secrets_envelope(side, config.secrets.as_ref(), data_dir)?;
boatramp_node::s3_credential::resolve_s3_credential(
cred_cfg,
Arc::new(boatramp_core::kv::MemoryKv::new()),
envelope,
posture.allow_env_secret_refs,
&boatramp_core::env::SystemEnv,
)
.await
.map_err(|source| Error::Credential { side, source })
}
fn backend_identity(args: &boatramp_node::blobs::BlobArgs, data_dir: &std::path::Path) -> String {
use boatramp_node::backends::BlobBackend;
match args.blobs {
BlobBackend::Fs => format!("fs {}", data_dir.join("blobs").display()),
BlobBackend::S3 => format!(
"s3 bucket={} endpoint={} region={} path_style={}",
args.s3_bucket.as_deref().unwrap_or("?"),
args.s3_endpoint.as_deref().unwrap_or("(default)"),
args.s3_region.as_deref().unwrap_or("(default)"),
args.s3_path_style
),
BlobBackend::Gcs => format!(
"gcs bucket={} endpoint={}",
args.gcs_bucket.as_deref().unwrap_or("?"),
args.gcs_endpoint.as_deref().unwrap_or("(default)")
),
BlobBackend::Azure => format!(
"azure account={} container={}",
args.azure_account.as_deref().unwrap_or("?"),
args.azure_container.as_deref().unwrap_or("?")
),
}
}
fn build_secrets_envelope(
side: &'static str,
secrets: Option<&boatramp_node::config::SecretsConfig>,
data_dir: &std::path::Path,
) -> Result<Option<std::sync::Arc<dyn boatramp_core::envelope::KeyEnvelope>>> {
use boatramp_server::envelope::{EnvelopeSpec, build_envelope};
let Some(cfg) = secrets else {
return Ok(None);
};
let err = |reason: String| Error::Envelope { side, reason };
let spec = match cfg.envelope.as_str() {
"" => EnvelopeSpec::None,
"local" => EnvelopeSpec::Local {
kek_file: cfg
.kek_file
.clone()
.unwrap_or_else(|| data_dir.join("secrets/kek")),
},
"vault" => {
let v = cfg.vault.as_ref().ok_or_else(|| {
err("secrets.envelope = \"vault\" needs a [secrets.vault] section".into())
})?;
let token = std::env::var(&v.token_env)
.map_err(|_| err(format!("Vault token env `{}` is not set", v.token_env)))?;
EnvelopeSpec::Vault {
addr: v.addr.clone(),
key: v.key.clone(),
token,
}
}
other => {
return Err(err(format!(
"unknown secrets.envelope {other:?} (want \"local\" or \"vault\")"
)));
}
};
build_envelope(spec).map_err(|e| err(e.to_string()))
}
#[cfg(feature = "blob-upload")]
async fn mint_upload(cp: &client::ControlPlane, a: MintUploadArgs) -> Result<()> {
if a.key.is_none() && a.prefix.is_none() {
return Err(Error::Usage(
"specify exactly one of --key or --prefix".into(),
));
}
let creds = cp
.mint_upload(
&a.site,
&a.container,
a.key.as_deref(),
a.prefix.as_deref(),
&a.perms,
a.ttl,
a.content_type.as_deref(),
a.max_bytes,
a.sha256,
)
.await?;
render(&creds, a.emit);
Ok(())
}
#[cfg(feature = "blob-upload")]
fn render(c: &client::MintUploadResult, emit: Emit) {
let is_presigned = c.kind == "presigned_put";
match emit {
Emit::Json => {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"kind": c.kind,
"url": c.url,
"method": c.method,
"required_headers": c.required_headers,
"access_key_id": c.access_key_id,
"secret": c.secret,
"session_token": c.session_token,
"endpoint": c.endpoint,
"region": c.region,
"bucket": c.bucket,
"force_path_style": c.force_path_style,
"expires_at": c.expires_at,
"expires_in_secs": c.expires_in_secs,
"enforced": c.enforced,
"advisory": c.advisory,
}))
.unwrap_or_default()
);
}
Emit::Aws => {
if is_presigned {
eprintln!(
"# a presigned-put credential is a single URL, not an SDK profile; use --emit env or json"
);
print_presigned(c);
} else {
println!("[boatramp-upload]");
println!("aws_access_key_id = {}", opt(&c.access_key_id));
println!("aws_secret_access_key = {}", opt(&c.secret));
println!("aws_session_token = {}", opt(&c.session_token));
println!("# region = {}", opt(&c.region));
println!("# endpoint_url = {}", opt(&c.endpoint));
println!("# addressing_style = path");
println!(
"# bucket = {} (expires in {}s)",
opt(&c.bucket),
c.expires_in_secs
);
}
}
Emit::Rclone => {
if is_presigned {
eprintln!(
"# a presigned-put credential is a single URL, not an rclone remote; use --emit env or json"
);
print_presigned(c);
} else {
println!("[boatramp-upload]");
println!("type = s3");
println!("provider = Other");
println!("access_key_id = {}", opt(&c.access_key_id));
println!("secret_access_key = {}", opt(&c.secret));
println!("session_token = {}", opt(&c.session_token));
println!("region = {}", opt(&c.region));
println!("endpoint = {}", opt(&c.endpoint));
println!("force_path_style = true");
}
}
Emit::Env => {
if is_presigned {
print_presigned(c);
} else {
print_env(c);
}
}
Emit::Human => {
if is_presigned {
println!("Minted a presigned PUT credential (single-key, browser-friendly):");
print_presigned(c);
} else {
println!(
"Minted temporary S3 credentials (path-style; expires in {}s):",
c.expires_in_secs
);
println!(" bucket {}", opt(&c.bucket));
println!(" endpoint {}", opt(&c.endpoint));
println!(" region {}", opt(&c.region));
if !c.enforced.is_empty() {
println!(" enforced {}", c.enforced.join(", "));
}
if !c.advisory.is_empty() {
println!(" advisory {}", c.advisory.join(", "));
}
println!();
print_env(c);
}
}
}
}
#[cfg(feature = "blob-upload")]
fn print_env(c: &client::MintUploadResult) {
println!("export AWS_ACCESS_KEY_ID={}", opt(&c.access_key_id));
println!("export AWS_SECRET_ACCESS_KEY={}", opt(&c.secret));
println!("export AWS_SESSION_TOKEN={}", opt(&c.session_token));
println!("export AWS_REGION={}", opt(&c.region));
println!("export AWS_ENDPOINT_URL={}", opt(&c.endpoint));
println!("export AWS_S3_FORCE_PATH_STYLE=true");
println!("# then: aws s3 cp <file> s3://{}/<key>", opt(&c.bucket));
}
#[cfg(feature = "blob-upload")]
fn print_presigned(c: &client::MintUploadResult) {
let url = opt(&c.url);
println!(" method {}", opt(&c.method));
println!(" url {url}");
for (k, v) in &c.required_headers {
println!(" header {k}: {v}");
}
println!(" expires in {}s", c.expires_in_secs);
let hdrs: String = c
.required_headers
.iter()
.map(|(k, v)| format!(" -H '{k}: {v}'"))
.collect();
println!(
"# curl -X {} '{url}'{hdrs} --data-binary @<file>",
opt(&c.method)
);
}
#[cfg(feature = "blob-upload")]
fn opt(v: &Option<String>) -> &str {
v.as_deref().unwrap_or("")
}