mod request_execution;
mod schedule_host;
#[cfg(not(target_arch = "wasm32"))]
mod stdio_json;
pub use meerkat_core::{
BUILD_ONLY_RECOVERY_OVERRIDE_ERROR, RecoveredSessionBuild, SurfaceSessionRecoveryContext,
SurfaceSessionRecoveryError, SurfaceSessionRecoveryOverrides, build_recovered_session,
has_build_only_turn_overrides, has_materialization_overrides,
session_allows_first_turn_build_overrides,
};
pub use request_execution::{
PreparedSurfaceSession, RequestAlreadyExists, RequestAsyncAction, RequestContext,
RequestTerminal, SurfaceRequestExecutor, noop_request_action, prepare_surface_session,
request_action,
};
pub use schedule_host::{
AcceptedScheduledInput, NoopScheduleMobHost, ScheduleHostHandle, ScheduledPromptDispatch,
SharedScheduleTargetAdapter, SurfaceScheduleMobHost, SurfaceScheduleSessionHost,
accepted_scheduled_input_from_runtime_outcome, async_completion_dispatch,
build_dispatch_from_accepted, immediate_completed_dispatch, immediate_delivery_failure,
schedule_attempt_idempotency_key, schedule_host_supported, spawn_schedule_host,
};
#[cfg(not(target_arch = "wasm32"))]
pub use stdio_json::{StdioJsonWriter, spawn_stdio_json_writer};
use meerkat_contracts::{
CapabilitiesResponse, CapabilityEntry, CapabilityStatus, ContractVersion, build_capabilities,
};
use meerkat_core::{AgentEvent, Config, PeerMeta};
#[cfg(not(target_arch = "wasm32"))]
use tokio::sync::mpsc;
#[cfg(target_arch = "wasm32")]
use tokio_with_wasm::alias::sync::mpsc;
#[cfg(feature = "skills")]
use meerkat_core::skills::{
SkillDocument, SkillError, SkillFilter, SkillId, SkillIntrospectionEntry, SkillRuntime,
};
#[cfg(feature = "mcp")]
use std::collections::HashMap;
#[cfg(feature = "mcp")]
use std::collections::VecDeque;
#[cfg(feature = "skills")]
use std::sync::Arc;
#[cfg(feature = "mcp")]
use std::sync::{Mutex, OnceLock};
pub fn build_capabilities_response(config: &Config) -> CapabilitiesResponse {
let registrations = build_capabilities();
let capabilities = registrations
.into_iter()
.map(|reg| {
let status = match reg.status_resolver {
Some(resolver) => resolver(config),
None => CapabilityStatus::Available,
};
CapabilityEntry {
id: reg.id,
description: reg.description.to_string(),
status,
}
})
.collect();
CapabilitiesResponse {
contract_version: ContractVersion::CURRENT,
capabilities,
}
}
pub fn build_models_catalog_response() -> meerkat_contracts::ModelsCatalogResponse {
let providers = meerkat_models::provider_defaults()
.iter()
.map(|pd| {
let models = pd
.models
.iter()
.map(|entry| {
let profile = meerkat_models::profile_for(entry.provider, entry.id).map(|p| {
meerkat_contracts::WireModelProfile {
model_family: p.model_family,
supports_temperature: p.supports_temperature,
supports_thinking: p.supports_thinking,
supports_reasoning: p.supports_reasoning,
inline_video: p.inline_video,
params_schema: p.params_schema,
}
});
meerkat_contracts::CatalogModelEntry {
id: entry.id.to_string(),
display_name: entry.display_name.to_string(),
tier: match entry.tier {
meerkat_models::ModelTier::Recommended => {
meerkat_contracts::WireModelTier::Recommended
}
meerkat_models::ModelTier::Supported => {
meerkat_contracts::WireModelTier::Supported
}
},
context_window: entry.context_window,
max_output_tokens: entry.max_output_tokens,
profile,
}
})
.collect();
meerkat_contracts::ProviderCatalog {
provider: pd.provider.to_string(),
default_model_id: pd.default_model_id.to_string(),
models,
}
})
.collect();
meerkat_contracts::ModelsCatalogResponse {
contract_version: ContractVersion::CURRENT,
providers,
}
}
pub fn resolve_keep_alive(requested: bool) -> Result<bool, String> {
#[cfg(feature = "comms")]
{
meerkat_comms::validate_keep_alive(requested)
}
#[cfg(not(feature = "comms"))]
{
if requested {
return Err(
"keep_alive requires comms support (build with --features comms)".to_string(),
);
}
Ok(false)
}
}
const RESERVED_MOB_PEER_META_LABELS: [&str; 3] = ["mob_id", "role", "meerkat_id"];
pub fn validate_public_peer_meta(peer_meta: Option<&PeerMeta>) -> Result<(), String> {
let Some(peer_meta) = peer_meta else {
return Ok(());
};
validate_raw_labels(Some(&peer_meta.labels))
}
pub fn validate_raw_labels(
labels: Option<&std::collections::BTreeMap<String, String>>,
) -> Result<(), String> {
let Some(labels) = labels else {
return Ok(());
};
for &label in &RESERVED_MOB_PEER_META_LABELS {
if labels.contains_key(label) {
return Err(format!(
"peer_meta label '{label}' is reserved for mob-managed sessions"
));
}
}
Ok(())
}
#[cfg(feature = "skills")]
pub async fn list_skills_introspection(
skill_runtime: &Option<Arc<SkillRuntime>>,
filter: &SkillFilter,
) -> Option<Result<Vec<SkillIntrospectionEntry>, SkillError>> {
let runtime = skill_runtime.as_ref()?;
Some(runtime.list_all_with_provenance(filter).await)
}
#[cfg(feature = "skills")]
pub async fn inspect_skill(
skill_runtime: &Option<Arc<SkillRuntime>>,
id: &SkillId,
source_name: Option<&str>,
) -> Option<Result<SkillDocument, SkillError>> {
let runtime = skill_runtime.as_ref()?;
Some(runtime.load_from_source(id, source_name).await)
}
pub fn spawn_event_forwarder<F>(callback: F) -> mpsc::Sender<AgentEvent>
where
F: Fn(AgentEvent) + Send + 'static,
{
let (tx, mut rx) = mpsc::channel::<AgentEvent>(256);
#[cfg(not(target_arch = "wasm32"))]
tokio::spawn(async move {
while let Some(event) = rx.recv().await {
callback(event);
}
});
#[cfg(target_arch = "wasm32")]
tokio_with_wasm::alias::task::spawn(async move {
while let Some(event) = rx.recv().await {
callback(event);
}
});
tx
}
#[cfg(feature = "mcp")]
pub fn mcp_live_response(
session_id: String,
operation: meerkat_contracts::McpLiveOperation,
server_name: Option<String>,
) -> meerkat_contracts::McpLiveOpResponse {
meerkat_contracts::McpLiveOpResponse {
session_id,
operation,
server_name,
status: meerkat_contracts::McpLiveOpStatus::Staged,
persisted: false,
applied_at_turn: None,
}
}
#[cfg(feature = "mcp")]
pub fn resolve_persisted(operation: &str, persisted: bool) -> bool {
if persisted {
tracing::warn!(
operation,
"caller sent persisted=true but config persistence is not yet implemented; \
responding with persisted=false"
);
}
false
}
#[cfg(feature = "mcp")]
pub async fn validate_reload_target(
adapter: &meerkat_mcp::McpRouterAdapter,
server_name: &str,
) -> Result<(), String> {
let active = adapter.active_server_names().await;
if active.iter().any(|n| n == server_name) {
Ok(())
} else {
Err(format!(
"MCP server '{server_name}' is not registered on this session"
))
}
}
#[cfg(feature = "mcp")]
pub async fn emit_mcp_lifecycle_events(
event_tx: &mpsc::Sender<meerkat_core::EventEnvelope<AgentEvent>>,
source_id: &str,
prompt: &mut String,
turn_number: u32,
actions: Vec<meerkat_mcp::McpLifecycleAction>,
) {
use meerkat_mcp::McpLifecyclePhase;
const MCP_SEQ_SOURCE_CAP: usize = 8192;
#[derive(Default)]
struct McpSeqState {
seq_by_source: HashMap<String, u64>,
source_order: VecDeque<String>,
}
static MCP_EVENT_SEQ_BY_SOURCE: OnceLock<Mutex<McpSeqState>> = OnceLock::new();
for action in actions {
let mut payload = action.to_tool_config_changed_payload();
payload.applied_at_turn = Some(turn_number);
let target = payload.target.clone();
let seq = {
let map = MCP_EVENT_SEQ_BY_SOURCE.get_or_init(|| Mutex::new(McpSeqState::default()));
let mut guard = match map.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
if !guard.seq_by_source.contains_key(source_id) {
let source_key = source_id.to_string();
guard.source_order.push_back(source_key.clone());
guard.seq_by_source.insert(source_key, 0);
while guard.seq_by_source.len() > MCP_SEQ_SOURCE_CAP {
if let Some(evicted) = guard.source_order.pop_front() {
guard.seq_by_source.remove(&evicted);
} else {
break;
}
}
}
let entry = guard
.seq_by_source
.entry(source_id.to_string())
.or_insert(0);
*entry += 1;
*entry
};
let _ = event_tx
.send(meerkat_core::EventEnvelope::new(
source_id,
seq,
None,
AgentEvent::ToolConfigChanged {
payload: payload.clone(),
},
))
.await;
if action.phase == McpLifecyclePhase::Forced {
if !prompt.is_empty() {
prompt.push('\n');
}
prompt.push_str(&format!(
"[system-notice] MCP server '{target}' removal forced after drain timeout."
));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn validate_public_peer_meta_rejects_reserved_labels() {
let peer_meta = PeerMeta::default().with_label("mob_id", "team");
let result = validate_public_peer_meta(Some(&peer_meta));
assert!(
result.is_err(),
"reserved mob labels must be rejected on public surfaces"
);
let Err(err) = result else {
unreachable!("asserted reserved labels are rejected above");
};
assert!(err.contains("mob-managed sessions"));
assert!(err.contains("mob_id"));
}
#[test]
fn validate_public_peer_meta_allows_unreserved_labels() {
let peer_meta = PeerMeta::default().with_label("team", "infra");
let result = validate_public_peer_meta(Some(&peer_meta));
assert!(
result.is_ok(),
"ordinary peer metadata should stay available"
);
}
}