use std::sync::Arc;
use kube::ResourceExt;
use crate::conversation::Conversation;
use crate::execution_backend::{
ActorControl, ActorControlError, ActorSnapshot, ActorState, DialAddress, EnsuredUnit,
ExecutionBackend, UnitReadiness,
};
use crate::reconcile::Error;
pub struct SubstrateActorBackend {
control: Arc<dyn ActorControl>,
}
impl SubstrateActorBackend {
#[must_use]
pub fn new(control: Arc<dyn ActorControl>) -> Self {
Self { control }
}
}
impl From<ActorControlError> for Error {
fn from(error: ActorControlError) -> Self {
Self::Backend(error.to_string())
}
}
#[must_use]
pub fn actor_readiness(atespace: &str, name: &str, actor: Option<&ActorSnapshot>) -> UnitReadiness {
let Some(actor) = actor else {
return UnitReadiness::default();
};
let crashed = actor.state == ActorState::Crashed;
let deleting = actor.state == ActorState::Deleting;
UnitReadiness {
unit_present: true,
address: Some(DialAddress::Actor {
atespace: atespace.to_owned(),
name: name.to_owned(),
}),
harness_ready: !crashed && !deleting,
crashed,
uid: Some(actor.uid.clone()),
}
}
#[async_trait::async_trait]
impl ExecutionBackend for SubstrateActorBackend {
async fn ensure(
&self,
conv: &Conversation,
name: &str,
template: &str,
ns: &str,
) -> Result<EnsuredUnit, Error> {
self.control.ensure_atespace(ns).await?;
if let Some(existing) = self.control.get(ns, name).await? {
return Ok(EnsuredUnit { uid: existing.uid });
}
if let Some(parent) = conv.spec.parent_conversation_id.as_deref() {
tracing::info!(
child_conversation = %name,
parent_conversation = parent,
"creating actor for child conversation (handoff)"
);
}
match self.control.create(ns, name, template).await {
Ok(created) => {
tracing::info!(actor = name, template, ns, "created actor");
Ok(EnsuredUnit { uid: created.uid })
}
Err(error) => match self.control.get(ns, name).await? {
Some(existing) => Ok(EnsuredUnit { uid: existing.uid }),
None => Err(error.into()),
},
}
}
async fn readiness(&self, name: &str, ns: &str) -> Result<UnitReadiness, Error> {
let actor = self.control.get(ns, name).await?;
Ok(actor_readiness(ns, name, actor.as_ref()))
}
async fn suspend(&self, name: &str, ns: &str) -> Result<(), Error> {
let Some(actor) = self.control.get(ns, name).await? else {
return Ok(());
};
match actor.state {
ActorState::Running => Ok(self.control.suspend(ns, name).await?),
ActorState::Suspended
| ActorState::Paused
| ActorState::Crashed
| ActorState::Deleting => Ok(()),
ActorState::Resuming
| ActorState::Suspending
| ActorState::Pausing
| ActorState::Reverting => Err(Error::Backend(format!(
"actor {ns}/{name} is {:?}; it cannot be suspended yet",
actor.state
))),
}
}
async fn retemplate(&self, name: &str, template: &str, ns: &str) -> Result<(), Error> {
let Some(actor) = self.control.get(ns, name).await? else {
return Err(Error::Backend(format!("actor {ns}/{name} does not exist")));
};
if actor.template == template {
return Ok(());
}
if matches!(actor.state, ActorState::Running | ActorState::Resuming) {
self.control.suspend(ns, name).await?;
}
Ok(self.control.set_template(ns, name, template).await?)
}
async fn teardown(&self, owner: &Conversation, name: &str, ns: &str) -> Result<(), Error> {
let Some(actor) = self.control.get(ns, name).await? else {
return Ok(());
};
let recorded = owner.status.as_ref().and_then(|s| s.actor_uid.as_deref());
if recorded.is_some_and(|uid| uid != actor.uid) {
tracing::warn!(
actor = name,
conversation = %owner.name_any(),
"refusing to delete an actor of another incarnation"
);
return Ok(());
}
Ok(self.control.delete(ns, name, &actor.uid).await?)
}
fn kind(&self) -> &'static str {
"substrate-actor"
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use std::collections::BTreeMap;
use std::sync::Mutex;
use super::*;
use crate::conversation::{ConversationSpec, ConversationStatus};
const NS: &str = "polychrome";
const TPL: &str = "polychrome-harness-abc";
#[derive(Default)]
struct Fake {
actors: Mutex<BTreeMap<(String, String), ActorSnapshot>>,
calls: Mutex<Vec<String>>,
fail: Mutex<bool>,
}
impl Fake {
fn with(actor: ActorSnapshot) -> Arc<Self> {
let fake = Arc::new(Self::default());
fake.actors
.lock()
.unwrap()
.insert((NS.to_owned(), "c1".to_owned()), actor);
fake
}
fn calls(&self) -> Vec<String> {
self.calls.lock().unwrap().clone()
}
fn record(&self, call: String) -> Result<(), ActorControlError> {
self.calls.lock().unwrap().push(call);
if *self.fail.lock().unwrap() {
return Err(ActorControlError("unavailable".to_owned()));
}
Ok(())
}
}
#[async_trait::async_trait]
impl ActorControl for Fake {
async fn ensure_atespace(&self, atespace: &str) -> Result<(), ActorControlError> {
self.record(format!("ensure_atespace {atespace}"))
}
async fn get(
&self,
atespace: &str,
name: &str,
) -> Result<Option<ActorSnapshot>, ActorControlError> {
self.record(format!("get {atespace}/{name}"))?;
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> {
self.record(format!("create {atespace}/{name} from {template}"))?;
let actor = ActorSnapshot {
uid: "uid-new".to_owned(),
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> {
self.record(format!("suspend {atespace}/{name}"))?;
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<crate::execution_backend::ResumeOutcome, ActorControlError> {
self.record(format!("resume {atespace}/{name}"))?;
Ok(crate::execution_backend::ResumeOutcome::Running)
}
async fn set_template(
&self,
atespace: &str,
name: &str,
template: &str,
) -> Result<(), ActorControlError> {
self.record(format!("set_template {atespace}/{name} to {template}"))?;
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> {
self.record(format!("delete {atespace}/{name} uid {uid}"))?;
self.actors
.lock()
.unwrap()
.remove(&(atespace.to_owned(), name.to_owned()));
Ok(())
}
}
fn actor(uid: &str, state: ActorState) -> ActorSnapshot {
ActorSnapshot {
uid: uid.to_owned(),
state,
template: "polychrome-harness-old".to_owned(),
}
}
fn conv(recorded_uid: Option<&str>) -> Conversation {
let mut conv = Conversation::new(
"c1",
ConversationSpec {
model: String::new(),
principal_ref: String::new(),
idle_timeout_seconds: 300,
tools_enabled: Vec::new(),
tools_disabled: Vec::new(),
parent_conversation_id: None,
agent_id: None,
},
);
conv.metadata.namespace = Some(NS.to_owned());
conv.status = Some(ConversationStatus {
actor_uid: recorded_uid.map(str::to_owned),
..Default::default()
});
conv
}
fn backend(fake: &Arc<Fake>) -> SubstrateActorBackend {
SubstrateActorBackend::new(Arc::clone(fake) as Arc<dyn ActorControl>)
}
#[tokio::test]
async fn not_found_creates_the_actor_in_the_namespace_atespace() {
let fake = Arc::new(Fake::default());
let ensured = backend(&fake)
.ensure(&conv(None), "c1", TPL, NS)
.await
.expect("a missing actor is created");
assert_eq!(ensured.uid, "uid-new");
assert_eq!(
fake.calls(),
vec![
"ensure_atespace polychrome".to_owned(),
"get polychrome/c1".to_owned(),
"create polychrome/c1 from polychrome-harness-abc".to_owned(),
]
);
}
#[tokio::test]
async fn an_existing_actor_is_returned_without_a_second_create() {
let fake = Fake::with(actor("uid-old", ActorState::Suspended));
let ensured = backend(&fake)
.ensure(&conv(Some("uid-old")), "c1", TPL, NS)
.await
.expect("an existing actor is adopted");
assert_eq!(ensured.uid, "uid-old");
assert!(
!fake.calls().iter().any(|call| call.starts_with("create")),
"{:?}",
fake.calls()
);
}
#[test]
fn an_absent_actor_reads_as_not_present() {
assert_eq!(actor_readiness(NS, "c1", None), UnitReadiness::default());
}
#[test]
fn crashed_actor_is_present_but_not_ready() {
let readiness = actor_readiness(NS, "c1", Some(&actor("u", ActorState::Crashed)));
assert!(readiness.unit_present);
assert!(!readiness.harness_ready);
assert!(readiness.crashed);
assert_eq!(readiness.uid.as_deref(), Some("u"));
}
#[test]
fn every_live_state_is_ready_because_the_router_resumes_the_actor() {
for state in [
ActorState::Running,
ActorState::Suspended,
ActorState::Paused,
ActorState::Resuming,
ActorState::Suspending,
ActorState::Pausing,
ActorState::Reverting,
] {
let readiness = actor_readiness(NS, "c1", Some(&actor("u", state)));
assert!(readiness.unit_present, "{state:?}");
assert!(readiness.harness_ready, "{state:?}");
assert!(!readiness.crashed, "{state:?}");
assert_eq!(
readiness.address,
Some(DialAddress::Actor {
atespace: NS.to_owned(),
name: "c1".to_owned()
})
);
}
}
#[test]
fn a_deleting_actor_is_present_and_not_ready() {
let readiness = actor_readiness(NS, "c1", Some(&actor("u", ActorState::Deleting)));
assert!(readiness.unit_present);
assert!(!readiness.harness_ready);
assert!(!readiness.crashed);
}
#[tokio::test]
async fn readiness_reads_the_actor_through_the_control_port() {
let fake = Fake::with(actor("u", ActorState::Suspended));
let readiness = backend(&fake).readiness("c1", NS).await.unwrap();
assert!(readiness.unit_present && readiness.harness_ready);
assert_eq!(fake.calls(), vec!["get polychrome/c1".to_owned()]);
}
#[tokio::test]
async fn teardown_deletes_the_recorded_incarnation_pinned_to_its_uid() {
let fake = Fake::with(actor("uid-1", ActorState::Running));
backend(&fake)
.teardown(&conv(Some("uid-1")), "c1", NS)
.await
.unwrap();
assert!(
fake.calls()
.contains(&"delete polychrome/c1 uid uid-1".to_owned()),
"{:?}",
fake.calls()
);
}
#[tokio::test]
async fn teardown_refuses_an_actor_of_another_incarnation() {
let fake = Fake::with(actor("uid-newer", ActorState::Running));
backend(&fake)
.teardown(&conv(Some("uid-stale")), "c1", NS)
.await
.expect("a refusal is not an error");
assert!(
!fake.calls().iter().any(|call| call.starts_with("delete")),
"a stale reconcile must not delete a newer actor: {:?}",
fake.calls()
);
}
#[tokio::test]
async fn teardown_of_an_absent_actor_succeeds() {
let fake = Arc::new(Fake::default());
backend(&fake)
.teardown(&conv(Some("uid-1")), "c1", NS)
.await
.unwrap();
assert!(!fake.calls().iter().any(|call| call.starts_with("delete")));
}
#[tokio::test]
async fn retemplate_suspends_a_running_actor_before_it_changes_the_template() {
let fake = Fake::with(actor("u", ActorState::Running));
backend(&fake).retemplate("c1", TPL, NS).await.unwrap();
assert_eq!(
fake.calls(),
vec![
"get polychrome/c1".to_owned(),
"suspend polychrome/c1".to_owned(),
"set_template polychrome/c1 to polychrome-harness-abc".to_owned(),
]
);
}
#[tokio::test]
async fn retemplate_of_a_suspended_actor_only_changes_the_template() {
let fake = Fake::with(actor("u", ActorState::Suspended));
backend(&fake).retemplate("c1", TPL, NS).await.unwrap();
assert!(!fake.calls().iter().any(|call| call.starts_with("suspend")));
assert!(
fake.calls()
.iter()
.any(|call| call.starts_with("set_template"))
);
}
#[tokio::test]
async fn retemplate_of_an_actor_already_on_the_template_writes_nothing() {
let mut current = actor("u", ActorState::Running);
current.template = TPL.to_owned();
let fake = Fake::with(current);
backend(&fake).retemplate("c1", TPL, NS).await.unwrap();
assert_eq!(
fake.calls(),
vec!["get polychrome/c1".to_owned()],
"an actor that holds the template needs no suspend and no update"
);
}
#[tokio::test]
async fn a_repeated_retemplate_updates_the_actor_once_per_template_change() {
let fake = Fake::with(actor("u", ActorState::Running));
let backend = backend(&fake);
for _ in 0..5 {
backend.retemplate("c1", TPL, NS).await.unwrap();
}
let count = |prefix: &str| {
fake.calls()
.iter()
.filter(|call| call.starts_with(prefix))
.count()
};
assert_eq!(count("set_template"), 1, "{:?}", fake.calls());
assert_eq!(count("suspend"), 1, "{:?}", fake.calls());
backend
.retemplate("c1", "polychrome-harness-next", NS)
.await
.unwrap();
assert_eq!(count("set_template"), 2, "a new template is a new change");
}
#[tokio::test]
async fn suspend_of_a_running_actor_suspends_it() {
let fake = Fake::with(actor("u", ActorState::Running));
backend(&fake).suspend("c1", NS).await.unwrap();
assert!(fake.calls().contains(&"suspend polychrome/c1".to_owned()));
}
#[tokio::test]
async fn suspend_of_an_actor_that_is_not_running_is_a_no_op() {
for state in [
ActorState::Suspended,
ActorState::Paused,
ActorState::Crashed,
ActorState::Deleting,
] {
let fake = Fake::with(actor("u", state));
backend(&fake).suspend("c1", NS).await.unwrap();
assert!(
!fake.calls().iter().any(|call| call.starts_with("suspend")),
"{state:?}"
);
}
let absent = Arc::new(Fake::default());
backend(&absent).suspend("c1", NS).await.unwrap();
}
#[tokio::test]
async fn suspend_of_a_transitioning_actor_fails_so_the_reconciler_retries() {
for state in [
ActorState::Resuming,
ActorState::Suspending,
ActorState::Pausing,
ActorState::Reverting,
] {
let fake = Fake::with(actor("u", state));
let error = backend(&fake).suspend("c1", NS).await.unwrap_err();
assert!(matches!(error, Error::Backend(_)), "{state:?}: {error}");
}
}
#[tokio::test]
async fn a_control_failure_surfaces_as_a_backend_error() {
let fake = Arc::new(Fake::default());
*fake.fail.lock().unwrap() = true;
let error = backend(&fake).readiness("c1", NS).await.unwrap_err();
assert!(matches!(error, Error::Backend(_)), "{error}");
}
}