#![allow(clippy::unwrap_used)] #![allow(clippy::too_many_lines)] #![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";
#[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() {
let _ = rustls::crypto::ring::default_provider().install_default();
let client = Client::try_default()
.await
.expect("kube client from current context (e.g. orbstack)");
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);
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>
)),
None,
5,
None,
Some(1),
));
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()));
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"));
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"));
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");
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());
}