use anyhow::{anyhow, Context, Result};
use clap::Parser;
use k8s_openapi::api::core::v1::ConfigMap;
use kube::api::{Patch, PatchParams, PostParams};
use kube::{Api, Client};
use serde_json::json;
use std::collections::BTreeMap;
use tatara_process::receipt::{ReceiptEnvelope, ReceiptKind};
use tracing::{info, warn};
mod probe;
#[derive(Parser, Debug)]
#[command(name = "closed-loop-probe")]
#[command(about = "Closed-loop authentication probe — emits a tatara-receipt/v1 envelope")]
struct Args {
#[arg(long, env = "ISSUER_SERVICE")]
issuer_service: String,
#[arg(long, env = "ISSUER_PORT", default_value_t = 8080)]
issuer_port: u16,
#[arg(long, env = "ISSUER_AUTH_PATH", default_value = "/v2/auth")]
issuer_auth_path: String,
#[arg(long, env = "ISSUER_JWKS_PATH", default_value = "/.well-known/jwks.json")]
issuer_jwks_path: String,
#[arg(long, env = "CONSUMER_SERVICE")]
consumer_service: String,
#[arg(long, env = "CONSUMER_PORT", default_value_t = 8000)]
consumer_port: u16,
#[arg(long, env = "CONSUMER_AUTH_PATH", default_value = "/v2/whoami")]
consumer_auth_path: String,
#[arg(long, env = "RECEIPT_CONFIG_MAP")]
receipt_config_map: String,
#[arg(long, env = "RECEIPT_NAMESPACE", default_value = "default")]
receipt_namespace: String,
#[arg(long, env = "RECEIPT_KIND", default_value_t = String::from(ReceiptKind::ClosedLoopAuth))]
receipt_kind: String,
#[arg(long, env = "TATARA_PROCESS_REF")]
process_ref: Option<String>,
#[arg(long, default_value = "10s")]
timeout: humantime::Duration,
}
#[tokio::main]
async fn main() -> Result<()> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let args = Args::parse();
let access_id =
std::env::var("ACCESS_ID").context("ACCESS_ID env var (from auth Secret) required")?;
let access_key =
std::env::var("ACCESS_KEY").context("ACCESS_KEY env var (from auth Secret) required")?;
info!(
issuer = %args.issuer_service,
consumer = %args.consumer_service,
receipt_cm = %args.receipt_config_map,
"starting closed-loop probe"
);
let probe_result = probe::run(
probe::ProbeConfig {
issuer: probe::ServiceEndpoint {
service: args.issuer_service,
port: args.issuer_port,
},
issuer_auth_path: args.issuer_auth_path,
issuer_jwks_path: args.issuer_jwks_path,
consumer: probe::ServiceEndpoint {
service: args.consumer_service,
port: args.consumer_port,
},
consumer_auth_path: args.consumer_auth_path,
access_id,
access_key,
http_timeout: args.timeout.into(),
},
)
.await?;
let mut envelope = ReceiptEnvelope::build(
&args.receipt_kind,
&probe_result.intent_hash,
&probe_result.artifact_hash,
&probe_result.control_hash,
None,
);
envelope.process_ref = args.process_ref.clone();
envelope.evidence = json!({
"issuer_url": probe_result.issuer_url,
"consumer_url": probe_result.consumer_url,
"token_present": probe_result.token_present,
"jwks_keys": probe_result.jwks_key_count,
"whoami_status": probe_result.whoami_status,
});
info!(
composed_root = %envelope.composed_root,
kind = %envelope.kind,
"writing receipt to ConfigMap"
);
write_receipt(&envelope, &args.receipt_config_map, &args.receipt_namespace).await?;
info!("closed-loop probe succeeded");
Ok(())
}
async fn write_receipt(envelope: &ReceiptEnvelope, cm_name: &str, ns: &str) -> Result<()> {
let client = Client::try_default()
.await
.context("create in-cluster kube client")?;
let api: Api<ConfigMap> = Api::namespaced(client, ns);
let payload = serde_json::to_string(envelope)?;
let mut data = BTreeMap::new();
data.insert("receipt.json".to_string(), payload.clone());
data.insert("receipt.yaml".to_string(), serde_yaml::to_string(envelope)?);
let cm = ConfigMap {
metadata: kube::core::ObjectMeta {
name: Some(cm_name.into()),
namespace: Some(ns.into()),
labels: Some(BTreeMap::from([(
"tatara.pleme.io/receipt".into(),
"tatara-receipt/v1".into(),
)])),
..Default::default()
},
data: Some(data),
..Default::default()
};
match api.create(&PostParams::default(), &cm).await {
Ok(_) => Ok(()),
Err(kube::Error::Api(e)) if e.code == 409 => {
let patch = json!({ "data": cm.data });
api.patch(cm_name, &PatchParams::default(), &Patch::Merge(&patch))
.await
.map_err(|e| anyhow!("patch ConfigMap {ns}/{cm_name}: {e}"))?;
Ok(())
}
Err(e) => {
warn!(error = %e, "create ConfigMap failed");
Err(anyhow!("create ConfigMap {ns}/{cm_name}: {e}"))
}
}
}
#[cfg(test)]
mod tests {
use super::Args;
use clap::Parser;
use tatara_process::receipt::ReceiptKind;
#[test]
fn args_parse_with_required_flags() {
let args = Args::try_parse_from([
"closed-loop-probe",
"--issuer-service",
"gator",
"--consumer-service",
"gateway",
"--receipt-config-map",
"my-receipt",
]);
assert!(args.is_ok(), "{:?}", args.err());
let a = args.unwrap();
assert_eq!(a.issuer_port, 8080);
assert_eq!(a.receipt_kind, ReceiptKind::ClosedLoopAuth.as_str());
}
}