use boatramp_core::cose::TokenPublicKey;
use k8s_openapi::api::core::v1::Secret;
use kube::api::{Api, Patch, PatchParams};
use kube::Client;
use serde_json::json;
use super::crd::BoatRampCluster;
use super::membership::{self, ApiMember, MembershipAction};
use super::{resources, Error, Result};
const CONTROL_PLANE_PORT: u16 = 8080;
const TOKEN_KEY: &str = "token";
fn join_secret_name(brc: &BoatRampCluster) -> String {
format!("{}-join", resources::instance(brc))
}
const JOIN_KEY: &str = "ticket";
use boatramp_core::time::now_unix;
fn pod_base(brc: &BoatRampCluster, ns: &str, ordinal: u32) -> String {
let inst = resources::instance(brc);
format!("https://{inst}-{ordinal}.{inst}-headless.{ns}.svc:{CONTROL_PLANE_PORT}")
}
struct ClusterApi {
http: reqwest::Client,
base: String,
token: String,
}
pub(super) async fn roll_partition(
client: &Client,
ns: &str,
brc: &BoatRampCluster,
pods: &[membership::PodState],
) -> i32 {
let replicas = brc.spec.replicas as i32;
let Ok(Some((token, roots))) = admin_creds(client, ns, brc).await else {
return 0;
};
let Ok(api) = ClusterApi::pin_pod(brc, ns, 0, &token, &roots).await else {
return 0;
};
let Ok(raw) = api.members().await else {
return 0;
};
let (members, _) = membership::members_from_api(&raw);
if membership::has_roll_margin(&members, pods) {
0
} else {
replicas
}
}
pub(super) async fn admin_creds(
client: &Client,
ns: &str,
brc: &BoatRampCluster,
) -> Result<Option<(String, Vec<TokenPublicKey>)>> {
let (Some(secret_name), Some(root)) = (
brc.spec.admin_token_secret.as_deref(),
brc.spec.root_pubkey.as_deref(),
) else {
return Ok(None);
};
let root = TokenPublicKey::from_hex(root.trim())
.map_err(|e| Error::Other(format!("invalid spec.rootPubkey: {e}")))?;
let secrets: Api<Secret> = Api::namespaced(client.clone(), ns);
let secret = secrets.get(secret_name).await?;
let token = secret
.data
.as_ref()
.and_then(|d| d.get(TOKEN_KEY))
.map(|b| String::from_utf8_lossy(&b.0).trim().to_string())
.filter(|t| !t.is_empty())
.ok_or_else(|| {
Error::Other(format!(
"admin token Secret {secret_name:?} has no non-empty `{TOKEN_KEY}` key"
))
})?;
Ok(Some((token, vec![root])))
}
pub(super) async fn pinned_admin_pod0(
client: &Client,
ns: &str,
brc: &BoatRampCluster,
) -> Result<Option<(reqwest::Client, String, String)>> {
let Some((token, roots)) = admin_creds(client, ns, brc).await? else {
return Ok(None);
};
let api = ClusterApi::pin_pod(brc, ns, 0, &token, &roots).await?;
Ok(Some((api.http, api.base, api.token)))
}
impl ClusterApi {
async fn pin_pod(
brc: &BoatRampCluster,
ns: &str,
ordinal: u32,
token: &str,
roots: &[TokenPublicKey],
) -> Result<Self> {
let base = pod_base(brc, ns, ordinal);
let http = crate::join::pinned_client(&base, roots, now_unix())
.await
.map_err(|e| Error::Other(format!("pin {base}: {e}")))?;
Ok(Self {
http,
base,
token: token.to_string(),
})
}
async fn members(&self) -> Result<Vec<ApiMember>> {
#[derive(serde::Deserialize)]
struct Row {
node: u64,
#[serde(default)]
voter: bool,
#[serde(default)]
caught_up: bool,
#[serde(default)]
leader: bool,
#[serde(default)]
addr: Option<String>,
}
let rows: Vec<Row> = self
.http
.get(format!("{}/api/cluster/members", self.base))
.bearer_auth(&self.token)
.send()
.await?
.error_for_status()?
.json()
.await?;
Ok(rows
.into_iter()
.map(|r| ApiMember {
node_id: r.node,
voter: r.voter,
caught_up: r.caught_up,
leader: r.leader,
addr: r.addr,
})
.collect())
}
async fn promote(&self, node: u64) -> Result<()> {
self.http
.post(format!("{}/api/cluster/promote", self.base))
.bearer_auth(&self.token)
.json(&json!({ "node_id": node }))
.send()
.await?
.error_for_status()?;
Ok(())
}
async fn remove(&self, node: u64) -> Result<()> {
self.http
.post(format!("{}/api/cluster/revoke", self.base))
.bearer_auth(&self.token)
.json(&json!({ "node_id": node }))
.send()
.await?
.error_for_status()?;
Ok(())
}
async fn mint_join_token(&self) -> Result<String> {
#[derive(serde::Deserialize)]
struct Resp {
token: String,
}
let resp: Resp = self
.http
.post(format!("{}/api/cluster/join-token", self.base))
.bearer_auth(&self.token)
.json(&json!({}))
.send()
.await?
.error_for_status()?
.json()
.await?;
Ok(resp.token)
}
}
pub async fn step(
client: &Client,
ns: &str,
brc: &BoatRampCluster,
pods: &[membership::PodState],
root_pubkey: &str,
) -> Result<Option<MembershipAction>> {
let Some((token, roots)) = admin_creds(client, ns, brc).await? else {
return Ok(None);
};
let api0 = ClusterApi::pin_pod(brc, ns, 0, &token, &roots).await?;
let raw = api0.members().await?;
let (members, ordinal_to_node) = membership::members_from_api(&raw);
let leader_ordinal = raw
.iter()
.find(|m| m.leader)
.and_then(|m| m.addr.as_deref())
.and_then(membership::ordinal_from_addr)
.unwrap_or(0);
if (members.len() as u32) < brc.spec.replicas {
let ticket_token = api0.mint_join_token().await?;
let seed = pod_base(brc, ns, leader_ordinal);
let ticket = crate::join::JoinTicket {
seeds: vec![seed],
root_pubkeys: vec![root_pubkey.to_string()],
token: ticket_token,
}
.encode()
.map_err(|e| Error::Other(e.to_string()))?;
store_join_ticket(client, ns, brc, &ticket).await?;
}
let Some(action) = membership::plan_next(brc.spec.replicas, pods, &members) else {
return Ok(None);
};
match action {
MembershipAction::PromoteToVoter { .. } | MembershipAction::Remove { .. } => {
let Some(node) = membership::action_node_id(&action, &ordinal_to_node) else {
return Ok(Some(action));
};
let leader = if leader_ordinal == 0 {
api0
} else {
ClusterApi::pin_pod(brc, ns, leader_ordinal, &token, &roots).await?
};
match action {
MembershipAction::PromoteToVoter { .. } => leader.promote(node).await?,
MembershipAction::Remove { .. } => leader.remove(node).await?,
MembershipAction::AddLearner { .. } => unreachable!(),
}
}
MembershipAction::AddLearner { .. } => {}
}
Ok(Some(action))
}
async fn store_join_ticket(
client: &Client,
ns: &str,
brc: &BoatRampCluster,
ticket: &str,
) -> Result<()> {
let name = join_secret_name(brc);
let secret = json!({
"apiVersion": "v1",
"kind": "Secret",
"metadata": { "name": name, "namespace": ns },
"type": "Opaque",
"stringData": { JOIN_KEY: ticket },
});
let api: Api<Secret> = Api::namespaced(client.clone(), ns);
api.patch(
&name,
&PatchParams::apply("boatramp-operator").force(),
&Patch::Apply(&secret),
)
.await?;
Ok(())
}
pub fn join_env_source(brc: &BoatRampCluster) -> (String, &'static str) {
(join_secret_name(brc), JOIN_KEY)
}