use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use async_trait::async_trait;
use meerkat_contracts::wire::WireAuthBindingRef;
use meerkat_contracts::wire::supervisor_bridge::{
BridgePeerIdentity, MaterializeLaunchMode, MaterializeLaunchOutcome,
};
use meerkat_contracts::wire::{
PortableMcpDecl, PortableMemberSpec, PortableSkillSource, PortableSystemPrompt,
WireResolvedToolAccessPolicy, WireSpawnContinuityIntent,
};
use meerkat_core::agent::CommsRuntime as CoreCommsRuntime;
use meerkat_core::comms::TrustedPeerDescriptor;
use meerkat_core::service::{SessionError, SessionProviderAuthFailure};
use meerkat_core::types::SessionId;
use meerkat_runtime::SessionServiceRuntimeExt as _;
use crate::MobRuntimeMode;
use crate::build::{BuildAgentConfigParams, BuildResumedAgentConfigParams};
use crate::definition::{MobDefinition, SkillSource};
use crate::profile::{Profile, ProfileBinding, ResumeOverrideField, ToolConfig};
use crate::runtime::SpawnSystemPromptOverride;
use crate::runtime::host_actor::{
HostCapabilityFacts, ProviderPresenceProbe, ProviderPresenceProbeError,
};
use crate::runtime::member_upcall::MemberUpcallBindingStamp;
use crate::runtime::provisioner::{
CommittedRuntimeSessionPublicationLease, MemberSessionDisposalArc,
MemberSessionDisposalVerdict, PreparedServiceActorTransaction, RuntimeSessionState,
RuntimeTurnFinalizationBoundaryLease,
};
use crate::runtime::session_service::{MobSessionService, ResumeSessionLoad};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MaterializeLlmPreflightOutcome {
Resolved,
ModelUnresolvable,
BindingUnresolvable,
}
#[async_trait]
pub trait MaterializePreflightProbe: ProviderPresenceProbe {
async fn preflight_llm_identity(
&self,
identity: &meerkat_core::SessionLlmIdentity,
custom_models: &BTreeMap<String, meerkat_core::config::CustomModelConfig>,
preferred_realm: Option<&meerkat_core::RealmId>,
auth_lease_handle: &meerkat_core::handles::GeneratedAuthLeaseHandle,
) -> Result<MaterializeLlmPreflightOutcome, ProviderPresenceProbeError>;
}
pub struct HostMemberSubstrate {
pub session_service: Arc<dyn MobSessionService>,
pub runtime_adapter: Arc<meerkat_runtime::MeerkatMachine>,
pub durable_event_log:
Option<Arc<dyn meerkat_runtime::member_observation::DurableEventLogRead>>,
pub realm_backend_persistent: bool,
pub member_identity_root: std::path::PathBuf,
pub preflight_probe: Arc<dyn MaterializePreflightProbe>,
}
#[derive(Debug, thiserror::Error)]
pub enum MaterializeDecompileError {
#[error("portable output_schema is not a valid schema: {detail}")]
OutputSchema { detail: String },
#[error("portable spec context is not valid JSON: {detail}")]
Context { detail: String },
#[error("mob tool authority context projection failed to rehydrate: {detail}")]
AuthorityContext { detail: String },
#[error(
"declared MCP stdio env key '{key}' for server '{server}' is absent from the host environment"
)]
McpEnvKeyMissing { server: String, key: String },
#[error(
"declared MCP HTTP header names for server '{server}' have no host-side value source \
in this engine version (v1 compiles them empty)"
)]
McpHeaderNamesUnsupported { server: String },
#[error("MCP connect timeout for server '{server}' exceeds the supported range")]
McpTimeoutOutOfRange { server: String },
}
#[derive(Debug, thiserror::Error)]
pub enum MaterializeServeError {
#[error(transparent)]
Decompile(#[from] MaterializeDecompileError),
#[error("resume session '{session_id}' not found on the member host")]
ResumeSessionNotFound { session_id: String },
#[error(
"recorded session '{session_id}' is terminal {state} and cannot be revived at the same tuple; the controller must advance the member generation"
)]
ResumeSessionNonRecoverable { session_id: String, state: String },
#[error("resume session id '{session_id}' is not a valid session id: {detail}")]
ResumeSessionIdInvalid { session_id: String, detail: String },
#[error("member build compile failed: {0}")]
Build(#[from] crate::error::MobError),
#[error("member comms runtime construction failed: {detail}")]
Comms { detail: String },
#[error("host-seeded supervisor bind failed: {detail}")]
SupervisorBind { detail: String },
#[error("member session service failed: {0}")]
SessionService(#[from] SessionError),
#[error("member LLM binding '{realm}/{binding}' became unresolvable after preflight: {detail}")]
AuthBindingUnresolvable {
realm: String,
binding: String,
detail: String,
},
#[error("{}", .0.public_message())]
ProviderAuth(SessionProviderAuthFailure),
#[error("session service returned '{created}' for admitted member session '{expected}'")]
IdentityMismatch { expected: String, created: String },
#[error("runtime session bindings preparation failed: {detail}")]
Bindings { detail: String },
#[error("member session executor attachment failed: {detail}")]
ExecutorAttach { detail: String },
#[error("member durable event floor read failed: {detail}")]
EventFloor { detail: String },
#[error(
"unrecorded member session '{session_id}' could not be proven quiescent after '{original}': {cleanup}"
)]
UnrecordedSessionCleanup {
session_id: String,
original: String,
cleanup: String,
},
#[error(
"revived member keypair diverged from the recorded pubkey (recorded '{recorded}', derived '{derived}')"
)]
RevivedIdentityDiverged { recorded: String, derived: String },
}
impl MaterializeServeError {
pub fn requires_actor_fail_stop(&self) -> bool {
matches!(self, Self::UnrecordedSessionCleanup { .. })
}
pub fn wire_cause(
&self,
) -> (
meerkat_contracts::wire::supervisor_bridge::BridgeRejectionCause,
String,
) {
use meerkat_contracts::wire::supervisor_bridge::BridgeRejectionCause as Cause;
match self {
Self::Decompile(err) => (Cause::Unsupported, err.to_string()),
Self::ResumeSessionNotFound { .. } | Self::ResumeSessionNonRecoverable { .. } => {
(Cause::ResumeSessionNotFound, self.to_string())
}
Self::ResumeSessionIdInvalid { .. } => (Cause::Unsupported, self.to_string()),
Self::AuthBindingUnresolvable { realm, binding, .. } => (
Cause::AuthBindingUnresolvable {
realm: realm.clone(),
binding: binding.clone(),
},
format!(
"materialize rejected: auth binding '{realm}/{binding}' is unresolvable on this host"
),
),
Self::ProviderAuth(failure) => (
Cause::MaterializeBuildRejected {
cause: meerkat_contracts::wire::supervisor_bridge::MemberBuildRejection::from_provider_auth_failure(
failure,
),
},
format!("materialize rejected: {}", failure.public_message()),
),
Self::Build(_)
| Self::Comms { .. }
| Self::SupervisorBind { .. }
| Self::SessionService(_)
| Self::IdentityMismatch { .. }
| Self::Bindings { .. }
| Self::ExecutorAttach { .. }
| Self::EventFloor { .. }
| Self::UnrecordedSessionCleanup { .. }
| Self::RevivedIdentityDiverged { .. } => (Cause::Internal, self.to_string()),
}
}
}
fn classify_materialize_create_error(
error: crate::error::MobError,
spec: &PortableMemberSpec,
) -> MaterializeServeError {
if let crate::error::MobError::SessionError(session_error) = &error
&& let Some(failure) = session_error.provider_auth_failure_data()
{
return MaterializeServeError::ProviderAuth(failure);
}
let llm_identity_unresolvable = matches!(
&error,
crate::error::MobError::SessionError(session_error)
if session_error.is_build_llm_identity_unresolvable()
);
if !llm_identity_unresolvable {
return MaterializeServeError::Build(error);
}
let (realm, binding) = spec
.overlay
.auth_binding
.as_ref()
.map(|binding| {
(
binding.realm.as_str().to_string(),
binding.binding.as_str().to_string(),
)
})
.unwrap_or_else(|| {
(
"env_default".to_string(),
spec.profile.provider.as_str().to_string(),
)
});
MaterializeServeError::AuthBindingUnresolvable {
realm,
binding,
detail: error.to_string(),
}
}
fn generation_start_seq_after(latest_seq: Option<u64>) -> Result<u64, MaterializeServeError> {
latest_seq.unwrap_or(0).checked_add(1).ok_or_else(|| {
MaterializeServeError::EventFloor {
detail: "resume event sequence space is exhausted at u64::MAX; no strictly later generation floor can be represented".to_string(),
}
})
}
fn combine_unregister_cleanup_result(
session_id: &SessionId,
context: &str,
unregister_error: Option<meerkat_runtime::RuntimeDriverError>,
fallback_result: Result<(), MaterializeServeError>,
) -> Result<(), MaterializeServeError> {
match (unregister_error, fallback_result) {
(None, result) => result,
(Some(meerkat_runtime::RuntimeDriverError::NotFound { .. }), Ok(())) => Ok(()),
(Some(unregister), Ok(())) => Err(MaterializeServeError::Bindings {
detail: format!(
"{context}: runtime unregister failed for session {session_id}: {unregister}"
),
}),
(Some(unregister), Err(fallback)) => Err(MaterializeServeError::Bindings {
detail: format!(
"{context}: runtime unregister failed for session {session_id}: {unregister}; fallback cleanup also failed: {fallback}"
),
}),
}
}
fn combine_prepared_session_unregister_result(
session_id: &SessionId,
original: MaterializeServeError,
unregister_result: Result<(), meerkat_runtime::RuntimeDriverError>,
) -> MaterializeServeError {
match unregister_result {
Ok(()) => original,
Err(cleanup) => MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: original.to_string(),
cleanup: format!("prepared runtime unregister failed: {cleanup}"),
},
}
}
pub fn validate_portable_spec_structure(
spec: &PortableMemberSpec,
) -> Result<(), MaterializeDecompileError> {
decompile_portable_spec_with_env(spec, &|_| Some(String::new())).map(drop)
}
pub fn decompile_portable_spec(
spec: &PortableMemberSpec,
) -> Result<DecompiledMemberBuild, MaterializeDecompileError> {
decompile_portable_spec_with_env(spec, &|key| std::env::var(key).ok())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DecompiledSystemPrompt {
Replace(String),
Disable,
}
pub struct DecompiledMemberBuild {
pub mob_id: crate::ids::MobId,
pub profile_name: crate::ids::ProfileName,
pub agent_identity: crate::ids::AgentIdentity,
pub profile: Profile,
pub definition: MobDefinition,
pub context: Option<serde_json::Value>,
pub labels: Option<BTreeMap<String, String>>,
pub additional_instructions: Option<Vec<String>>,
pub system_prompt: DecompiledSystemPrompt,
pub tool_access_policy: Option<meerkat_core::ops::ToolAccessPolicy>,
pub mob_tool_authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
pub auth_binding: Option<meerkat_core::AuthBindingRef>,
pub budget_limits: Option<meerkat_core::BudgetLimits>,
pub requester_resolved_policy: Option<WireResolvedToolAccessPolicy>,
}
fn wire_runtime_mode_to_domain(
mode: meerkat_contracts::wire::WireMobRuntimeMode,
) -> MobRuntimeMode {
match mode {
meerkat_contracts::wire::WireMobRuntimeMode::AutonomousHost => {
MobRuntimeMode::AutonomousHost
}
meerkat_contracts::wire::WireMobRuntimeMode::TurnDriven => MobRuntimeMode::TurnDriven,
}
}
fn wire_resume_override_to_domain(
field: meerkat_contracts::wire::WireMobResumeOverrideField,
) -> ResumeOverrideField {
match field {
meerkat_contracts::wire::WireMobResumeOverrideField::Model => ResumeOverrideField::Model,
meerkat_contracts::wire::WireMobResumeOverrideField::Provider => {
ResumeOverrideField::Provider
}
meerkat_contracts::wire::WireMobResumeOverrideField::ProviderParams => {
ResumeOverrideField::ProviderParams
}
}
}
fn decompile_mcp_servers(
decls: &BTreeMap<String, PortableMcpDecl>,
env_lookup: &impl Fn(&str) -> Option<String>,
) -> Result<Vec<meerkat_core::mcp_config::McpServerConfig>, MaterializeDecompileError> {
let mut servers = Vec::with_capacity(decls.len());
for (name, decl) in decls {
match decl {
PortableMcpDecl::Stdio {
command,
args,
required_env_keys,
connect_timeout_secs,
} => {
let mut env = HashMap::with_capacity(required_env_keys.len());
for key in required_env_keys {
let value = env_lookup(key).ok_or_else(|| {
MaterializeDecompileError::McpEnvKeyMissing {
server: name.clone(),
key: key.clone(),
}
})?;
env.insert(key.clone(), value);
}
let mut server = meerkat_core::mcp_config::McpServerConfig::stdio(
name.clone(),
command.clone(),
args.clone(),
env,
);
if let Some(secs) = connect_timeout_secs {
server.connect_timeout_secs = Some(u32::try_from(*secs).map_err(|_| {
MaterializeDecompileError::McpTimeoutOutOfRange {
server: name.clone(),
}
})?);
}
servers.push(server);
}
PortableMcpDecl::Http {
url,
http_transport,
required_header_names,
connect_timeout_secs,
} => {
if !required_header_names.is_empty() {
return Err(MaterializeDecompileError::McpHeaderNamesUnsupported {
server: name.clone(),
});
}
servers.push(meerkat_core::mcp_config::McpServerConfig {
name: name.clone(),
transport: meerkat_core::mcp_config::McpTransportConfig::Http(
meerkat_core::mcp_config::McpHttpConfig {
url: url.clone(),
headers: HashMap::new(),
transport: *http_transport,
},
),
connect_timeout_secs: connect_timeout_secs
.map(|seconds| {
u32::try_from(seconds).map_err(|_| {
MaterializeDecompileError::McpTimeoutOutOfRange {
server: name.clone(),
}
})
})
.transpose()?,
});
}
}
}
Ok(servers)
}
fn decompile_portable_spec_with_env(
spec: &PortableMemberSpec,
env_lookup: &impl Fn(&str) -> Option<String>,
) -> Result<DecompiledMemberBuild, MaterializeDecompileError> {
let output_schema = spec
.profile
.output_schema
.as_ref()
.map(|envelope| {
let value =
envelope
.to_value()
.map_err(|err| MaterializeDecompileError::OutputSchema {
detail: err.to_string(),
})?;
meerkat_core::MeerkatSchema::new(value).map_err(|err| {
MaterializeDecompileError::OutputSchema {
detail: err.to_string(),
}
})
})
.transpose()?;
let profile = Profile {
model: spec.profile.model.clone(),
provider: Some(spec.profile.provider),
self_hosted_server_id: spec.profile.self_hosted_server_id.clone(),
image_generation_provider: spec.profile.image_generation_provider,
auto_compact_threshold: spec.profile.auto_compact_threshold,
resume_overrides: spec
.profile
.resume_overrides
.iter()
.copied()
.map(wire_resume_override_to_domain)
.collect(),
skills: spec.profile.skills.clone(),
tools: ToolConfig {
builtins: spec.profile.tools.builtins,
shell: spec.profile.tools.shell,
comms: spec.profile.tools.comms,
memory: spec.profile.tools.memory,
workgraph: spec.profile.tools.workgraph,
mob: spec.profile.tools.mob,
schedule: spec.profile.tools.schedule,
image_generation: spec.profile.tools.image_generation,
mcp: Vec::new(),
mcp_servers: decompile_mcp_servers(&spec.profile.tools.mcp_servers, env_lookup)?,
rust_bundles: Vec::new(),
},
peer_description: spec.profile.peer_description.clone(),
external_addressable: spec.profile.external_addressable,
backend: None,
runtime_mode: wire_runtime_mode_to_domain(spec.profile.runtime_mode),
max_inline_peer_notifications: spec.profile.max_inline_peer_notifications,
output_schema,
provider_params: spec.profile.provider_params.clone().map(Into::into),
};
let mut definition = MobDefinition::explicit(spec.mob_id.as_str());
definition.models = spec.definition_extract.models.clone();
definition.image_generation_provider = spec.definition_extract.image_generation_provider;
definition.skills = spec
.definition_extract
.skills
.iter()
.map(|(name, source)| {
let PortableSkillSource::Inline { content } = source;
(
name.clone(),
SkillSource::Inline {
content: content.clone(),
},
)
})
.collect();
for name in &spec.definition_extract.profile_names {
definition.profiles.insert(
crate::ids::ProfileName::from(name.as_str()),
ProfileBinding::RealmRef {
realm_profile: name.clone(),
},
);
}
let context = spec
.overlay
.context
.as_ref()
.map(|envelope| {
envelope
.to_value()
.map_err(|err| MaterializeDecompileError::Context {
detail: err.to_string(),
})
})
.transpose()?;
let system_prompt = match &spec.overlay.system_prompt {
PortableSystemPrompt::Set { text } => DecompiledSystemPrompt::Replace(text.clone()),
PortableSystemPrompt::Disable => DecompiledSystemPrompt::Disable,
};
let tool_access_policy = spec.overlay.tool_access_policy.as_ref().map(|policy| {
match policy {
WireResolvedToolAccessPolicy::AllowList(names) => {
meerkat_core::ops::ToolAccessPolicy::AllowList(names.iter().cloned().collect())
}
WireResolvedToolAccessPolicy::DenyList(names) => {
meerkat_core::ops::ToolAccessPolicy::DenyList(names.iter().cloned().collect())
}
}
});
let mob_tool_authority_context = spec
.overlay
.mob_tool_authority_context
.as_ref()
.map(|wire| {
serde_json::to_value(wire)
.and_then(serde_json::from_value::<meerkat_core::service::MobToolAuthorityContext>)
.map_err(|err| MaterializeDecompileError::AuthorityContext {
detail: err.to_string(),
})
})
.transpose()?;
Ok(DecompiledMemberBuild {
mob_id: crate::ids::MobId::from(spec.mob_id.as_str()),
profile_name: crate::ids::ProfileName::from(spec.profile_name.as_str()),
agent_identity: crate::ids::AgentIdentity::from(spec.agent_identity.as_str()),
profile,
definition,
context,
labels: spec.overlay.labels.clone(),
additional_instructions: spec.overlay.additional_instructions.clone(),
system_prompt,
tool_access_policy,
mob_tool_authority_context,
auth_binding: spec
.overlay
.auth_binding
.clone()
.map(meerkat_core::AuthBindingRef::from),
budget_limits: spec.overlay.budget_limits.clone(),
requester_resolved_policy: spec.overlay.tool_access_policy.clone(),
})
}
pub struct MaterializePreflightObservations {
pub model_resolvable: bool,
pub binding_resolvable: bool,
pub env_keys_present: bool,
pub stdio_commands_present: bool,
pub engine_protocol_supported: bool,
pub durable_sessions_required: bool,
pub realm_backend_persistent: bool,
pub memory_required: bool,
pub memory_capability: bool,
pub first_missing_env_key: Option<String>,
pub first_missing_stdio_server: Option<String>,
}
fn stdio_command_discoverable(command: &str) -> bool {
let path = std::path::Path::new(command);
if path.is_absolute() || command.contains(std::path::MAIN_SEPARATOR) {
return path.is_file();
}
let Some(search) = std::env::var_os("PATH") else {
return false;
};
std::env::split_paths(&search).any(|dir| dir.join(command).is_file())
}
pub async fn assemble_preflight_observations(
spec: &PortableMemberSpec,
launch: &MaterializeLaunchMode,
protocol_version: meerkat_contracts::wire::supervisor_bridge::BridgeProtocolVersion,
substrate: &HostMemberSubstrate,
capability_facts: HostCapabilityFacts,
) -> Result<MaterializePreflightObservations, ProviderPresenceProbeError> {
let candidate_identity = meerkat_core::SessionLlmIdentity {
model: spec.profile.model.clone(),
provider: spec.profile.provider,
self_hosted_server_id: spec.profile.self_hosted_server_id.clone(),
provider_params: spec.profile.provider_params.clone().map(Into::into),
auth_binding: spec
.overlay
.auth_binding
.clone()
.map(meerkat_core::AuthBindingRef::from),
};
let identity = match launch {
MaterializeLaunchMode::Fresh {} => Some(candidate_identity),
MaterializeLaunchMode::Resume { session_id } => match SessionId::parse(session_id) {
Err(_) => {
None
}
Ok(session_id) => match substrate
.session_service
.load_persisted_session_metadata(&session_id)
.await
.map_err(|error| ProviderPresenceProbeError::PreflightInput {
detail: format!(
"resume session '{session_id}' metadata preflight failed: {error}"
),
})? {
Some(view) => {
let metadata = view.session_metadata.ok_or_else(|| {
ProviderPresenceProbeError::PreflightInput {
detail: format!(
"resume session '{session_id}' has no durable session metadata"
),
}
})?;
let mut mask = meerkat_core::service::ResumeOverrideMask::default();
for field in &spec.profile.resume_overrides {
match field {
meerkat_contracts::wire::WireMobResumeOverrideField::Model => {
mask.model = true;
}
meerkat_contracts::wire::WireMobResumeOverrideField::Provider => {
mask.provider = true;
}
meerkat_contracts::wire::WireMobResumeOverrideField::ProviderParams => {
mask.provider_params = true;
}
}
}
Some(crate::build::effective_resumed_session_llm_identity(
candidate_identity,
&metadata,
mask,
))
}
None => None,
},
},
};
let preferred_realm = meerkat_core::mob_realm_id(&spec.mob_id).map_err(|error| {
ProviderPresenceProbeError::PreflightInput {
detail: format!("invalid mob identity '{}': {error}", spec.mob_id),
}
})?;
let auth_lease_handle = substrate.runtime_adapter.generated_auth_lease_handle();
let (model_resolvable, binding_resolvable) = match identity {
Some(identity) => match substrate
.preflight_probe
.preflight_llm_identity(
&identity,
&spec.definition_extract.models,
Some(&preferred_realm),
&auth_lease_handle,
)
.await?
{
MaterializeLlmPreflightOutcome::Resolved => (true, true),
MaterializeLlmPreflightOutcome::ModelUnresolvable => (false, true),
MaterializeLlmPreflightOutcome::BindingUnresolvable => (true, false),
},
None => (true, true),
};
let mut first_missing_env_key = None;
for key in &spec.required_env_keys {
if std::env::var_os(key).is_none() {
first_missing_env_key = Some(key.clone());
break;
}
}
let mut first_missing_stdio_server = None;
for (server, decl) in &spec.profile.tools.mcp_servers {
if let PortableMcpDecl::Stdio { command, .. } = decl
&& !stdio_command_discoverable(command)
{
first_missing_stdio_server = Some(server.clone());
break;
}
}
Ok(MaterializePreflightObservations {
model_resolvable,
binding_resolvable,
env_keys_present: first_missing_env_key.is_none(),
stdio_commands_present: first_missing_stdio_server.is_none(),
first_missing_env_key,
first_missing_stdio_server,
engine_protocol_supported: protocol_version.supports_multi_host(),
durable_sessions_required: matches!(launch, MaterializeLaunchMode::Resume { .. })
|| matches!(
spec.overlay.continuity_intent,
WireSpawnContinuityIntent::DurableIdentity { .. }
),
realm_backend_persistent: substrate.realm_backend_persistent,
memory_required: spec.profile.tools.memory,
memory_capability: capability_facts.memory_store,
})
}
#[derive(Clone)]
pub struct LiveMemberRuntime {
pub runtime: Arc<meerkat_comms::CommsRuntime>,
pub ack_keypair: Arc<meerkat_comms::Keypair>,
}
pub struct MaterializedBuildOutcome {
pub session_id: SessionId,
pub residency_update: meerkat_runtime::meerkat_machine::MemberResidencyUpdate,
pub(super) runtime_publication: Option<CommittedRuntimeSessionPublicationLease>,
pub member_runtime: Arc<meerkat_comms::CommsRuntime>,
pub ack_keypair: Arc<meerkat_comms::Keypair>,
pub member_pubkey: String,
pub member_peer_id: String,
pub launch_outcome: MaterializeLaunchOutcome,
pub resolved_auth_binding: Option<WireAuthBindingRef>,
pub generation_start_seq: u64,
}
pub struct MaterializeServingContext<'a> {
pub generation: u64,
pub fence_token: u64,
pub host_id: &'a str,
pub host_binding_generation: u64,
pub supervisor: &'a BridgePeerIdentity,
pub epoch: u64,
}
pub struct RevivedMemberOutcome {
pub member: LiveMemberRuntime,
pub residency_update: Option<meerkat_runtime::meerkat_machine::MemberResidencyUpdate>,
pub(super) runtime_publication: Option<CommittedRuntimeSessionPublicationLease>,
}
pub struct HostMemberMaterializer {
substrate: HostMemberSubstrate,
disposal: MemberSessionDisposalArc,
live_runtimes: HashMap<SessionId, LiveMemberRuntime>,
upcall_binding_stamps: HashMap<SessionId, Arc<MemberUpcallBindingStamp>>,
runtime_sessions: Arc<tokio::sync::RwLock<HashMap<SessionId, Arc<RuntimeSessionState>>>>,
}
const fn member_runtime_is_healthy(
materializer_runtime: bool,
executor_resident: bool,
machine_serving: bool,
service_live: bool,
) -> bool {
materializer_runtime && executor_resident && machine_serving && service_live
}
impl HostMemberMaterializer {
pub fn new(substrate: HostMemberSubstrate) -> Self {
let runtime_sessions = Arc::new(tokio::sync::RwLock::new(HashMap::new()));
let disposal = MemberSessionDisposalArc::with_runtime_sessions(
Arc::clone(&substrate.session_service),
Some(Arc::clone(&substrate.runtime_adapter)),
Arc::clone(&runtime_sessions),
);
Self {
substrate,
disposal,
live_runtimes: HashMap::new(),
upcall_binding_stamps: HashMap::new(),
runtime_sessions,
}
}
async fn acquire_vacant_materialization_boundary(
&self,
session_id: &SessionId,
) -> Result<RuntimeTurnFinalizationBoundaryLease, MaterializeServeError> {
let boundary = RuntimeTurnFinalizationBoundaryLease::acquire(
&self.substrate.session_service,
session_id,
)
.await
.map_err(MaterializeServeError::Build)?;
let machine_registered = self
.substrate
.runtime_adapter
.contains_session(session_id)
.await;
let sidecar_registered = self.runtime_sessions.read().await.contains_key(session_id);
let actor_registered = self
.substrate
.session_service
.live_session_actor_registered(session_id)
.await
.map_err(MaterializeServeError::SessionService)?;
if machine_registered || sidecar_registered || actor_registered {
return Err(MaterializeServeError::Bindings {
detail: format!(
"session {session_id} acquired a replacement process carrier before exact host materialization could claim B"
),
});
}
Ok(boundary)
}
async fn attach_member_executor(
&self,
transaction: PreparedServiceActorTransaction,
) -> Result<CommittedRuntimeSessionPublicationLease, MaterializeServeError> {
transaction
.prepare_host_runtime_publication_owned(
Arc::clone(&self.substrate.runtime_adapter),
Arc::clone(&self.runtime_sessions),
)
.await
.map_err(|error| MaterializeServeError::ExecutorAttach {
detail: error.to_string(),
})
}
async fn prepare_member_runtime_placement(
&self,
session_id: &SessionId,
agent_identity: &crate::ids::AgentIdentity,
generation: u64,
fence_token: u64,
) -> Result<(), MaterializeServeError> {
let agent_runtime_id = crate::ids::AgentRuntimeId::new(
agent_identity.clone(),
crate::ids::Generation::new(generation),
);
self.substrate
.runtime_adapter
.prepare_runtime_placement_binding(
session_id.clone(),
meerkat_runtime::identifiers::LogicalRuntimeId::new(agent_runtime_id.to_string()),
fence_token,
generation,
)
.await
.map_err(|error| MaterializeServeError::Bindings {
detail: error.to_string(),
})
}
pub fn substrate(&self) -> &HostMemberSubstrate {
&self.substrate
}
pub fn live_runtime(&self, session_id: &SessionId) -> Option<&LiveMemberRuntime> {
self.live_runtimes.get(session_id)
}
pub async fn refresh_live_supervisor(
&mut self,
recorded_session_id: &str,
supervisor: &TrustedPeerDescriptor,
epoch: u64,
host_id: &str,
host_binding_generation: u64,
) -> Result<bool, MaterializeServeError> {
let session_id = SessionId::parse(recorded_session_id).map_err(|error| {
MaterializeServeError::ResumeSessionIdInvalid {
session_id: recorded_session_id.to_string(),
detail: error.to_string(),
}
})?;
let Some(live) = self.live_runtimes.get(&session_id).cloned() else {
return Ok(false);
};
meerkat_runtime::comms_drain::authorize_supervisor_for_materialized_session(
self.substrate.runtime_adapter.as_ref(),
&session_id,
live.runtime.as_ref() as &dyn CoreCommsRuntime,
supervisor,
epoch,
)
.await
.map_err(|error| MaterializeServeError::SupervisorBind {
detail: error.to_string(),
})?;
if let Some(binding_stamp) = self.upcall_binding_stamps.get(&session_id) {
let current = binding_stamp.snapshot();
binding_stamp.update(
current.generation,
current.fence_token,
host_id,
host_binding_generation,
current.member_session_id,
);
}
Ok(true)
}
pub async fn session_live(&self, session_id: &SessionId) -> bool {
let Some(live) = self.live_runtimes.get(session_id) else {
return false;
};
let comms_runtime: Arc<dyn CoreCommsRuntime> = live.runtime.clone();
let sidecar = self.runtime_sessions.read().await.get(session_id).cloned();
let machine_witness = self
.substrate
.runtime_adapter
.current_executor_attachment_witness(session_id)
.await;
let executor_resident = sidecar.as_ref().is_some_and(|state| {
machine_witness.as_ref().is_some_and(|witness| {
state.witness() == witness && state.attachment_is_active(witness)
})
});
let service_live = match self
.substrate
.session_service
.live_session_actor_registered(session_id)
.await
{
Ok(live) => live,
Err(error) => {
tracing::warn!(
%session_id,
error = %error,
"mob host: failed to read live member session-service witness"
);
false
}
};
let machine_serving = match self
.substrate
.runtime_adapter
.materialized_member_runtime_is_serving(session_id, &comms_runtime)
.await
{
Ok(serving) => serving,
Err(error) => {
tracing::warn!(
%session_id,
error = %error,
"mob host: failed to read materialized member serving witness"
);
false
}
};
member_runtime_is_healthy(true, executor_resident, machine_serving, service_live)
}
async fn quiesce_nonserving_incarnation(
&mut self,
session_id: &SessionId,
) -> Result<(), MaterializeServeError> {
let old_state = self.runtime_sessions.read().await.get(session_id).cloned();
let unregister_error = if let Some(old_state) = old_state.as_ref() {
match self
.substrate
.runtime_adapter
.unregister_executor_attachment_if_current(old_state.witness())
.await
{
Ok(true) => None,
Ok(false)
if self
.substrate
.runtime_adapter
.current_executor_attachment_witness(session_id)
.await
.is_none() =>
{
None
}
Ok(false) => {
return Err(MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} changed executor attachment before quiescence; refusing to touch the replacement"
),
});
}
Err(error) => Some(error),
}
} else {
let current_attachment = self
.substrate
.runtime_adapter
.current_executor_attachment_witness(session_id)
.await;
let service_live = self
.substrate
.session_service
.live_session_actor_registered(session_id)
.await
.map_err(MaterializeServeError::SessionService)?;
if current_attachment.is_some() {
return Err(MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} has a current executor attachment without its exact host sidecar; refusing to touch a possible replacement"
),
});
}
let registration = self
.substrate
.runtime_adapter
.current_session_registration_witness(session_id)
.await;
match registration {
Some(registration) => match self
.substrate
.runtime_adapter
.unregister_terminal_session_registration_if_current(®istration)
.await
{
Ok(true) => None,
Ok(false) => {
return Err(MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} changed terminal registration before exact quiescence"
),
});
}
Err(
meerkat_runtime::RuntimeDriverError::ValidationFailed { .. }
| meerkat_runtime::RuntimeDriverError::UnregisterInProgress { .. }
| meerkat_runtime::RuntimeDriverError::RuntimeStopInProgress { .. },
) => None,
Err(error) => Some(error),
},
None if service_live => {
return Err(MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} has a live service actor without exact machine registration authority"
),
});
}
None => None,
}
};
let fallback_result = async {
const CLEANUP_GRACE: std::time::Duration = std::time::Duration::from_secs(30);
tokio::time::timeout(CLEANUP_GRACE, async {
loop {
let runtime_registered = self
.substrate
.runtime_adapter
.contains_session(session_id)
.await;
let service_live = self
.substrate
.session_service
.live_session_actor_registered(session_id)
.await
.map_err(MaterializeServeError::SessionService)?;
let executor_resident =
self.runtime_sessions.read().await.contains_key(session_id);
if !runtime_registered && !service_live && !executor_resident {
return Ok::<(), MaterializeServeError>(());
}
if runtime_registered && !service_live && !executor_resident {
let current_attachment = self
.substrate
.runtime_adapter
.current_executor_attachment_witness(session_id)
.await;
if current_attachment.is_none()
&& let Some(registration) = self
.substrate
.runtime_adapter
.current_session_registration_witness(session_id)
.await
{
match self
.substrate
.runtime_adapter
.unregister_terminal_session_registration_if_current(®istration)
.await
{
Ok(_)
| Err(
meerkat_runtime::RuntimeDriverError::NotFound { .. }
| meerkat_runtime::RuntimeDriverError::ValidationFailed {
..
}
| meerkat_runtime::RuntimeDriverError::UnregisterInProgress {
..
}
| meerkat_runtime::RuntimeDriverError::RuntimeStopInProgress {
..
},
) => {}
Err(error) => {
return Err(MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} exact terminal registration cleanup failed: {error}"
),
});
}
}
}
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.map_err(|_| MaterializeServeError::Bindings {
detail: format!(
"non-serving incarnation {session_id} did not quiesce within {}s",
CLEANUP_GRACE.as_secs()
),
})??;
self.live_runtimes.remove(session_id);
self.upcall_binding_stamps.remove(session_id);
Ok(())
}
.await;
combine_unregister_cleanup_result(
session_id,
"non-serving incarnation quiescence",
unregister_error,
fallback_result,
)
}
pub async fn dispose(
&mut self,
session_id: &SessionId,
) -> Result<MemberSessionDisposalVerdict, SessionError> {
let verdict = self.disposal.dispose(session_id).await?;
self.live_runtimes.remove(session_id);
self.upcall_binding_stamps.remove(session_id);
Ok(verdict)
}
pub async fn release_runtime_only(
&mut self,
session_id: &SessionId,
) -> Result<(), SessionError> {
self.disposal.release_runtime_only(session_id).await?;
self.live_runtimes.remove(session_id);
self.upcall_binding_stamps.remove(session_id);
Ok(())
}
async fn quiesce_generation_cutover(
&mut self,
session_id: &SessionId,
projection_witness_required: bool,
) -> Result<(), MaterializeServeError> {
self.quiesce_nonserving_incarnation(session_id).await?;
const PROJECTION_DRAIN_GRACE: std::time::Duration = std::time::Duration::from_secs(30);
let witnessed = tokio::time::timeout(
PROJECTION_DRAIN_GRACE,
self.substrate
.session_service
.await_event_projection_drain(session_id),
)
.await
.map_err(|_| MaterializeServeError::EventFloor {
detail: format!(
"generation cutover event projection did not drain within {}s",
PROJECTION_DRAIN_GRACE.as_secs()
),
})?
.map_err(|error| MaterializeServeError::EventFloor {
detail: format!("generation cutover event projection failed: {error}"),
})?;
if projection_witness_required && !witnessed {
return Err(MaterializeServeError::EventFloor {
detail: "generation cutover lost the live projection drain witness".to_string(),
});
}
Ok(())
}
async fn quiesce_unrecorded_session(&mut self, session_id: &SessionId) -> Result<(), String> {
self.quiesce_nonserving_incarnation(session_id)
.await
.map_err(|error| {
format!("failed to quiesce session while preserving its durable snapshot: {error}")
})
}
pub fn forget_runtime_after_exact_publication_abort(&mut self, session_id: &SessionId) {
self.live_runtimes.remove(session_id);
self.upcall_binding_stamps.remove(session_id);
}
pub async fn quiesce_before_superseding_record(
&mut self,
session_id: &SessionId,
) -> Result<(), String> {
self.quiesce_unrecorded_session(session_id).await
}
pub async fn dispose_after_superseding_commit(
&mut self,
session_id: &SessionId,
) -> Result<(), String> {
let first_disposal_error = self
.dispose(session_id)
.await
.err()
.map(|error| error.to_string());
if let Err(error) = self.quiesce_nonserving_incarnation(session_id).await {
return Err(match first_disposal_error {
Some(disposal) => format!(
"initial superseded-session disposal failed ({disposal}); forced runtime quiescence also failed ({error})"
),
None => format!(
"forced runtime quiescence failed after superseded-session disposal: {error}"
),
});
}
if let Some(first_error) = first_disposal_error
&& let Err(retry_error) = self.dispose(session_id).await
{
return Err(format!(
"initial superseded-session disposal failed ({first_error}); runtime quiesced, but archive retry failed ({retry_error})"
));
}
Ok(())
}
async fn preserve_snapshot_after_revival_failure(
&mut self,
session_id: &SessionId,
original: MaterializeServeError,
) -> MaterializeServeError {
match self.quiesce_unrecorded_session(session_id).await {
Ok(()) => original,
Err(cleanup) => MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: original.to_string(),
cleanup,
},
}
}
async fn rollback_prepared_session_after_failure(
&self,
prepared: &mut meerkat_runtime::PreparedSessionMaterialization,
original: MaterializeServeError,
) -> MaterializeServeError {
let session_id = prepared.session_id().clone();
let rollback_result = prepared
.rollback_now_under_turn_finalization_boundary()
.await;
combine_prepared_session_unregister_result(
&session_id,
original,
rollback_result.map(|_| ()),
)
}
pub async fn materialize(
&mut self,
spec: &PortableMemberSpec,
launch: &MaterializeLaunchMode,
context: MaterializeServingContext<'_>,
) -> Result<MaterializedBuildOutcome, MaterializeServeError> {
let MaterializeServingContext {
generation,
fence_token,
host_id,
host_binding_generation,
supervisor,
epoch,
} = context;
let decompiled = decompile_portable_spec(spec)?;
let (
session_id,
resumed_session,
launch_outcome,
generation_start_seq,
residency_update,
mut materialization_boundary,
) = match launch {
MaterializeLaunchMode::Fresh {} => {
let id = SessionId::new();
let boundary = self.acquire_vacant_materialization_boundary(&id).await?;
let residency_update = self
.substrate
.runtime_adapter
.begin_member_residency_update(id.clone())
.await;
(
id,
None,
MaterializeLaunchOutcome::Fresh,
1,
residency_update,
Some(boundary),
)
}
MaterializeLaunchMode::Resume { session_id } => {
let id = SessionId::parse(session_id).map_err(|err| {
MaterializeServeError::ResumeSessionIdInvalid {
session_id: session_id.clone(),
detail: err.to_string(),
}
})?;
let residency_update = self
.substrate
.runtime_adapter
.begin_member_residency_update(id.clone())
.await;
let was_live = if self.live_runtimes.contains_key(&id) {
true
} else {
self.substrate
.session_service
.live_session_actor_registered(&id)
.await
.map_err(MaterializeServeError::SessionService)?
};
self.quiesce_generation_cutover(&id, was_live).await?;
let boundary = self.acquire_vacant_materialization_boundary(&id).await?;
let log = self.substrate.durable_event_log.as_ref().ok_or_else(|| {
MaterializeServeError::EventFloor {
detail: "resume requires a durable event projection".to_string(),
}
})?;
let latest_seq = log.latest_seq(&id).await.map_err(|error| {
MaterializeServeError::EventFloor {
detail: error.to_string(),
}
})?;
let generation_start_seq = generation_start_seq_after(latest_seq)?;
match self
.substrate
.session_service
.materialize_session_for_resume(&id)
.await?
{
ResumeSessionLoad::Active(session) | ResumeSessionLoad::Revivable(session) => (
id,
Some(*session),
MaterializeLaunchOutcome::ResumedFromSnapshot,
generation_start_seq,
residency_update,
Some(boundary),
),
ResumeSessionLoad::Absent => {
return Err(MaterializeServeError::ResumeSessionNotFound {
session_id: session_id.clone(),
});
}
ResumeSessionLoad::ArchivedNotRevivable { runtime_state } => {
return Err(MaterializeServeError::ResumeSessionNonRecoverable {
session_id: session_id.clone(),
state: runtime_state.map_or_else(
|| "<no runtime record>".to_string(),
|state| state.to_string(),
),
});
}
}
}
};
let mut config = self
.compile_member_config(&decompiled, &session_id, resumed_session)
.await?;
let mut prepared = self
.substrate
.runtime_adapter
.prepare_local_session_materialization(session_id.clone())
.await
.map_err(|err| MaterializeServeError::Bindings {
detail: err.to_string(),
})?;
let actor_witness_slot = meerkat_session::LiveSessionActorWitnessSlot::default();
if let Err(error) =
crate::runtime::provisioner::install_prepared_mob_session_executor_handles(
Arc::clone(&self.substrate.session_service),
Arc::clone(&self.substrate.runtime_adapter),
&prepared,
actor_witness_slot.clone(),
)
.await
{
let rollback_error = prepared
.rollback_now_under_turn_finalization_boundary()
.await
.err();
return Err(MaterializeServeError::Bindings {
detail: rollback_error.map_or_else(
|| format!("failed to install prepared session handles: {error}"),
|rollback| {
format!(
"failed to install prepared session handles: {error}; exact rollback failed: {rollback}"
)
},
),
});
}
let bindings = prepared.bindings_clone();
let claim_handle = Arc::clone(bindings.session_claim_handle());
let live = self
.build_member_comms_runtime(&config, &session_id, claim_handle)
.await;
let live = match live {
Ok(live) => live,
Err(error) => {
return Err(self
.rollback_prepared_session_after_failure(&mut prepared, error)
.await);
}
};
let member_runtime = Arc::clone(&live.runtime);
if let Err(detail) = bindings.install_peer_comms_on(member_runtime.as_ref()) {
let original = MaterializeServeError::Comms { detail };
return Err(self
.rollback_prepared_session_after_failure(&mut prepared, original)
.await);
}
config.runtime_build_mode =
meerkat_core::runtime_epoch::RuntimeBuildMode::SessionOwned(bindings);
let supervisor_desc = supervisor.clone().into_trusted_peer_descriptor();
if let Err(error) = meerkat_runtime::comms_drain::bind_supervisor_for_materialized_session(
self.substrate.runtime_adapter.as_ref(),
&session_id,
member_runtime.as_ref() as &dyn CoreCommsRuntime,
&supervisor_desc,
epoch,
)
.await
{
let original = MaterializeServeError::SupervisorBind {
detail: error.to_string(),
};
return Err(self
.rollback_prepared_session_after_failure(&mut prepared, original)
.await);
}
self.mount_member_operator_tools(
&mut config,
&decompiled,
&session_id,
generation,
fence_token,
host_id,
host_binding_generation,
&supervisor_desc,
&member_runtime,
);
config.session_comms_runtime_override = Some(
meerkat::encode_session_comms_runtime_override_for_service(Arc::clone(&member_runtime)),
);
let req = crate::build::to_create_session_request(
&config,
meerkat_core::types::ContentInput::Text(String::new()),
);
let boundary = materialization_boundary.take().ok_or_else(|| {
MaterializeServeError::Bindings {
detail: format!(
"materialization for session {session_id} lost its exact service boundary before actor creation"
),
}
})?;
let actor_transaction = PreparedServiceActorTransaction::new(
session_id.clone(),
Arc::clone(&self.substrate.session_service),
prepared,
boundary,
actor_witness_slot,
)
.map_err(MaterializeServeError::Build)?;
let (created, actor_transaction) = actor_transaction
.create_owned(req, None)
.await
.map_err(|error| classify_materialize_create_error(error, spec))?;
if created.session_id != session_id {
let original = MaterializeServeError::IdentityMismatch {
expected: session_id.to_string(),
created: created.session_id.to_string(),
};
drop(config);
drop(member_runtime);
drop(live);
if let Err(cleanup) = actor_transaction.abort().await {
return Err(MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: original.to_string(),
cleanup: format!(
"exact actor/materialization transaction cleanup failed: {cleanup}"
),
});
}
return Err(original);
}
let pre_attach = async {
self.prepare_member_runtime_placement(
&session_id,
&decompiled.agent_identity,
generation,
fence_token,
)
.await?;
self.substrate
.runtime_adapter
.maybe_spawn_mob_comms_drain(
&session_id,
Arc::clone(&member_runtime) as Arc<dyn CoreCommsRuntime>,
meerkat_runtime::meerkat_machine::dsl::MobId::from(spec.mob_id.clone()),
)
.await
.map_err(|error| MaterializeServeError::Comms {
detail: format!("member peer-ingress drain spawn failed: {error}"),
})?;
self.read_resolved_auth_binding(&session_id).await
}
.await;
let resolved_auth_binding = match pre_attach {
Ok(resolved) => resolved,
Err(original) => {
if let Err(cleanup) = actor_transaction.abort().await {
return Err(MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: original.to_string(),
cleanup: format!(
"exact actor/materialization transaction cleanup failed: {cleanup}"
),
});
}
return Err(original);
}
};
let runtime_publication = match self.attach_member_executor(actor_transaction).await {
Ok(publication) => publication,
Err(original) => return Err(original),
};
let member_pubkey = member_runtime.public_key().to_pubkey_string();
let member_peer_id = member_runtime.public_key().to_peer_id().as_str();
let ack_keypair = Arc::clone(&live.ack_keypair);
self.live_runtimes.insert(session_id.clone(), live);
Ok(MaterializedBuildOutcome {
session_id,
residency_update,
runtime_publication: Some(runtime_publication),
member_runtime,
ack_keypair,
member_pubkey,
member_peer_id,
launch_outcome,
resolved_auth_binding,
generation_start_seq,
})
}
#[allow(clippy::too_many_arguments)] pub async fn revive_from_row(
&mut self,
spec: &PortableMemberSpec,
recorded_session_id: &str,
recorded_member_pubkey: &str,
generation: u64,
fence_token: u64,
host_id: &str,
host_binding_generation: u64,
supervisor: &TrustedPeerDescriptor,
epoch: u64,
) -> Result<RevivedMemberOutcome, MaterializeServeError> {
let decompiled = decompile_portable_spec(spec)?;
self.revive_decompiled_from_row(
spec,
decompiled,
recorded_session_id,
recorded_member_pubkey,
generation,
fence_token,
host_id,
host_binding_generation,
supervisor,
epoch,
)
.await
}
#[allow(clippy::too_many_arguments)] pub async fn revive_prepared_from_row(
&mut self,
spec: &PortableMemberSpec,
decompiled: DecompiledMemberBuild,
recorded_session_id: &str,
recorded_member_pubkey: &str,
generation: u64,
fence_token: u64,
host_id: &str,
host_binding_generation: u64,
supervisor: &TrustedPeerDescriptor,
epoch: u64,
) -> Result<RevivedMemberOutcome, MaterializeServeError> {
self.revive_decompiled_from_row(
spec,
decompiled,
recorded_session_id,
recorded_member_pubkey,
generation,
fence_token,
host_id,
host_binding_generation,
supervisor,
epoch,
)
.await
}
#[allow(clippy::too_many_arguments)] async fn revive_decompiled_from_row(
&mut self,
spec: &PortableMemberSpec,
decompiled: DecompiledMemberBuild,
recorded_session_id: &str,
recorded_member_pubkey: &str,
generation: u64,
fence_token: u64,
host_id: &str,
host_binding_generation: u64,
supervisor: &TrustedPeerDescriptor,
epoch: u64,
) -> Result<RevivedMemberOutcome, MaterializeServeError> {
let session_id = SessionId::parse(recorded_session_id).map_err(|err| {
MaterializeServeError::ResumeSessionIdInvalid {
session_id: recorded_session_id.to_string(),
detail: err.to_string(),
}
})?;
if self.session_live(&session_id).await {
if let Some(live) = self.live_runtimes.get(&session_id) {
if let Some(stamp) = self.upcall_binding_stamps.get(&session_id) {
stamp.update(
generation,
fence_token,
host_id,
host_binding_generation,
session_id.to_string(),
);
}
let expected =
meerkat_contracts::wire::supervisor_bridge::BridgeMemberIncarnation {
mob_id: spec.mob_id.clone(),
agent_identity: spec.agent_identity.clone(),
host_id: host_id.to_string(),
binding_generation: host_binding_generation,
member_session_id: session_id.to_string(),
generation,
fence_token,
};
if self
.substrate
.runtime_adapter
.member_incarnation(&session_id)
.as_ref()
== Some(&expected)
{
return Ok(RevivedMemberOutcome {
member: live.clone(),
residency_update: None,
runtime_publication: None,
});
}
}
}
if let Ok(
state @ (meerkat_runtime::RuntimeState::Retired
| meerkat_runtime::RuntimeState::Destroyed),
) = self
.substrate
.runtime_adapter
.runtime_state(&session_id)
.await
{
return Err(MaterializeServeError::ResumeSessionNonRecoverable {
session_id: recorded_session_id.to_string(),
state: state.to_string(),
});
}
let residency_update = self
.substrate
.runtime_adapter
.begin_member_residency_update(session_id.clone())
.await;
self.quiesce_nonserving_incarnation(&session_id).await?;
let mut materialization_boundary = Some(
self.acquire_vacant_materialization_boundary(&session_id)
.await?,
);
let loaded = self
.substrate
.session_service
.materialize_session_for_resume(&session_id)
.await?;
let session = match loaded {
ResumeSessionLoad::Active(session) => *session,
ResumeSessionLoad::Revivable(session) => {
if !session
.lifecycle_terminal()
.is_some_and(meerkat_core::SessionLifecycleTerminal::is_archived)
{
return Err(MaterializeServeError::ResumeSessionNonRecoverable {
session_id: recorded_session_id.to_string(),
state: meerkat_runtime::RuntimeState::Retired.to_string(),
});
}
*session
}
ResumeSessionLoad::ArchivedNotRevivable { runtime_state } => {
return Err(MaterializeServeError::ResumeSessionNonRecoverable {
session_id: recorded_session_id.to_string(),
state: runtime_state.map_or_else(
|| "<no runtime record>".to_string(),
|state| state.to_string(),
),
});
}
ResumeSessionLoad::Absent => {
if self
.substrate
.session_service
.session_known_to_archive_authority(&session_id)
.await?
{
return Err(MaterializeServeError::ResumeSessionNonRecoverable {
session_id: recorded_session_id.to_string(),
state: "Archived/Retired".to_string(),
});
}
return Err(MaterializeServeError::ResumeSessionNotFound {
session_id: recorded_session_id.to_string(),
});
}
};
let mut config = self
.compile_member_config(&decompiled, &session_id, Some(session))
.await?;
let mut prepared = match self
.substrate
.runtime_adapter
.prepare_local_session_materialization(session_id.clone())
.await
{
Ok(prepared) => prepared,
Err(error) => {
let original = MaterializeServeError::Bindings {
detail: error.to_string(),
};
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, original)
.await);
}
};
let actor_witness_slot = meerkat_session::LiveSessionActorWitnessSlot::default();
if let Err(error) =
crate::runtime::provisioner::install_prepared_mob_session_executor_handles(
Arc::clone(&self.substrate.session_service),
Arc::clone(&self.substrate.runtime_adapter),
&prepared,
actor_witness_slot.clone(),
)
.await
{
let rollback_error = prepared
.rollback_now_under_turn_finalization_boundary()
.await
.err();
let original = MaterializeServeError::Bindings {
detail: rollback_error.map_or_else(
|| format!("failed to install prepared session handles: {error}"),
|rollback| {
format!(
"failed to install prepared session handles: {error}; exact rollback failed: {rollback}"
)
},
),
};
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, original)
.await);
}
let bindings = prepared.bindings_clone();
let claim_handle = Arc::clone(bindings.session_claim_handle());
let live = match self
.build_member_comms_runtime(&config, &session_id, claim_handle)
.await
{
Ok(live) => live,
Err(error) => {
let error = self
.rollback_prepared_session_after_failure(&mut prepared, error)
.await;
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, error)
.await);
}
};
let member_runtime = Arc::clone(&live.runtime);
let derived = member_runtime.public_key().to_pubkey_string();
if derived != recorded_member_pubkey {
let original = MaterializeServeError::RevivedIdentityDiverged {
recorded: recorded_member_pubkey.to_string(),
derived,
};
drop(member_runtime);
drop(live);
let original = self
.rollback_prepared_session_after_failure(&mut prepared, original)
.await;
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, original)
.await);
}
if let Err(detail) = bindings.install_peer_comms_on(member_runtime.as_ref()) {
drop(member_runtime);
drop(live);
let original = self
.rollback_prepared_session_after_failure(
&mut prepared,
MaterializeServeError::Comms { detail },
)
.await;
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, original)
.await);
}
config.runtime_build_mode =
meerkat_core::runtime_epoch::RuntimeBuildMode::SessionOwned(bindings);
if let Err(error) = meerkat_runtime::comms_drain::bind_supervisor_for_materialized_session(
self.substrate.runtime_adapter.as_ref(),
&session_id,
member_runtime.as_ref() as &dyn CoreCommsRuntime,
supervisor,
epoch,
)
.await
{
drop(config);
drop(member_runtime);
drop(live);
let original = self
.rollback_prepared_session_after_failure(
&mut prepared,
MaterializeServeError::SupervisorBind {
detail: error.to_string(),
},
)
.await;
drop(materialization_boundary.take());
return Err(self
.preserve_snapshot_after_revival_failure(&session_id, original)
.await);
}
self.mount_member_operator_tools(
&mut config,
&decompiled,
&session_id,
generation,
fence_token,
host_id,
host_binding_generation,
supervisor,
&member_runtime,
);
config.session_comms_runtime_override = Some(
meerkat::encode_session_comms_runtime_override_for_service(Arc::clone(&member_runtime)),
);
let req = crate::build::to_create_session_request(
&config,
meerkat_core::types::ContentInput::Text(String::new()),
);
let boundary = materialization_boundary.take().ok_or_else(|| {
MaterializeServeError::Bindings {
detail: format!(
"revival for session {session_id} lost its exact service boundary before actor creation"
),
}
})?;
let actor_transaction = PreparedServiceActorTransaction::new(
session_id.clone(),
Arc::clone(&self.substrate.session_service),
prepared,
boundary,
actor_witness_slot,
)
.map_err(MaterializeServeError::Build)?;
let (created, actor_transaction) = actor_transaction
.create_owned(req, None)
.await
.map_err(|error| classify_materialize_create_error(error, spec))?;
if created.session_id != session_id {
let original = MaterializeServeError::IdentityMismatch {
expected: session_id.to_string(),
created: created.session_id.to_string(),
};
drop(config);
drop(member_runtime);
drop(live);
if let Err(cleanup) = actor_transaction.abort().await {
return Err(MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: original.to_string(),
cleanup: format!(
"exact actor/materialization transaction cleanup failed: {cleanup}"
),
});
}
return Err(original);
}
let pre_attach = async {
self.prepare_member_runtime_placement(
&session_id,
&decompiled.agent_identity,
generation,
fence_token,
)
.await?;
self.substrate
.runtime_adapter
.maybe_spawn_mob_comms_drain(
&session_id,
Arc::clone(&member_runtime) as Arc<dyn CoreCommsRuntime>,
meerkat_runtime::meerkat_machine::dsl::MobId::from(spec.mob_id.clone()),
)
.await
.map_err(|error| MaterializeServeError::Comms {
detail: format!("revived member peer-ingress drain spawn failed: {error}"),
})?;
Ok::<(), MaterializeServeError>(())
}
.await;
if let Err(error) = pre_attach {
drop(config);
drop(member_runtime);
drop(live);
if let Err(cleanup) = actor_transaction.abort().await {
return Err(MaterializeServeError::UnrecordedSessionCleanup {
session_id: session_id.to_string(),
original: error.to_string(),
cleanup: format!(
"exact actor/materialization transaction cleanup failed: {cleanup}"
),
});
}
return Err(error);
}
let mut runtime_publication = match self.attach_member_executor(actor_transaction).await {
Ok(publication) => publication,
Err(original) => return Err(original),
};
let abandoned_predecessor_inputs = runtime_publication
.abandon_recovered_predecessor_inputs()
.await
.map_err(|error| MaterializeServeError::ExecutorAttach {
detail: error.to_string(),
})?;
if abandoned_predecessor_inputs != 0 {
tracing::info!(
%session_id,
abandoned_predecessor_inputs,
"mob host revival abandoned requests owned by the interrupted attachment"
);
}
self.live_runtimes.insert(session_id, live.clone());
Ok(RevivedMemberOutcome {
member: live,
residency_update: Some(residency_update),
runtime_publication: Some(runtime_publication),
})
}
async fn compile_member_config(
&self,
decompiled: &DecompiledMemberBuild,
session_id: &SessionId,
resumed_session: Option<meerkat_core::Session>,
) -> Result<meerkat::AgentBuildConfig, MaterializeServeError> {
let base = BuildAgentConfigParams {
mob_id: &decompiled.mob_id,
profile_name: &decompiled.profile_name,
agent_identity: &decompiled.agent_identity,
profile: &decompiled.profile,
definition: &decompiled.definition,
external_tools: None,
context: decompiled.context.clone(),
labels: decompiled.labels.clone(),
additional_instructions: decompiled.additional_instructions.clone(),
shell_env: None,
mob_tool_authority_context: decompiled.mob_tool_authority_context.clone(),
inherited_tool_filter: None,
tool_access_policy: decompiled.tool_access_policy.clone(),
system_prompt_override: match &decompiled.system_prompt {
DecompiledSystemPrompt::Replace(text) => {
Some(SpawnSystemPromptOverride::Replace(text.clone()))
}
DecompiledSystemPrompt::Disable => Some(SpawnSystemPromptOverride::Disable),
},
};
let mut config = match resumed_session {
Some(session) => {
crate::build::build_resumed_agent_config(BuildResumedAgentConfigParams {
base,
expected_session_id: session_id,
resumed_session: session,
})
.await?
}
None => crate::build::build_agent_config(base).await?,
};
if matches!(decompiled.system_prompt, DecompiledSystemPrompt::Disable)
&& config
.resume_session
.as_ref()
.is_none_or(|s| s.messages().is_empty())
{
config.system_prompt = meerkat::SystemPromptOverride::Disable;
}
config.host_prompt_sections = meerkat_core::service::HostPromptSections::SpecPinned;
config.budget_limits = decompiled.budget_limits.clone();
config.keep_alive = true;
config.auth_binding = decompiled.auth_binding.clone();
Ok(config)
}
#[allow(clippy::too_many_arguments)] fn mount_member_operator_tools(
&mut self,
config: &mut meerkat::AgentBuildConfig,
decompiled: &DecompiledMemberBuild,
session_id: &SessionId,
generation: u64,
fence_token: u64,
host_id: &str,
host_binding_generation: u64,
supervisor: &TrustedPeerDescriptor,
member_runtime: &Arc<meerkat_comms::CommsRuntime>,
) {
if !decompiled.profile.tools.mob {
return;
}
let supervisor_route = meerkat_core::comms::PeerRoute::with_display_name(
supervisor.peer_id,
supervisor.name.clone(),
);
let binding_stamp = self
.upcall_binding_stamps
.entry(session_id.clone())
.or_insert_with(|| {
Arc::new(MemberUpcallBindingStamp::new(
generation,
fence_token,
host_id,
host_binding_generation,
session_id.to_string(),
))
});
binding_stamp.update(
generation,
fence_token,
host_id,
host_binding_generation,
session_id.to_string(),
);
let binding_stamp = Arc::clone(binding_stamp);
let dispatcher = crate::runtime::member_upcall::MemberUpcallToolDispatcher::new(
decompiled.agent_identity.clone(),
binding_stamp,
supervisor_route,
Arc::clone(member_runtime),
decompiled.requester_resolved_policy.clone(),
);
config.external_tools =
Some(Arc::new(dispatcher) as Arc<dyn meerkat_core::AgentToolDispatcher>);
}
async fn build_member_comms_runtime(
&self,
config: &meerkat::AgentBuildConfig,
session_id: &SessionId,
claim_handle: Arc<dyn meerkat_core::handles::SessionClaimHandle>,
) -> Result<LiveMemberRuntime, MaterializeServeError> {
let comms_name =
config
.comms_name
.as_deref()
.ok_or_else(|| MaterializeServeError::Comms {
detail: "compiled member config carries no comms name".to_string(),
})?;
let namespace = config
.realm_id
.as_ref()
.map(|realm| realm.as_str().to_string());
let silent_intents = Arc::new(
config
.silent_comms_intents
.iter()
.cloned()
.collect::<std::collections::HashSet<String>>(),
);
let runtime = meerkat_comms::CommsRuntime::inproc_only_session_scoped_with_silent_intents(
comms_name,
namespace,
self.substrate.member_identity_root.clone(),
session_id,
silent_intents,
claim_handle,
)
.await
.map_err(|err| MaterializeServeError::Comms {
detail: err.to_string(),
})?;
runtime.require_peer_comms_machine_authority();
if let Some(meta) = config.peer_meta.clone() {
runtime.set_peer_meta(meta);
}
let identity_dir = self
.substrate
.member_identity_root
.join(session_id.to_string());
let ack_keypair = meerkat_comms::Keypair::load_or_generate(&identity_dir)
.await
.map_err(|err| MaterializeServeError::Comms {
detail: format!("member ack keypair load failed: {err}"),
})?;
if ack_keypair.public_key() != runtime.public_key() {
return Err(MaterializeServeError::Comms {
detail: format!(
"member ack keypair for session '{session_id}' does not match the \
runtime identity (durable key material diverged)"
),
});
}
Ok(LiveMemberRuntime {
runtime: Arc::new(runtime),
ack_keypair: Arc::new(ack_keypair),
})
}
async fn read_resolved_auth_binding(
&self,
session_id: &SessionId,
) -> Result<Option<WireAuthBindingRef>, MaterializeServeError> {
let session = self
.substrate
.session_service
.load_persisted_session(session_id)
.await?
.ok_or_else(|| {
MaterializeServeError::SessionService(SessionError::NotFound {
id: session_id.clone(),
})
})?;
Ok(session
.session_metadata()
.and_then(|metadata| metadata.auth_binding)
.filter(|binding| {
matches!(
binding.origin,
meerkat_core::connection::BindingOrigin::Configured
)
})
.map(WireAuthBindingRef::from))
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use meerkat_contracts::wire::{
PortableDefinitionExtract, PortableProfile, PortableSpawnOverlay, PortableToolConfig,
WireMobToolAuthorityContext,
};
use std::collections::BTreeSet;
#[test]
fn resume_generation_floor_rejects_sequence_exhaustion_without_aliasing_max() {
assert_eq!(generation_start_seq_after(None).unwrap(), 1);
assert_eq!(
generation_start_seq_after(Some(u64::MAX - 1)).unwrap(),
u64::MAX
);
let error = generation_start_seq_after(Some(u64::MAX))
.expect_err("MAX has no strictly later generation floor");
assert!(matches!(error, MaterializeServeError::EventFloor { .. }));
assert!(error.to_string().contains("sequence space is exhausted"));
}
#[test]
fn unregister_cleanup_result_composition_preserves_every_failure() {
let session_id = SessionId::new();
assert!(
combine_unregister_cleanup_result(&session_id, "test cleanup", None, Ok(())).is_ok()
);
let fallback_only = combine_unregister_cleanup_result(
&session_id,
"test cleanup",
None,
Err(MaterializeServeError::Comms {
detail: "fallback-only failure".to_string(),
}),
)
.expect_err("a fallback failure must remain visible");
assert!(matches!(
fallback_only,
MaterializeServeError::Comms { ref detail }
if detail == "fallback-only failure"
));
let missing_runtime = meerkat_runtime::RuntimeDriverError::NotFound {
runtime_id: meerkat_runtime::identifiers::LogicalRuntimeId::new("missing-runtime"),
};
assert!(
combine_unregister_cleanup_result(
&session_id,
"test cleanup",
Some(missing_runtime.clone()),
Ok(()),
)
.is_ok(),
"NotFound is benign only after the fallback proves every carrier absent"
);
let missing_with_failed_fallback = combine_unregister_cleanup_result(
&session_id,
"test cleanup",
Some(missing_runtime),
Err(MaterializeServeError::Comms {
detail: "absence proof failed".to_string(),
}),
)
.expect_err("NotFound cannot hide a failed fallback absence proof");
let missing_with_failed_fallback = missing_with_failed_fallback.to_string();
assert!(missing_with_failed_fallback.contains("Runtime not found: missing-runtime"));
assert!(missing_with_failed_fallback.contains("absence proof failed"));
let unregister_only = combine_unregister_cleanup_result(
&session_id,
"test cleanup",
Some(meerkat_runtime::RuntimeDriverError::Internal(
"unregister-only failure".to_string(),
)),
Ok(()),
)
.expect_err("a non-NotFound unregister failure must remain visible");
let unregister_only = unregister_only.to_string();
assert!(unregister_only.contains("test cleanup"));
assert!(unregister_only.contains("unregister-only failure"));
let dual_failure = combine_unregister_cleanup_result(
&session_id,
"test cleanup",
Some(meerkat_runtime::RuntimeDriverError::Internal(
"primary unregister failure".to_string(),
)),
Err(MaterializeServeError::Comms {
detail: "secondary fallback failure".to_string(),
}),
)
.expect_err("both cleanup failures must remain visible");
let dual_failure = dual_failure.to_string();
assert!(dual_failure.contains("primary unregister failure"));
assert!(dual_failure.contains("secondary fallback failure"));
}
#[test]
fn prepared_session_unregister_failure_retains_original_error() {
let session_id = SessionId::new();
let combined = combine_prepared_session_unregister_result(
&session_id,
MaterializeServeError::Comms {
detail: "original comms failure".to_string(),
},
Err(meerkat_runtime::RuntimeDriverError::Internal(
"cleanup unregister failure".to_string(),
)),
);
match combined {
MaterializeServeError::UnrecordedSessionCleanup {
session_id: recorded_session_id,
original,
cleanup,
} => {
assert_eq!(recorded_session_id, session_id.to_string());
assert!(original.contains("original comms failure"));
assert!(cleanup.contains("cleanup unregister failure"));
}
other => panic!("expected composed unrecorded-session cleanup, got {other:?}"),
}
}
#[test]
fn host_health_requires_runtime_executor_and_admissibly_live_machine() {
assert!(member_runtime_is_healthy(true, true, true, true));
assert!(!member_runtime_is_healthy(false, true, true, true));
assert!(!member_runtime_is_healthy(true, false, true, true));
assert!(!member_runtime_is_healthy(true, true, false, true));
assert!(!member_runtime_is_healthy(true, true, true, false));
}
fn sample_spec() -> PortableMemberSpec {
PortableMemberSpec {
mob_id: "mob-1".to_string(),
profile_name: "worker".to_string(),
agent_identity: "worker-1".to_string(),
profile: PortableProfile {
model: "claude-opus-4-8".to_string(),
provider: meerkat_core::Provider::Anthropic,
self_hosted_server_id: None,
image_generation_provider: Some(meerkat_core::Provider::Gemini),
auto_compact_threshold: std::num::NonZeroU64::new(60_000),
resume_overrides: vec![meerkat_contracts::wire::WireMobResumeOverrideField::Model],
skills: vec!["review".to_string()],
tools: PortableToolConfig {
builtins: true,
shell: false,
comms: true,
memory: false,
workgraph: false,
mob: true,
schedule: false,
image_generation: false,
mcp_servers: BTreeMap::new(),
non_portable_disabled: Vec::new(),
},
peer_description: "reviews things".to_string(),
external_addressable: false,
runtime_mode: meerkat_contracts::wire::WireMobRuntimeMode::TurnDriven,
max_inline_peer_notifications: Some(3),
output_schema: None,
provider_params: None,
},
definition_extract: PortableDefinitionExtract {
models: BTreeMap::new(),
image_generation_provider: None,
skills: BTreeMap::from([(
"review".to_string(),
PortableSkillSource::Inline {
content: "# Review skill".to_string(),
},
)]),
profile_names: vec!["worker".to_string()],
},
overlay: PortableSpawnOverlay {
context: None,
labels: Some(BTreeMap::from([("team".to_string(), "alpha".to_string())])),
additional_instructions: Some(vec!["Be terse.".to_string()]),
system_prompt: PortableSystemPrompt::Set {
text: "You are worker-1.".to_string(),
},
tool_access_policy: Some(WireResolvedToolAccessPolicy::AllowList(vec![
"member_status".to_string(),
])),
mob_tool_authority_context: Some(WireMobToolAuthorityContext {
principal_token: "principal-token-1".to_string(),
can_create_mobs: false,
can_mutate_profiles: false,
can_run_adaptive_packs: false,
managed_mob_scope: BTreeSet::from(["mob-1".to_string()]),
spawn_profile_scope: BTreeMap::from([(
"mob-1".to_string(),
BTreeSet::from(["worker".to_string()]),
)]),
caller_provenance: None,
audit_invocation_id: Some("audit-1".to_string()),
}),
auth_binding: None,
budget_limits: Some(meerkat_core::BudgetLimits {
max_tokens: Some(100_000),
max_duration: None,
max_tool_calls: Some(50),
}),
runtime_mode: meerkat_contracts::wire::WireMobRuntimeMode::TurnDriven,
continuity_intent: WireSpawnContinuityIntent::Ephemeral,
},
required_env_keys: Vec::new(),
}
}
#[test]
fn post_preflight_llm_failure_is_typed_and_wire_diagnostic_is_sanitized() {
let spec = sample_spec();
let secret_stderr = "sk-secret-from-command-stderr";
let error = classify_materialize_create_error(
crate::error::MobError::SessionError(SessionError::build_llm_identity_unresolvable(
format!("command failed: {secret_stderr}"),
)),
&spec,
);
assert!(matches!(
&error,
MaterializeServeError::AuthBindingUnresolvable { .. }
));
let (cause, reason) = error.wire_cause();
assert_eq!(
cause,
meerkat_contracts::wire::supervisor_bridge::BridgeRejectionCause::AuthBindingUnresolvable {
realm: "env_default".to_string(),
binding: "anthropic".to_string(),
}
);
assert!(
!reason.contains(secret_stderr),
"raw provider/command diagnostics must remain local"
);
}
#[test]
fn post_preflight_provider_auth_failure_preserves_typed_wire_cause() {
let spec = sample_spec();
let failure = SessionProviderAuthFailure {
kind: meerkat_core::AuthErrorKind::InteractiveLoginRequired,
provider: meerkat_core::Provider::Anthropic,
realm_id: Some(meerkat_core::RealmId::parse("global").expect("realm")),
binding_id: Some(meerkat_core::BindingId::parse("anthropic").expect("binding")),
};
let error = classify_materialize_create_error(
crate::error::MobError::SessionError(SessionError::provider_auth_failure(
failure.clone(),
)),
&spec,
);
assert!(matches!(
&error,
MaterializeServeError::ProviderAuth(actual) if actual == &failure
));
let (cause, reason) = error.wire_cause();
assert_eq!(
cause,
meerkat_contracts::wire::supervisor_bridge::BridgeRejectionCause::MaterializeBuildRejected {
cause: meerkat_contracts::wire::supervisor_bridge::MemberBuildRejection::from_provider_auth_failure(
&failure,
),
}
);
assert!(reason.contains("interactive_login_required"));
assert!(reason.contains("global/anthropic"));
assert!(!reason.contains("credential"));
}
#[test]
fn decompile_maps_every_portable_field() {
let spec = sample_spec();
let decompiled = decompile_portable_spec(&spec).expect("decompile");
assert_eq!(decompiled.mob_id.as_str(), spec.mob_id);
assert_eq!(decompiled.profile_name.as_str(), spec.profile_name);
assert_eq!(decompiled.agent_identity.as_str(), spec.agent_identity);
assert_eq!(decompiled.profile.model, spec.profile.model);
assert_eq!(decompiled.profile.provider, Some(spec.profile.provider));
assert_eq!(
decompiled.profile.image_generation_provider,
spec.profile.image_generation_provider
);
assert_eq!(
decompiled.profile.auto_compact_threshold,
spec.profile.auto_compact_threshold
);
assert_eq!(decompiled.profile.skills, spec.profile.skills);
assert_eq!(
decompiled.profile.tools.builtins,
spec.profile.tools.builtins
);
assert_eq!(decompiled.profile.tools.mob, spec.profile.tools.mob);
assert!(decompiled.profile.tools.mcp.is_empty());
assert!(decompiled.profile.tools.rust_bundles.is_empty());
assert_eq!(
decompiled.profile.peer_description,
spec.profile.peer_description
);
assert_eq!(decompiled.profile.backend, None);
assert_eq!(
decompiled.profile.max_inline_peer_notifications,
spec.profile.max_inline_peer_notifications
);
assert_eq!(
decompiled.definition.models.len(),
spec.definition_extract.models.len()
);
assert_eq!(
decompiled.definition.skills.len(),
spec.definition_extract.skills.len()
);
for name in &spec.definition_extract.profile_names {
assert!(
decompiled
.definition
.profiles
.contains_key(&crate::ids::ProfileName::from(name.as_str()))
);
}
assert_eq!(decompiled.labels, spec.overlay.labels);
assert_eq!(
decompiled.additional_instructions,
spec.overlay.additional_instructions
);
assert_eq!(decompiled.budget_limits, spec.overlay.budget_limits);
assert!(decompiled.mob_tool_authority_context.is_some());
assert_eq!(decompiled.auth_binding, None);
}
#[test]
fn decompiled_tool_access_policy_is_never_inherit() {
let spec = sample_spec();
let decompiled = decompile_portable_spec(&spec).expect("decompile");
match decompiled.tool_access_policy {
Some(meerkat_core::ops::ToolAccessPolicy::AllowList(names)) => {
assert_eq!(names.len(), 1);
assert!(names.contains("member_status"));
}
other => panic!("expected AllowList, got {other:?}"),
}
}
#[test]
fn decompiled_system_prompt_maps_set_and_disable() {
let mut spec = sample_spec();
let decompiled = decompile_portable_spec(&spec).expect("decompile");
assert!(matches!(
decompiled.system_prompt,
DecompiledSystemPrompt::Replace(ref text) if text == "You are worker-1."
));
spec.overlay.system_prompt = PortableSystemPrompt::Disable;
let decompiled = decompile_portable_spec(&spec).expect("decompile");
assert!(matches!(
decompiled.system_prompt,
DecompiledSystemPrompt::Disable
));
}
#[test]
fn decompile_http_header_names_fail_closed() {
let mut spec = sample_spec();
spec.profile.tools.mcp_servers.insert(
"remote".to_string(),
PortableMcpDecl::Http {
url: "https://mcp.example".to_string(),
http_transport: None,
required_header_names: vec!["authorization".to_string()],
connect_timeout_secs: None,
},
);
assert!(matches!(
decompile_portable_spec(&spec),
Err(MaterializeDecompileError::McpHeaderNamesUnsupported { .. })
));
}
#[test]
fn decompile_preserves_sse_transport_and_nondefault_timeout() {
let mut spec = sample_spec();
spec.profile.tools.mcp_servers.insert(
"remote".to_string(),
PortableMcpDecl::Http {
url: "https://mcp.example/events".to_string(),
http_transport: Some(meerkat_core::mcp_config::McpHttpTransport::Sse),
required_header_names: Vec::new(),
connect_timeout_secs: Some(23),
},
);
let decompiled = decompile_portable_spec(&spec).expect("SSE declaration decompiles");
let server = decompiled
.profile
.tools
.mcp_servers
.iter()
.find(|server| server.name == "remote")
.expect("decompiled remote MCP server");
assert_eq!(
server.transport_kind(),
meerkat_core::mcp_config::McpTransportKind::Sse
);
assert_eq!(server.connect_timeout_secs, Some(23));
}
#[test]
fn stdio_command_path_walk_finds_absolute_and_path_entries() {
assert!(!stdio_command_discoverable("/definitely/not/a/real/binary"));
assert!(!stdio_command_discoverable("definitely-not-a-real-binary"));
}
}