use std::sync::{Arc, Mutex};
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use boatramp_core::cose::{self, TokenPublicKey};
use boatramp_rpktls::RpkIdentity;
const TICKET_MAGIC: &str = "brjoin1";
#[derive(Debug, thiserror::Error)]
pub enum JoinError {
#[error("invalid join ticket: {0}")]
Ticket(String),
#[error("join token: {0}")]
Token(String),
#[error("possession proof: {0}")]
Proof(String),
#[error("member assertion: {0}")]
Member(String),
#[error("seed {0}")]
Seed(String),
#[error("join refused by {seed}: {status} {body}")]
Refused {
seed: String,
status: u16,
body: String,
},
#[error(transparent)]
Http(#[from] reqwest::Error),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct JoinTicket {
pub seeds: Vec<String>,
pub root_pubkeys: Vec<String>,
pub token: String,
}
impl JoinTicket {
pub fn encode(&self) -> Result<String, JoinError> {
let json = serde_json::to_vec(self).map_err(|e| JoinError::Ticket(e.to_string()))?;
let b64 = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(json);
Ok(format!("{TICKET_MAGIC}.{b64}"))
}
pub fn decode(blob: &str) -> Result<Self, JoinError> {
let b64 = blob
.trim()
.strip_prefix(TICKET_MAGIC)
.and_then(|s| s.strip_prefix('.'))
.ok_or_else(|| JoinError::Ticket("not a boatramp join ticket (bad prefix)".into()))?;
let json = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(b64)
.map_err(|e| JoinError::Ticket(e.to_string()))?;
let ticket: Self =
serde_json::from_slice(&json).map_err(|e| JoinError::Ticket(e.to_string()))?;
if ticket.seeds.is_empty() || ticket.root_pubkeys.is_empty() || ticket.token.is_empty() {
return Err(JoinError::Ticket(
"ticket is missing seeds, root pubkeys, or token".into(),
));
}
Ok(ticket)
}
pub fn roots(&self) -> Result<Vec<TokenPublicKey>, JoinError> {
self.root_pubkeys
.iter()
.map(|s| {
TokenPublicKey::from_hex(s.trim())
.map_err(|e| JoinError::Ticket(format!("bad root pubkey: {e}")))
})
.collect()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AdoptedMember {
pub node_id: u64,
pub mesh_pubkey_hex: String,
pub mesh_addr: Option<String>,
}
pub fn verify_members(
assertions: &[String],
addrs: &std::collections::BTreeMap<u64, String>,
roots: &[TokenPublicKey],
now: u64,
) -> Result<Vec<AdoptedMember>, JoinError> {
let mut adopted = Vec::with_capacity(assertions.len());
for token in assertions {
let member = roots
.iter()
.find_map(|root| cose::verify_member_assertion(token, root, now).ok())
.ok_or_else(|| {
JoinError::Member(
"a returned member did not verify against the cluster root anchor".into(),
)
})?;
adopted.push(AdoptedMember {
mesh_addr: addrs.get(&member.node_id).cloned(),
node_id: member.node_id,
mesh_pubkey_hex: member.pubkey_hex,
});
}
Ok(adopted)
}
fn verify_token(token: &str, roots: &[TokenPublicKey], now: u64) -> Result<String, JoinError> {
roots
.iter()
.find_map(|root| cose::verify_join(token, root, now).ok())
.ok_or_else(|| JoinError::Token("did not verify against the cluster root anchor".into()))
}
fn possession_proof(
identity: &RpkIdentity,
jti: &str,
mesh_pubkey_hex: &str,
proof_iat: u64,
) -> Result<String, JoinError> {
let challenge = cose::join_challenge(jti, mesh_pubkey_hex, proof_iat);
let sig = identity
.sign(&challenge)
.map_err(|e| JoinError::Proof(e.to_string()))?;
Ok(hex::encode(sig))
}
#[derive(Serialize)]
struct JoinRequestBody {
token: String,
mesh_pubkey: String,
possession_proof: String,
proof_iat: u64,
#[serde(skip_serializing_if = "Option::is_none")]
advertise_addr: Option<String>,
}
#[derive(Deserialize)]
struct JoinResponseBody {
members: Vec<String>,
#[serde(default)]
member_addrs: std::collections::BTreeMap<u64, String>,
}
pub async fn join_cluster(
ticket: &JoinTicket,
identity: &RpkIdentity,
advertise_addr: Option<&str>,
now: u64,
) -> Result<Vec<AdoptedMember>, JoinError> {
let roots = ticket.roots()?;
let jti = verify_token(&ticket.token, &roots, now)?;
let mesh_pubkey_hex = identity.public_key_hex();
let proof_iat = now;
let proof = possession_proof(identity, &jti, &mesh_pubkey_hex, proof_iat)?;
let mut last_err: Option<JoinError> = None;
for seed in &ticket.seeds {
match join_via_seed(
seed,
&roots,
&ticket.token,
&mesh_pubkey_hex,
&proof,
proof_iat,
advertise_addr,
now,
)
.await
{
Ok(body) => return verify_members(&body.members, &body.member_addrs, &roots, now),
Err(err) => last_err = Some(err),
}
}
Err(last_err.unwrap_or_else(|| JoinError::Seed("no seeds in ticket".into())))
}
fn seed_base_url(seed: &str) -> String {
let s = seed.trim();
if s.starts_with("http://") || s.starts_with("https://") {
s.trim_end_matches('/').to_string()
} else {
format!("https://{}", s.trim_end_matches('/'))
}
}
#[allow(clippy::too_many_arguments)]
async fn join_via_seed(
seed: &str,
roots: &[TokenPublicKey],
token: &str,
mesh_pubkey_hex: &str,
proof: &str,
proof_iat: u64,
advertise_addr: Option<&str>,
now: u64,
) -> Result<JoinResponseBody, JoinError> {
let base = seed_base_url(seed);
let http = pinned_client(&base, roots, now).await?;
let resp = http
.post(format!("{base}/api/cluster/join"))
.bearer_auth(token)
.json(&JoinRequestBody {
token: token.to_string(),
mesh_pubkey: mesh_pubkey_hex.to_string(),
possession_proof: proof.to_string(),
proof_iat,
advertise_addr: advertise_addr.map(str::to_string),
})
.send()
.await?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
return Err(JoinError::Refused {
seed: base,
status: status.as_u16(),
body: body.trim().to_string(),
});
}
Ok(resp.json().await?)
}
pub(crate) async fn pinned_client(
base: &str,
roots: &[TokenPublicKey],
now: u64,
) -> Result<reqwest::Client, JoinError> {
let attested_spki = pin_seed(base, roots, now).await?;
let peer: boatramp_rpktls::PeerId = 1;
let trust =
boatramp_rpktls::TrustSet::from_map(std::iter::once((peer, attested_spki)).collect());
let tls = boatramp_rpktls::client_config_server_auth(trust, peer)
.map_err(|e| JoinError::Seed(format!("{base}: {e}")))?;
reqwest::Client::builder()
.use_preconfigured_tls(tls)
.build()
.map_err(JoinError::Http)
}
async fn pin_seed(base: &str, roots: &[TokenPublicKey], now: u64) -> Result<Vec<u8>, JoinError> {
let captured = Arc::new(Mutex::new(None));
let tls = boatramp_rpktls::client_config_capturing(captured.clone())
.map_err(|e| JoinError::Seed(format!("{base}: {e}")))?;
let http = reqwest::Client::builder()
.use_preconfigured_tls(tls)
.build()?;
let attestation = http
.get(format!("{base}/.well-known/boatramp-bootstrap-identity"))
.send()
.await?
.error_for_status()
.map_err(|_| {
JoinError::Seed(format!(
"{base} served no bootstrap attestation (is it `--tls rpk` under this root?)"
))
})?
.text()
.await?;
let attested_hex = roots
.iter()
.find_map(|root| cose::verify_attestation(attestation.trim(), root, now).ok())
.ok_or_else(|| {
JoinError::Seed(format!(
"{base}: attestation did not verify against the cluster root anchor"
))
})?;
let attested = boatramp_rpktls::parse_public_key(&attested_hex)
.map_err(|e| JoinError::Seed(format!("{base}: {e}")))?;
let presented = captured
.lock()
.expect("capture slot")
.clone()
.ok_or_else(|| JoinError::Seed(format!("{base}: presented no key")))?;
if presented != attested {
return Err(JoinError::Seed(format!(
"{base}: the attestation does not match the key the seed presented"
)));
}
Ok(attested)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StartupAction {
Resume,
Found,
Join,
FailClosed(String),
}
#[derive(Debug, Clone, Copy)]
pub struct StartupInputs {
pub has_committed_state: bool,
pub ever_member: bool,
pub seeds_present: bool,
pub init_requested: bool,
}
pub fn decide_startup(i: &StartupInputs) -> StartupAction {
if i.has_committed_state {
return StartupAction::Resume;
}
if i.init_requested && i.seeds_present {
return StartupAction::FailClosed(
"both cluster-init and seeds are configured — a node either founds \
(init, no seeds) or joins (seeds, no init), not both"
.into(),
);
}
if i.seeds_present {
StartupAction::Join
} else if i.init_requested {
StartupAction::Found
} else if i.ever_member {
StartupAction::FailClosed(
"durable store present but holds no committed cluster state, and no seeds \
to rejoin: refusing to self-found (F5) — provide a join ticket/seeds"
.into(),
)
} else {
StartupAction::FailClosed(
"no durable state, no seeds, and no cluster-init: refusing to self-found \
(F5) — pass --cluster-init to found a new cluster or seeds to join one"
.into(),
)
}
}
pub fn resolve_join_token(spec: &str) -> Result<Option<String>, JoinError> {
let spec = spec.trim();
if spec.is_empty() {
return Ok(None);
}
if let Some(var) = spec.strip_prefix("env:") {
return match std::env::var(var) {
Ok(v) if !v.trim().is_empty() => Ok(Some(v.trim().to_string())),
_ => Err(JoinError::Token(format!(
"join_token env var {var} is unset or empty"
))),
};
}
if let Some(path) = spec.strip_prefix("path:") {
let raw = std::fs::read_to_string(path)
.map_err(|e| JoinError::Token(format!("join_token file {path}: {e}")))?;
let token = raw.trim().to_string();
return if token.is_empty() {
Err(JoinError::Token(format!("join_token file {path} is empty")))
} else {
Ok(Some(token))
};
}
Ok(Some(spec.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
use boatramp_core::cose::{LocalSigner, Signer, TokenAlg};
fn now() -> u64 {
1_800_000_000 }
#[test]
fn ticket_round_trips_and_rejects_garbage() {
let ticket = JoinTicket {
seeds: vec!["node-1.internal:7000".into(), "node-2.internal:7000".into()],
root_pubkeys: vec!["es256:0342".into()],
token: "join-token-blob".into(),
};
let blob = ticket.encode().unwrap();
assert!(blob.starts_with("brjoin1."));
assert_eq!(JoinTicket::decode(&blob).unwrap(), ticket);
assert!(JoinTicket::decode("not-a-ticket").is_err());
let empty = JoinTicket {
seeds: vec![],
root_pubkeys: vec![],
token: String::new(),
}
.encode()
.unwrap();
assert!(JoinTicket::decode(&empty).is_err());
}
#[tokio::test]
async fn verify_members_adopts_genuine_and_rejects_forged() {
let root = LocalSigner::generate(TokenAlg::Es256);
let roots = vec![root.public_key()];
let now = now();
let good_a = cose::mint_member_assertion(11, "aa11", 300, now, &root)
.await
.unwrap();
let good_b = cose::mint_member_assertion(22, "bb22", 300, now, &root)
.await
.unwrap();
let addrs = std::collections::BTreeMap::from([(11u64, "https://a:7000".to_string())]);
let adopted =
verify_members(&[good_a.clone(), good_b.clone()], &addrs, &roots, now).unwrap();
assert_eq!(
adopted,
vec![
AdoptedMember {
node_id: 11,
mesh_pubkey_hex: "aa11".into(),
mesh_addr: Some("https://a:7000".into()),
},
AdoptedMember {
node_id: 22,
mesh_pubkey_hex: "bb22".into(),
mesh_addr: None,
},
]
);
let attacker = LocalSigner::generate(TokenAlg::Es256);
let forged = cose::mint_member_assertion(33, "cc33", 300, now, &attacker)
.await
.unwrap();
assert!(verify_members(&[good_a, forged], &addrs, &roots, now).is_err());
}
#[test]
fn possession_proof_binds_the_mesh_key_and_challenge() {
let identity = RpkIdentity::generate().unwrap();
let mesh_pubkey_hex = identity.public_key_hex();
let jti = "single-use-jti";
let proof_iat = now();
let proof_hex = possession_proof(&identity, jti, &mesh_pubkey_hex, proof_iat).unwrap();
let proof = hex::decode(proof_hex).unwrap();
let challenge = cose::join_challenge(jti, &mesh_pubkey_hex, proof_iat);
let spki = boatramp_rpktls::parse_public_key(&mesh_pubkey_hex).unwrap();
assert!(boatramp_rpktls::verify_signature(&spki, &challenge, &proof));
let other =
boatramp_rpktls::parse_public_key(&RpkIdentity::generate().unwrap().public_key_hex())
.unwrap();
assert!(!boatramp_rpktls::verify_signature(
&other, &challenge, &proof
));
}
#[test]
fn startup_decision_enforces_explicit_single_shot_founding() {
let base = StartupInputs {
has_committed_state: false,
ever_member: false,
seeds_present: false,
init_requested: false,
};
let with = |f: fn(&mut StartupInputs)| {
let mut i = base;
f(&mut i);
decide_startup(&i)
};
assert_eq!(
with(|i| {
i.has_committed_state = true;
i.seeds_present = true;
i.init_requested = true;
}),
StartupAction::Resume
);
assert_eq!(with(|i| i.init_requested = true), StartupAction::Found);
assert_eq!(with(|i| i.seeds_present = true), StartupAction::Join);
assert!(matches!(
with(|i| {
i.init_requested = true;
i.seeds_present = true;
}),
StartupAction::FailClosed(_)
));
assert_eq!(
with(|i| {
i.ever_member = true;
i.seeds_present = true;
}),
StartupAction::Join
);
assert_eq!(
with(|i| {
i.ever_member = true;
i.init_requested = true;
}),
StartupAction::Found
);
assert!(matches!(
with(|i| i.ever_member = true),
StartupAction::FailClosed(_)
));
assert!(matches!(with(|_| {}), StartupAction::FailClosed(_)));
}
#[test]
fn join_token_resolves_env_path_and_inline() {
assert_eq!(resolve_join_token("").unwrap(), None);
assert_eq!(
resolve_join_token("inline-token").unwrap(),
Some("inline-token".to_string())
);
std::env::set_var("BOATRAMP_TEST_JOIN_TOKEN", " tok-from-env ");
assert_eq!(
resolve_join_token("env:BOATRAMP_TEST_JOIN_TOKEN").unwrap(),
Some("tok-from-env".to_string())
);
assert!(resolve_join_token("env:BOATRAMP_TEST_JOIN_TOKEN_UNSET").is_err());
let dir = std::env::temp_dir();
let file = dir.join("boatramp-test-join-token");
std::fs::write(&file, "tok-from-file\n").unwrap();
assert_eq!(
resolve_join_token(&format!("path:{}", file.display())).unwrap(),
Some("tok-from-file".to_string())
);
std::fs::write(&file, " \n").unwrap();
assert!(resolve_join_token(&format!("path:{}", file.display())).is_err());
let _ = std::fs::remove_file(&file);
}
}