use std::{
sync::{Arc, LazyLock},
time::{Duration, SystemTime},
};
use chrono::{DateTime, SecondsFormat, Utc};
use futures::StreamExt;
use k8s_openapi::api::core::v1::Secret;
use kube::{
Api, Client, Resource, ResourceExt,
api::{Patch, PatchParams},
runtime::{
controller::{Action, Controller},
watcher,
},
};
use polyc_agent::ToolExecutor;
use polyc_tools::mcp_client::McpToolSource;
use prometheus::{IntCounter, register_int_counter};
use serde_json::json;
use crate::reserved_secrets::{ReservedSecret, check_secret_ref};
use crate::toolservice::{Auth, BearerSecretRef, FINALIZER, ToolDescriptor, ToolService};
pub const CHECK_INTERVAL: Duration = Duration::from_mins(5);
static TOOLSERVICE_RECONCILE_TOTAL: LazyLock<IntCounter> = LazyLock::new(|| {
register_int_counter!(
"polychrome_toolservice_reconcile_total",
"Total ToolService reconcile passes"
)
.expect("register polychrome_toolservice_reconcile_total")
});
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("kube api: {0}")]
Kube(#[from] kube::Error),
#[error("toolservice has no namespace")]
NoNamespace,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ToolServiceAction {
AddFinalizer,
Cleanup,
CheckHealth {
url: String,
},
Noop,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ToolServiceReadiness {
pub healthy: bool,
pub available_tools: Vec<ToolDescriptor>,
pub message: Option<String>,
}
#[must_use]
pub fn plan(ts: &ToolService, now: SystemTime, interval: Duration) -> ToolServiceAction {
let being_deleted = ts.meta().deletion_timestamp.is_some();
let has_finalizer = ts.finalizers().iter().any(|f| f == FINALIZER);
if being_deleted {
return if has_finalizer {
ToolServiceAction::Cleanup
} else {
ToolServiceAction::Noop
};
}
if !has_finalizer {
return ToolServiceAction::AddFinalizer;
}
if check_due(ts, now, interval) {
let url = ts
.spec
.remotes
.first()
.map(|r| r.url.clone())
.unwrap_or_default();
return ToolServiceAction::CheckHealth { url };
}
ToolServiceAction::Noop
}
fn check_due(ts: &ToolService, now: SystemTime, interval: Duration) -> bool {
let Some(status) = ts.status.as_ref() else {
return true; };
if status.observed_version.as_deref() != Some(ts.spec.version.as_str()) {
return true; }
let Some(last) = status.last_checked_at.as_deref() else {
return true; };
let Ok(parsed) = DateTime::parse_from_rfc3339(last) else {
return true; };
let now: DateTime<Utc> = now.into();
now.signed_duration_since(parsed.with_timezone(&Utc))
.to_std()
.is_ok_and(|elapsed| elapsed >= interval)
}
pub async fn health_check(url: &str, bearer: Option<String>) -> ToolServiceReadiness {
if url.is_empty() {
return ToolServiceReadiness {
healthy: false,
available_tools: Vec::new(),
message: Some("no remote URL configured".to_owned()),
};
}
let bearer = match bearer {
Some(token) => match polyc_tools::AudienceBoundToken::new(token, url) {
Ok(token) => Some(token),
Err(err) => {
return ToolServiceReadiness {
healthy: false,
available_tools: Vec::new(),
message: Some(err.to_string()),
};
}
},
None => None,
};
match McpToolSource::connect(
url.to_owned(),
polyc_tools::ConnectOptions {
bearer,
..polyc_tools::ConnectOptions::default()
},
)
.await
{
Ok(source) => {
let available_tools = source
.specs()
.into_iter()
.map(descriptor_from_spec)
.collect();
source.shutdown();
ToolServiceReadiness {
healthy: true,
available_tools,
message: None,
}
}
Err(err) => ToolServiceReadiness {
healthy: false,
available_tools: Vec::new(),
message: Some(err.to_string()),
},
}
}
fn descriptor_from_spec(spec: polyc_llm::ToolSpec) -> ToolDescriptor {
ToolDescriptor {
name: spec.name,
description: (!spec.description.is_empty()).then_some(spec.description),
input_schema: serde_json::to_string(&spec.schema_json).unwrap_or_default(),
title: spec.title,
read_only: spec.read_only,
destructive: spec.destructive,
open_world: spec.open_world,
}
}
#[must_use]
fn version_keyed_catalog(
current: &crate::toolservice::ToolServiceStatus,
spec_version: &str,
readiness: &ToolServiceReadiness,
) -> Vec<ToolDescriptor> {
let version_matches = current.observed_version.as_deref() == Some(spec_version);
if version_matches || !readiness.healthy {
current.available_tools.clone()
} else {
readiness.available_tools.clone()
}
}
#[must_use]
fn observed_version_after_check(
current: &crate::toolservice::ToolServiceStatus,
spec_version: &str,
healthy: bool,
) -> Option<String> {
let version_matches = current.observed_version.as_deref() == Some(spec_version);
if version_matches || healthy {
Some(spec_version.to_owned())
} else {
current.observed_version.clone()
}
}
pub struct Context {
pub client: Client,
pub check_interval: Duration,
}
fn admissible_bearer_ref(auth: Option<&Auth>) -> Result<Option<&BearerSecretRef>, ReservedSecret> {
let Some(auth) = auth else {
return Ok(None);
};
let r = &auth.bearer_secret_ref;
check_secret_ref(&r.name)?;
Ok(Some(r))
}
async fn resolve_bearer(
client: &Client,
namespace: &str,
auth: Option<&Auth>,
) -> Result<Option<String>, ReservedSecret> {
let Some(r) = admissible_bearer_ref(auth)? else {
return Ok(None);
};
let secrets: Api<Secret> = Api::namespaced(client.clone(), namespace);
match secrets.get(&r.name).await {
Ok(secret) => Ok(secret
.data
.as_ref()
.and_then(|d| d.get(&r.key))
.and_then(|b| String::from_utf8(b.0.clone()).ok())
.map(|s| s.trim().to_owned())),
Err(e) => {
tracing::warn!(secret = %r.name, key = %r.key, error = %e,
"toolservice health check: failed to read bearer secret");
Ok(None)
}
}
}
fn refused_readiness(err: &ReservedSecret) -> ToolServiceReadiness {
ToolServiceReadiness {
healthy: false,
available_tools: Vec::new(),
message: Some(err.to_string()),
}
}
#[tracing::instrument(skip_all, fields(toolservice = %ts.name_any()))]
pub async fn reconcile(ts: Arc<ToolService>, ctx: Arc<Context>) -> Result<Action, Error> {
TOOLSERVICE_RECONCILE_TOTAL.inc();
let ns = ts.namespace().ok_or(Error::NoNamespace)?;
let name = ts.name_any();
let api: Api<ToolService> = Api::namespaced(ctx.client.clone(), &ns);
let pp = PatchParams::apply("polychrome.dev/controller");
match plan(&ts, SystemTime::now(), ctx.check_interval) {
ToolServiceAction::AddFinalizer => {
let patch = json!({ "metadata": { "finalizers": [FINALIZER] } });
api.patch(&name, &pp, &Patch::Merge(&patch)).await?;
Ok(Action::requeue(Duration::from_secs(1)))
}
ToolServiceAction::Cleanup => {
let patch = json!({ "metadata": { "finalizers": [] } });
api.patch(&name, &pp, &Patch::Merge(&patch)).await?;
Ok(Action::await_change())
}
ToolServiceAction::CheckHealth { url } => {
let readiness = match resolve_bearer(&ctx.client, &ns, ts.spec.auth.as_ref()).await {
Ok(bearer) => health_check(&url, bearer).await,
Err(err) => {
tracing::error!(
toolservice = %name,
secret = %err.secret,
"refused a toolservice bearer reference to platform key material"
);
refused_readiness(&err)
}
};
let version = ts.spec.version.clone();
let now = Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true);
let current = ts.status.clone().unwrap_or_default();
let catalog = version_keyed_catalog(¤t, &version, &readiness);
let observed_version =
observed_version_after_check(¤t, &version, readiness.healthy);
let changed = current.healthy != readiness.healthy
|| current.available_tools != catalog
|| current.message != readiness.message
|| current.observed_version != observed_version;
if changed {
let patch = json!({ "status": {
"healthy": readiness.healthy,
"lastCheckedAt": now,
"availableTools": catalog,
"message": readiness.message,
"observedVersion": observed_version,
} });
api.patch_status(&name, &PatchParams::default(), &Patch::Merge(&patch))
.await?;
tracing::info!(
healthy = readiness.healthy,
tools = catalog.len(),
message = ?readiness.message,
"synced toolservice status from MCP health check"
);
} else {
let patch = json!({ "status": { "lastCheckedAt": now } });
api.patch_status(&name, &PatchParams::default(), &Patch::Merge(&patch))
.await?;
}
Ok(Action::requeue(ctx.check_interval))
}
ToolServiceAction::Noop => Ok(Action::requeue(ctx.check_interval)),
}
}
#[must_use]
pub fn error_policy(_ts: Arc<ToolService>, err: &Error, _ctx: Arc<Context>) -> Action {
tracing::warn!(error = %err, "toolservice reconcile failed; requeuing");
Action::requeue(Duration::from_secs(10))
}
pub async fn run_toolservice(
client: Client,
watch_client: Client,
namespace: &str,
) -> Result<(), Error> {
let api: Api<ToolService> = Api::namespaced(watch_client, namespace);
let ctx = Arc::new(Context {
client,
check_interval: CHECK_INTERVAL,
});
Controller::new(api, watcher::Config::default())
.run(reconcile, error_policy, ctx)
.for_each(|res| async move {
if let Err(e) = res {
tracing::warn!(error = %e, "toolservice reconcile stream item errored");
}
})
.await;
Ok(())
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use std::time::{Duration, SystemTime};
use chrono::{SecondsFormat, Utc};
use kube::api::ObjectMeta;
use super::*;
use crate::toolservice::{Remote, ToolServiceSpec, ToolServiceStatus};
const INTERVAL: Duration = Duration::from_secs(300);
fn ts(name: &str) -> ToolService {
ToolService::new(
name,
ToolServiceSpec {
server_name: Some("io.github.acme/search".to_owned()),
description: Some("Search connector".to_owned()),
version: "1.0.0".to_owned(),
remotes: vec![Remote {
url: "https://acme.example/mcp".to_owned(),
transport: "streamable-http".to_owned(),
}],
needs_approval: false,
approval_tools: Vec::new(),
default_enabled: false,
auth: None,
tool_expansions: std::collections::BTreeMap::new(),
},
)
}
fn finalized(name: &str) -> ToolService {
let mut t = ts(name);
t.metadata.finalizers = Some(vec![FINALIZER.to_owned()]);
t
}
fn rfc3339(at: SystemTime) -> String {
let dt: chrono::DateTime<Utc> = at.into();
dt.to_rfc3339_opts(SecondsFormat::Secs, true)
}
#[test]
fn fresh_toolservice_gets_a_finalizer_first() {
assert_eq!(
plan(&ts("s1"), SystemTime::now(), INTERVAL),
ToolServiceAction::AddFinalizer
);
}
#[test]
fn finalized_without_status_checks_health() {
assert_eq!(
plan(&finalized("s1"), SystemTime::now(), INTERVAL),
ToolServiceAction::CheckHealth {
url: "https://acme.example/mcp".to_owned(),
}
);
}
#[test]
fn fresh_status_settles_to_noop() {
let mut t = finalized("s1");
let now = SystemTime::now();
t.status = Some(ToolServiceStatus {
healthy: true,
last_checked_at: Some(rfc3339(now)),
observed_version: Some("1.0.0".to_owned()),
..Default::default()
});
assert_eq!(plan(&t, now, INTERVAL), ToolServiceAction::Noop);
}
#[test]
fn stale_status_checks_health() {
let mut t = finalized("s1");
let checked = SystemTime::now() - Duration::from_secs(600);
t.status = Some(ToolServiceStatus {
healthy: true,
last_checked_at: Some(rfc3339(checked)),
observed_version: Some("1.0.0".to_owned()),
..Default::default()
});
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::CheckHealth {
url: "https://acme.example/mcp".to_owned(),
}
);
}
#[test]
fn version_bump_checks_health_even_when_timestamp_is_fresh() {
let mut t = finalized("s1");
let now = SystemTime::now();
t.status = Some(ToolServiceStatus {
healthy: true,
last_checked_at: Some(rfc3339(now)),
observed_version: Some("0.9.0".to_owned()), ..Default::default()
});
assert_eq!(
plan(&t, now, INTERVAL),
ToolServiceAction::CheckHealth {
url: "https://acme.example/mcp".to_owned(),
}
);
}
#[test]
fn missing_timestamp_checks_health() {
let mut t = finalized("s1");
t.status = Some(ToolServiceStatus {
healthy: true,
last_checked_at: None,
observed_version: Some("1.0.0".to_owned()),
..Default::default()
});
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::CheckHealth {
url: "https://acme.example/mcp".to_owned(),
}
);
}
#[test]
fn unparseable_timestamp_checks_health() {
let mut t = finalized("s1");
t.status = Some(ToolServiceStatus {
healthy: true,
last_checked_at: Some("not-a-timestamp".to_owned()),
observed_version: Some("1.0.0".to_owned()),
..Default::default()
});
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::CheckHealth {
url: "https://acme.example/mcp".to_owned(),
}
);
}
#[test]
fn check_due_with_no_remotes_yields_empty_url() {
let mut t = finalized("s1");
t.spec.remotes.clear();
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::CheckHealth { url: String::new() }
);
}
#[test]
fn deletion_with_finalizer_cleans_up() {
let mut t = finalized("s1");
t.metadata.deletion_timestamp = Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
"2026-05-27T00:00:00Z".parse().unwrap(),
));
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::Cleanup
);
}
#[test]
fn deletion_without_finalizer_is_noop() {
let mut t = ts("s1");
t.metadata = ObjectMeta {
name: Some("s1".to_owned()),
deletion_timestamp: Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
"2026-05-27T00:00:00Z".parse().unwrap(),
)),
..Default::default()
};
assert_eq!(
plan(&t, SystemTime::now(), INTERVAL),
ToolServiceAction::Noop
);
}
fn descriptor(name: &str) -> ToolDescriptor {
ToolDescriptor {
name: name.to_owned(),
description: Some("d".to_owned()),
input_schema: r#"{"type":"object"}"#.to_owned(),
title: None,
read_only: false,
destructive: false,
open_world: true,
}
}
#[test]
fn descriptor_from_spec_carries_schema_and_annotations() {
let mut spec = polyc_llm::ToolSpec::new(
"delete_file",
"Delete a file.",
serde_json::json!({ "type": "object", "properties": { "path": { "type": "string" } } }),
);
spec.title = Some("Delete a file".to_owned());
spec.destructive = true;
spec.open_world = true;
let d = descriptor_from_spec(spec);
assert_eq!(d.name, "delete_file");
assert_eq!(d.description.as_deref(), Some("Delete a file."));
assert_eq!(d.title.as_deref(), Some("Delete a file"));
assert!(d.destructive);
assert!(d.open_world);
assert!(!d.read_only);
let schema: serde_json::Value = serde_json::from_str(&d.input_schema).unwrap();
assert_eq!(schema["properties"]["path"]["type"], "string");
}
fn healthy_readiness(tools: Vec<ToolDescriptor>) -> ToolServiceReadiness {
ToolServiceReadiness {
healthy: true,
available_tools: tools,
message: None,
}
}
fn failed_readiness() -> ToolServiceReadiness {
ToolServiceReadiness {
healthy: false,
available_tools: Vec::new(),
message: Some("dial failed".to_owned()),
}
}
#[test]
fn version_keyed_catalog_freezes_on_unchanged_version() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
available_tools: vec![descriptor("echo")],
..Default::default()
};
let drifted = healthy_readiness(vec![descriptor("echo"), descriptor("sneaky_new_tool")]);
let kept = version_keyed_catalog(¤t, "1.0.0", &drifted);
assert_eq!(
kept,
vec![descriptor("echo")],
"catalog frozen to the reviewed version"
);
}
#[test]
fn version_keyed_catalog_adopts_on_version_bump() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
available_tools: vec![descriptor("echo")],
..Default::default()
};
let fresh = vec![descriptor("echo"), descriptor("new_tool")];
let adopted = version_keyed_catalog(¤t, "2.0.0", &healthy_readiness(fresh.clone()));
assert_eq!(
adopted, fresh,
"a version bump adopts the freshly listed catalog"
);
}
#[test]
fn version_keyed_catalog_adopts_on_first_check() {
let current = ToolServiceStatus::default();
let fresh = vec![descriptor("echo")];
assert_eq!(
version_keyed_catalog(¤t, "1.0.0", &healthy_readiness(fresh.clone())),
fresh
);
}
#[test]
fn version_keyed_catalog_keeps_catalog_on_failed_same_version_recheck() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
available_tools: vec![descriptor("echo")],
healthy: true,
..Default::default()
};
assert_eq!(
version_keyed_catalog(¤t, "1.0.0", &failed_readiness()),
vec![descriptor("echo")],
"an outage must not remove tools from the catalog"
);
}
#[test]
fn version_keyed_catalog_freezes_prior_catalog_on_failed_check_during_version_bump() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
available_tools: vec![descriptor("echo")],
healthy: true,
..Default::default()
};
let kept = version_keyed_catalog(¤t, "2.0.0", &failed_readiness());
assert_eq!(
kept,
vec![descriptor("echo")],
"a failed check on a version bump must not wipe the prior catalog"
);
let recovered = healthy_readiness(vec![descriptor("echo"), descriptor("new_tool")]);
let adopted = version_keyed_catalog(¤t, "2.0.0", &recovered);
assert_eq!(adopted, recovered.available_tools);
}
#[test]
fn version_keyed_catalog_stays_empty_on_failed_first_check() {
let current = ToolServiceStatus::default();
assert_eq!(
version_keyed_catalog(¤t, "1.0.0", &failed_readiness()),
Vec::<ToolDescriptor>::new()
);
}
#[test]
fn observed_version_settles_on_healthy_check() {
let current = ToolServiceStatus::default();
assert_eq!(
observed_version_after_check(¤t, "1.0.0", true),
Some("1.0.0".to_owned())
);
}
#[test]
fn observed_version_does_not_settle_on_failed_version_bump() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
..Default::default()
};
assert_eq!(
observed_version_after_check(¤t, "2.0.0", false),
Some("1.0.0".to_owned()),
"keep the prior observed_version so check_due keeps forcing retries"
);
}
#[test]
fn observed_version_stays_unset_on_failed_first_check() {
let current = ToolServiceStatus::default();
assert_eq!(observed_version_after_check(¤t, "1.0.0", false), None);
}
#[test]
fn observed_version_settles_on_same_version_even_if_unhealthy() {
let current = ToolServiceStatus {
observed_version: Some("1.0.0".to_owned()),
..Default::default()
};
assert_eq!(
observed_version_after_check(¤t, "1.0.0", false),
Some("1.0.0".to_owned())
);
}
#[tokio::test]
async fn health_check_empty_url_is_unhealthy_not_an_error() {
let r = health_check("", None).await;
assert!(!r.healthy);
assert!(r.available_tools.is_empty());
assert_eq!(r.message.as_deref(), Some("no remote URL configured"));
}
fn auth_for(secret: &str) -> Auth {
Auth {
bearer_secret_ref: BearerSecretRef {
name: secret.to_owned(),
key: "token".to_owned(),
},
}
}
#[test]
fn the_bearer_ref_screen_refuses_every_reserved_secret() {
let wallet = format!(
"{}{}",
crate::reserved_secrets::CONTROL_KEY_SECRET_PREFIX,
"ab".repeat(32)
);
let reserved = crate::reserved_secrets::RESERVED_SECRET_NAMES
.iter()
.map(|s| (*s).to_owned())
.chain(std::iter::once(wallet));
for secret in reserved {
let auth = auth_for(&secret);
let err = admissible_bearer_ref(Some(&auth))
.expect_err("a reserved reference must be refused, not admitted");
assert_eq!(err.secret, secret);
assert!(
err.to_string().contains(&secret),
"the refusal must name the Secret it refused"
);
let readiness = refused_readiness(&err);
assert!(!readiness.healthy);
assert!(readiness.available_tools.is_empty());
assert_eq!(readiness.message.as_deref(), Some(err.to_string().as_str()));
}
}
#[test]
fn the_bearer_ref_screen_admits_an_author_chosen_secret() {
let auth = auth_for("acme-mcp-token");
let admitted = admissible_bearer_ref(Some(&auth))
.expect("an ordinary reference is never refused")
.expect("configured auth yields a reference to read");
assert_eq!(admitted.name, "acme-mcp-token");
assert_eq!(admitted.key, "token");
assert!(
admissible_bearer_ref(None)
.expect("no auth is never a refusal")
.is_none()
);
}
#[test]
fn a_refused_toolservice_settles_unhealthy_with_its_catalog_frozen() {
let auth = auth_for(crate::reserved_secrets::STATE_ATTESTATION_SECRET);
let err = admissible_bearer_ref(Some(&auth)).expect_err("refused");
let readiness = refused_readiness(&err);
let current = ToolServiceStatus {
healthy: true,
observed_version: Some("1.0.0".to_owned()),
available_tools: vec![ToolDescriptor {
name: "find".to_owned(),
description: None,
input_schema: String::new(),
title: None,
read_only: true,
destructive: false,
open_world: false,
}],
..Default::default()
};
let catalog = version_keyed_catalog(¤t, "2.0.0", &readiness);
assert_eq!(catalog, current.available_tools);
assert_eq!(
observed_version_after_check(¤t, "2.0.0", readiness.healthy),
Some("1.0.0".to_owned())
);
}
#[tokio::test]
async fn health_check_unreachable_url_is_unhealthy_not_an_error() {
let r = health_check("http://127.0.0.1:1/mcp", None).await;
assert!(!r.healthy);
assert!(r.available_tools.is_empty());
assert!(r.message.is_some());
}
}