use std::sync::Arc;
use async_trait::async_trait;
#[cfg(feature = "comms")]
use super::configure_peer_ingress;
use super::{
AcceptedScheduledInput, NoopScheduleMobHost, ScheduledPromptDispatch,
SharedScheduleTargetAdapter, SurfaceScheduleMobHost, SurfaceScheduleSessionHost,
build_dispatch_from_accepted, default_persistent_executor, immediate_delivery_failure,
materialize_session, schedule_attempt_idempotency_key, schedule_host_supported,
spawn_schedule_host,
};
use crate::{
Config, CreateSessionRequest, FactoryAgentBuilder, PersistentSessionService,
ScheduleDomainError, ScheduleService, Session, SessionMaterializationSpec, SessionService,
SessionTargetBinding, TargetProbeOutcome,
};
use meerkat_core::service::{DeferredPromptPolicy, InitialTurnPolicy, SessionBuildOptions};
use meerkat_core::types::{ContentInput, SessionId};
use meerkat_runtime::MeerkatMachine;
use meerkat_schedule::DeliveryFailureReason;
pub fn spawn_runtime_backed_schedule_host(
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
runtime_adapter: Arc<MeerkatMachine>,
config: Config,
schedule_service: ScheduleService,
build_template: SessionBuildOptions,
owner_id: impl Into<String>,
) -> Option<super::ScheduleHostHandle> {
let mob_host: Arc<dyn SurfaceScheduleMobHost> = Arc::new(NoopScheduleMobHost::new(
"scheduled mob targets are not supported by this runtime host",
));
spawn_runtime_backed_schedule_host_with_mobs(
service,
runtime_adapter,
config,
schedule_service,
build_template,
mob_host,
owner_id,
)
}
pub fn spawn_runtime_backed_schedule_host_with_mobs(
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
runtime_adapter: Arc<MeerkatMachine>,
_config: Config,
schedule_service: ScheduleService,
build_template: SessionBuildOptions,
mob_host: Arc<dyn SurfaceScheduleMobHost>,
owner_id: impl Into<String>,
) -> Option<super::ScheduleHostHandle> {
if !schedule_host_supported(schedule_service.store().kind()) {
return None;
}
let session_host: Arc<dyn SurfaceScheduleSessionHost> = Arc::new(
RuntimeBackedScheduleSessionHost::new(service, runtime_adapter, build_template),
);
let adapter = Arc::new(SharedScheduleTargetAdapter::new(
schedule_service.clone(),
session_host,
mob_host,
));
Some(spawn_schedule_host(schedule_service, adapter, owner_id))
}
struct RuntimeBackedScheduleSessionHost {
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
runtime_adapter: Arc<MeerkatMachine>,
build_template: SessionBuildOptions,
}
fn materialized_build_options(
template: &SessionBuildOptions,
create: &SessionMaterializationSpec,
) -> SessionBuildOptions {
let mut build = template.clone();
build.provider = create.provider;
build.output_schema = create.output_schema.clone();
build.structured_output_retries = create.structured_output_retries;
build.comms_name = create.comms_name.clone();
build.peer_meta = create.peer_meta.clone();
build.provider_params = create.provider_params.clone();
if !create.preload_skills.is_empty() {
build.preload_skills = Some(create.preload_skills.clone());
}
build.realm_id = create.realm_id.clone();
build.instance_id = create.instance_id.clone();
build.backend = create.backend.clone();
build.config_generation = create.config_generation;
build.keep_alive = create.keep_alive;
build.app_context = create.app_context.clone();
build.additional_instructions = (!create.additional_instructions.is_empty())
.then(|| create.additional_instructions.clone())
.or(build.additional_instructions);
build
}
fn accepted_scheduled_input_from_runtime_handle(
correlation_id: Option<String>,
handle: Option<meerkat_runtime::CompletionHandle>,
) -> AcceptedScheduledInput {
match handle {
Some(handle) => AcceptedScheduledInput::with_runtime_handle(correlation_id, handle),
None => AcceptedScheduledInput::with_authority_unavailable(
correlation_id,
"runtime completion handle missing after accepted dispatch",
),
}
}
fn runtime_delivery_dispatch(
occurrence: &crate::Occurrence,
outcome: meerkat_runtime::accept::AcceptOutcome,
handle: Option<meerkat_runtime::CompletionHandle>,
materialized_session_id: Option<SessionId>,
) -> Result<crate::DeliveryDispatch, ScheduleDomainError> {
match outcome {
meerkat_runtime::accept::AcceptOutcome::Accepted { input_id, .. } => {
let accepted =
accepted_scheduled_input_from_runtime_handle(Some(input_id.to_string()), handle);
Ok(build_dispatch_from_accepted(
occurrence,
accepted,
materialized_session_id,
))
}
meerkat_runtime::accept::AcceptOutcome::Deduplicated { existing_id, .. } => {
let accepted = match handle {
Some(handle) => AcceptedScheduledInput::with_runtime_handle(
Some(existing_id.to_string()),
handle,
),
None => AcceptedScheduledInput::with_authority_unavailable(
Some(existing_id.to_string()),
format!(
"runtime completion authority unavailable for terminal deduplicated input {existing_id}"
),
),
};
Ok(build_dispatch_from_accepted(
occurrence,
accepted,
materialized_session_id,
))
}
meerkat_runtime::accept::AcceptOutcome::Rejected { reason } => {
Ok(immediate_delivery_failure(
occurrence,
reason.to_string(),
DeliveryFailureReason::RuntimeRejected,
None,
materialized_session_id,
))
}
_ => Ok(immediate_delivery_failure(
occurrence,
"runtime returned an unknown admission outcome".to_string(),
DeliveryFailureReason::RuntimeRejected,
None,
materialized_session_id,
)),
}
}
impl RuntimeBackedScheduleSessionHost {
fn new(
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
runtime_adapter: Arc<MeerkatMachine>,
build_template: SessionBuildOptions,
) -> Self {
Self {
service,
runtime_adapter,
build_template,
}
}
async fn ensure_runtime_session_registered(
&self,
session_id: &SessionId,
) -> Result<(), ScheduleDomainError> {
self.ensure_session_target_exists(session_id).await?;
self.runtime_adapter
.ensure_session_with_executor(
session_id.clone(),
default_persistent_executor(
Arc::clone(&self.service),
Arc::clone(&self.runtime_adapter),
session_id.clone(),
),
)
.await
.map_err(schedule_internal)?;
self.update_peer_ingress_context(session_id).await?;
Ok(())
}
async fn update_peer_ingress_context(
&self,
session_id: &SessionId,
) -> Result<(), ScheduleDomainError> {
#[cfg(feature = "comms")]
{
let session = self
.service
.load_authoritative_session(session_id)
.await
.map_err(schedule_internal)?
.ok_or_else(|| {
ScheduleDomainError::InvalidSchedule(format!("session not found: {session_id}"))
})?;
let keep_alive = session
.session_metadata()
.ok_or_else(|| {
ScheduleDomainError::Internal(format!(
"session {session_id} is missing session metadata"
))
})?
.keep_alive;
configure_peer_ingress(&self.runtime_adapter, &self.service, session_id, keep_alive)
.await;
}
#[cfg(not(feature = "comms"))]
let _ = session_id;
Ok(())
}
async fn ensure_session_target_exists(
&self,
session_id: &SessionId,
) -> Result<(), ScheduleDomainError> {
match self.service.read(session_id).await {
Ok(_) => Ok(()),
Err(meerkat_core::service::SessionError::NotFound { .. }) => Err(
ScheduleDomainError::InvalidSchedule(format!("session not found: {session_id}")),
),
Err(error) => Err(ScheduleDomainError::Internal(format!(
"failed to read session target {session_id}: {error}"
))),
}
}
fn build_materialized_request(
&self,
create: &SessionMaterializationSpec,
prompt_system_prompt: Option<&str>,
) -> CreateSessionRequest {
let build = materialized_build_options(&self.build_template, create);
CreateSessionRequest {
model: create.model.clone(),
prompt: ContentInput::Text(String::new()),
render_metadata: None,
system_prompt: prompt_system_prompt
.map(str::to_owned)
.or_else(|| create.system_prompt.clone()),
max_tokens: create.max_tokens,
event_tx: None,
skill_references: None,
initial_turn: InitialTurnPolicy::Defer,
deferred_prompt_policy: DeferredPromptPolicy::Discard,
build: Some(build),
labels: Some(create.labels.clone()),
}
}
}
#[async_trait]
impl SurfaceScheduleSessionHost for RuntimeBackedScheduleSessionHost {
async fn probe_session_target(
&self,
binding: &SessionTargetBinding,
) -> Result<TargetProbeOutcome, ScheduleDomainError> {
let Some(session_id) = binding.resolved_session_id() else {
return Ok(TargetProbeOutcome::Ready);
};
match self.service.read(session_id).await {
Ok(view) if view.state.is_active => Ok(TargetProbeOutcome::Busy {
detail: Some(format!("session still running: {session_id}")),
}),
Ok(_) => Ok(TargetProbeOutcome::Ready),
Err(meerkat_core::service::SessionError::NotFound { .. }) => {
Ok(TargetProbeOutcome::Missing {
detail: Some(format!("session not found: {session_id}")),
})
}
Err(error) => Err(ScheduleDomainError::Internal(format!(
"failed to read session target {session_id}: {error}"
))),
}
}
async fn materialize_session(
&self,
create: &SessionMaterializationSpec,
prompt_system_prompt: Option<&str>,
) -> Result<SessionId, ScheduleDomainError> {
let request = self.build_materialized_request(create, prompt_system_prompt);
let keep_alive = request.build.as_ref().is_some_and(|build| build.keep_alive);
let result = Box::pin(materialize_session(
&self.service,
&self.runtime_adapter,
Session::new(),
request,
{
let service = Arc::clone(&self.service);
let runtime_adapter = Arc::clone(&self.runtime_adapter);
move |session_id| default_persistent_executor(service, runtime_adapter, session_id)
},
))
.await
.map_err(schedule_internal)?;
#[cfg(feature = "comms")]
configure_peer_ingress(
&self.runtime_adapter,
&self.service,
&result.session_id,
keep_alive,
)
.await;
#[cfg(not(feature = "comms"))]
let _ = keep_alive;
Ok(result.session_id)
}
async fn deliver_prompt(
&self,
session_id: &SessionId,
occurrence: &crate::Occurrence,
dispatch: ScheduledPromptDispatch,
) -> Result<crate::DeliveryDispatch, ScheduleDomainError> {
self.ensure_runtime_session_registered(session_id).await?;
let turn_metadata = Some(
meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata {
handling_mode: None,
keep_alive: None,
skill_references: (!dispatch.skill_refs.is_empty()).then(|| {
dispatch
.skill_refs
.iter()
.map(|skill_ref| skill_ref.key().clone())
.collect()
}),
flow_tool_overlay: None,
additional_instructions: (!dispatch.additional_instructions.is_empty()).then(
|| {
dispatch
.additional_instructions
.iter()
.map(|body| {
meerkat_core::lifecycle::run_primitive::TurnInstruction {
kind: meerkat_core::lifecycle::run_primitive::TurnInstructionKind::Host,
body: body.clone(),
}
})
.collect()
},
),
model: None,
provider: None,
provider_params: None,
render_metadata: dispatch.render_metadata.clone(),
execution_kind: None,
peer_response_terminal_apply_intent: None,
auth_binding: None,
},
);
let mut prompt_input =
meerkat_runtime::PromptInput::from_content_input(dispatch.prompt, turn_metadata);
prompt_input.header.source = meerkat_runtime::InputOrigin::System;
prompt_input.header.idempotency_key = Some(meerkat_runtime::IdempotencyKey::new(
schedule_attempt_idempotency_key(occurrence),
));
prompt_input.header.correlation_id = Some(meerkat_runtime::CorrelationId::from_uuid(
occurrence.occurrence_id.0,
));
let (outcome, handle) = self
.runtime_adapter
.accept_input_with_completion(session_id, meerkat_runtime::Input::Prompt(prompt_input))
.await
.map_err(|error| ScheduleDomainError::Internal(error.to_string()))?;
runtime_delivery_dispatch(
occurrence,
outcome,
handle,
dispatch.materialized_session_id,
)
}
async fn deliver_event(
&self,
session_id: &SessionId,
occurrence: &crate::Occurrence,
event_type: String,
payload: serde_json::Value,
render_metadata: Option<meerkat_core::types::RenderMetadata>,
materialized_session_id: Option<SessionId>,
) -> Result<crate::DeliveryDispatch, ScheduleDomainError> {
self.ensure_runtime_session_registered(session_id).await?;
let input = meerkat_runtime::Input::ExternalEvent(meerkat_runtime::ExternalEventInput {
header: meerkat_runtime::input::InputHeader {
id: meerkat_core::lifecycle::InputId::new(),
timestamp: chrono::Utc::now(),
source: meerkat_runtime::InputOrigin::External {
source_name: format!("schedule:{}", occurrence.schedule_id),
},
durability: meerkat_runtime::input::InputDurability::Durable,
visibility: meerkat_runtime::input::InputVisibility::default(),
idempotency_key: Some(meerkat_runtime::IdempotencyKey::new(
schedule_attempt_idempotency_key(occurrence),
)),
supersession_key: None,
correlation_id: Some(meerkat_runtime::CorrelationId::from_uuid(
occurrence.occurrence_id.0,
)),
},
event_type,
payload,
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
render_metadata,
});
let (outcome, handle) = self
.runtime_adapter
.accept_input_with_completion(session_id, input)
.await
.map_err(|error| ScheduleDomainError::Internal(error.to_string()))?;
runtime_delivery_dispatch(occurrence, outcome, handle, materialized_session_id)
}
}
fn schedule_internal(error: impl std::fmt::Display) -> ScheduleDomainError {
ScheduleDomainError::Internal(error.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
use meerkat_core::skills::{SkillKey, SkillName, SourceUuid};
fn fixture_skill_key(name: &str) -> SkillKey {
let skill_name = match SkillName::parse(name) {
Ok(skill_name) => skill_name,
Err(err) => unreachable!("static skill fixture is invalid: {err}"),
};
SkillKey::new(SourceUuid::builtin(), skill_name)
}
#[test]
fn materialized_build_options_forwards_preload_skill_keys() {
let key = fixture_skill_key("email");
let create = SessionMaterializationSpec {
model: "claude-sonnet-4-6".to_string(),
system_prompt: None,
max_tokens: None,
provider: None,
output_schema: None,
structured_output_retries: 0,
provider_params: None,
comms_name: None,
peer_meta: None,
labels: Default::default(),
preload_skills: vec![key.clone()],
additional_instructions: Vec::new(),
realm_id: None,
instance_id: None,
backend: None,
config_generation: None,
keep_alive: false,
app_context: None,
};
let build = materialized_build_options(&SessionBuildOptions::default(), &create);
assert_eq!(build.preload_skills, Some(vec![key]));
}
}