polyc-controller 2026.10.0

Conversation CRD + kube reconciler for the polychrome control plane.
#![allow(clippy::unwrap_used)] // test/example/bench: panics are acceptable
#![allow(clippy::too_many_lines)] // test/example/bench: long linear setup is fine
//! Live-cluster end-to-end test for the reconciler.
//!
//! Gated behind the `e2e` feature. It is clippy-checked in CI (the lint step
//! runs `--all-features`) but never executed there — running it needs a live
//! cluster for the `Conversation` objects. The substrate control API is an
//! in-memory fake, so no substrate install is needed. Run locally against the
//! current kube context (e.g. OrbStack):
//!
//! ```text
//! # one-time: apply the CRDs
//! cargo run -p polyc-controller --bin crdgen | kubectl apply -f -
//! # then:
//! cargo test -p polyc-controller --features e2e -- --nocapture
//! ```
//!
//! It drives [`reconcile`] directly (rather than the full watch loop) so the
//! sequence is deterministic: create a `Conversation` → reconcile adds the
//! finalizer → reconcile creates the actor → assert → suspend when idle →
//! recreate after a crash → delete and reconcile once more to run cleanup.
#![cfg(feature = "e2e")]

use std::{
    collections::BTreeMap,
    sync::{
        Arc, Mutex,
        atomic::{AtomicU32, Ordering},
    },
    time::{SystemTime, UNIX_EPOCH},
};

use k8s_openapi::api::core::v1::Namespace;
use kube::{
    Api, Client, ResourceExt,
    api::{DeleteParams, ListParams, ObjectMeta, Patch, PatchParams, PostParams},
};
use polyc_controller::{
    ActorControl, ActorControlError, ActorSnapshot, ActorState, Conversation, ConversationSpec,
    SubstrateActorBackend,
    reconcile::{Context, reconcile},
};
use serde_json::json;

const NS: &str = "polychrome-e2e";
const TEMPLATE: &str = "polychrome-harness-e2e";

/// An in-memory stand-in for the substrate control API.
#[derive(Default)]
struct FakeAteApi {
    actors: Mutex<BTreeMap<(String, String), ActorSnapshot>>,
    next_uid: AtomicU32,
}

impl FakeAteApi {
    fn actor(&self, name: &str) -> Option<ActorSnapshot> {
        self.actors
            .lock()
            .unwrap()
            .get(&(NS.to_owned(), name.to_owned()))
            .cloned()
    }

    fn set_state(&self, name: &str, state: ActorState) {
        self.actors
            .lock()
            .unwrap()
            .get_mut(&(NS.to_owned(), name.to_owned()))
            .expect("the actor exists")
            .state = state;
    }
}

#[async_trait::async_trait]
impl ActorControl for FakeAteApi {
    async fn ensure_atespace(&self, _atespace: &str) -> Result<(), ActorControlError> {
        Ok(())
    }

    async fn get(
        &self,
        atespace: &str,
        name: &str,
    ) -> Result<Option<ActorSnapshot>, ActorControlError> {
        Ok(self
            .actors
            .lock()
            .unwrap()
            .get(&(atespace.to_owned(), name.to_owned()))
            .cloned())
    }

    async fn create(
        &self,
        atespace: &str,
        name: &str,
        template: &str,
    ) -> Result<ActorSnapshot, ActorControlError> {
        assert_eq!(
            template, TEMPLATE,
            "the actor comes from the configured template"
        );
        let next = self.next_uid.fetch_add(1, Ordering::SeqCst) + 1;
        let actor = ActorSnapshot {
            uid: format!("uid-{next}"),
            state: ActorState::Running,
            template: template.to_owned(),
        };
        self.actors
            .lock()
            .unwrap()
            .insert((atespace.to_owned(), name.to_owned()), actor.clone());
        Ok(actor)
    }

    async fn suspend(&self, atespace: &str, name: &str) -> Result<(), ActorControlError> {
        if let Some(actor) = self
            .actors
            .lock()
            .unwrap()
            .get_mut(&(atespace.to_owned(), name.to_owned()))
        {
            actor.state = ActorState::Suspended;
        }
        Ok(())
    }

    async fn resume(
        &self,
        _atespace: &str,
        _name: &str,
    ) -> Result<polyc_controller::ResumeOutcome, ActorControlError> {
        Ok(polyc_controller::ResumeOutcome::Running)
    }

    async fn set_template(
        &self,
        atespace: &str,
        name: &str,
        template: &str,
    ) -> Result<(), ActorControlError> {
        if let Some(actor) = self
            .actors
            .lock()
            .unwrap()
            .get_mut(&(atespace.to_owned(), name.to_owned()))
        {
            template.clone_into(&mut actor.template);
        }
        Ok(())
    }

    async fn delete(&self, atespace: &str, name: &str, uid: &str) -> Result<(), ActorControlError> {
        let key = (atespace.to_owned(), name.to_owned());
        let mut actors = self.actors.lock().unwrap();
        if actors.get(&key).is_some_and(|actor| actor.uid == uid) {
            actors.remove(&key);
        }
        drop(actors);
        Ok(())
    }
}

async fn ensure_namespace(client: &Client) {
    let api: Api<Namespace> = Api::all(client.clone());
    let ns = Namespace {
        metadata: ObjectMeta {
            name: Some(NS.to_owned()),
            ..ObjectMeta::default()
        },
        ..Namespace::default()
    };
    api.patch(
        NS,
        &PatchParams::apply("polychrome-e2e"),
        &Patch::Apply(&ns),
    )
    .await
    .expect("create/ensure namespace");
}

async fn reconcile_once(convs: &Api<Conversation>, name: &str, ctx: &Arc<Context>) {
    let latest = convs.get(name).await.expect("get conv");
    reconcile(Arc::new(latest), ctx.clone())
        .await
        .expect("reconcile pass");
}

#[tokio::test]
async fn conversation_reconciles_into_a_substrate_actor() {
    // rustls 0.23 won't auto-select a crypto provider when several are linked.
    let _ = rustls::crypto::ring::default_provider().install_default();

    let client = Client::try_default()
        .await
        .expect("kube client from current context (e.g. orbstack)");

    // Fail loudly if the CRDs aren't applied — that's a setup error, not a bug.
    assert!(
        Api::<Conversation>::all(client.clone())
            .list(&ListParams::default())
            .await
            .is_ok(),
        "Conversation CRD not installed; run: cargo run --bin crdgen | kubectl apply -f -"
    );

    ensure_namespace(&client).await;
    let convs: Api<Conversation> = Api::namespaced(client.clone(), NS);

    // Unique name per run so reruns don't collide with a finalizer-blocked object.
    let suffix = SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .unwrap()
        .as_millis();
    let name = format!("e2e-{suffix}");

    let mut conv = Conversation::new(
        &name,
        ConversationSpec {
            model: "fast-2".to_owned(),
            principal_ref: "persona-e2e".to_owned(),
            idle_timeout_seconds: 0,
            tools_enabled: vec![],
            tools_disabled: vec![],
            parent_conversation_id: None,
            agent_id: None,
        },
    );
    conv.metadata.namespace = Some(NS.to_owned());
    convs
        .create(&PostParams::default(), &conv)
        .await
        .expect("create Conversation");

    let ate_api = Arc::new(FakeAteApi::default());
    let ctx = Arc::new(Context::new(
        client.clone(),
        TEMPLATE.to_owned(),
        Arc::new(SubstrateActorBackend::new(
            Arc::clone(&ate_api) as Arc<dyn ActorControl>
        )),
        // Roll-on-image-change is off: this test exercises the ordinary
        // create, sync, suspend, and crash lifecycle.
        None,
        5,
        // Closed-conversation GC off — this test never closes the conversation.
        None,
        // Suspend one second after a turn ends.
        Some(1),
    ));

    // Pass 1 adds the finalizer. Pass 2 creates the actor and records it.
    for _ in 0..2 {
        reconcile_once(&convs, &name, &ctx).await;
    }
    let first = ate_api.actor(&name).expect("the actor was created");
    let recorded = convs
        .get_status(&name)
        .await
        .expect("get status")
        .status
        .unwrap_or_default();
    assert_eq!(recorded.actor_name.as_deref(), Some(name.as_str()));
    assert_eq!(recorded.actor_uid.as_deref(), Some(first.uid.as_str()));

    // Pass 3 mirrors readiness into the status.
    reconcile_once(&convs, &name, &ctx).await;
    let synced = convs
        .get_status(&name)
        .await
        .expect("get synced status")
        .status
        .unwrap_or_default();
    assert!(synced.harness_ready, "harness_ready should propagate");
    assert_eq!(synced.phase.as_deref(), Some("Ready"));

    // A turn starts and ends long ago: the actor is idle, so the next pass
    // suspends it. The actor keeps its identity.
    let long_ago = 1_000_i64;
    convs
        .patch_status(
            &name,
            &PatchParams::default(),
            &Patch::Merge(&json!({ "status": {
                "lastActivityUnix": long_ago,
                "lastTerminalUnix": long_ago + 1,
            } })),
        )
        .await
        .expect("stamp a finished turn");
    reconcile_once(&convs, &name, &ctx).await;
    assert_eq!(
        ate_api.actor(&name).expect("still exists").state,
        ActorState::Suspended
    );
    let paused = convs
        .get_status(&name)
        .await
        .unwrap()
        .status
        .unwrap_or_default();
    assert!(paused.idle_reclaimed);
    assert_eq!(paused.phase.as_deref(), Some("Paused"));

    // A crashed actor is deleted, and a new one takes its place.
    ate_api.set_state(&name, ActorState::Crashed);
    reconcile_once(&convs, &name, &ctx).await;
    assert!(
        ate_api.actor(&name).is_none(),
        "the crashed actor is deleted"
    );
    reconcile_once(&convs, &name, &ctx).await;
    let replacement = ate_api.actor(&name).expect("a new actor was created");
    assert_ne!(replacement.uid, first.uid, "a new incarnation");

    // Cleanup: delete sets a deletionTimestamp (finalizer blocks removal); one
    // more reconcile runs the cleanup branch, deletes the actor, drops the
    // finalizer, and lets the object go.
    convs
        .delete(&name, &DeleteParams::default())
        .await
        .expect("delete conv");
    if let Ok(terminating) = convs.get(&name).await {
        reconcile(Arc::new(terminating), ctx.clone())
            .await
            .expect("cleanup reconcile");
    }
    assert!(
        ate_api.actor(&name).is_none(),
        "the actor should be deleted after cleanup"
    );
    let _ = convs.get(&name).await.map(|c| c.name_any());
}