#![allow(
clippy::expect_used,
clippy::unwrap_used,
clippy::panic,
clippy::uninlined_format_args,
clippy::collapsible_if,
clippy::redundant_clone,
clippy::needless_raw_string_hashes,
clippy::single_match,
clippy::redundant_closure_for_method_calls,
clippy::redundant_pattern_matching,
clippy::ignored_unit_patterns,
clippy::clone_on_copy,
clippy::manual_assert,
clippy::unwrap_in_result,
clippy::useless_vec
)]
const DEFAULT_TRACING_FILTER: &str = "warn,meerkat_mobkit=info,rpc_gateway=info";
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::io::Write;
use std::path::PathBuf;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use base64::Engine;
use meerkat_mobkit::unified_runtime::EventLogError;
use meerkat_mobkit::{
AuthPolicy, AuthProvider, Base64BlobStoreAdapter, BigQueryNaming, BinaryBlobStore,
ConsolePolicy, ConsoleUiConfig, DiscoverySpec, EventLogConfig, EventLogStore, EventQuery,
InMemoryMetadataStore, LocalJsonMemoryBackendConfig, MOBKIT_CONTRACT_VERSION,
MemoryBackendConfig, MobBootstrapOptions, MobBootstrapSpec, MobKitConfig, ModuleConfig,
ObjectStoreBlobStore, PersistedEvent, PersistentMetadataStore, PreSpawnData, ReleaseMetadata,
RestartPolicy, RuntimeDecisionState, RuntimeOpsPolicy, RuntimeOptions, RuntimeRoute,
STORAGE_RESOLUTION_CODE, ScheduleDefinition, SqliteConsoleLogStore, SqliteMetadataStore,
TrustedOidcRuntimeConfig, UnifiedRuntime, UnifiedRuntimeShutdownReport, handle_mobkit_rpc_json,
load_console_ui_config_from_path_for_realm,
mob_handle_runtime::{
ensure_shell_tooling_build_substrate, mob_definition_may_use_image_generation,
mob_definition_may_use_shell,
},
start_mobkit_runtime,
};
use sha2::{Digest, Sha256};
use async_trait::async_trait;
use meerkat::{
AgentEvent, AgentFactory, Config, CreateSessionRequest, EphemeralSessionService, FactoryAgent,
FactoryAgentBuilder, SessionAgentBuilder, SessionError,
};
use meerkat_core::ContentBlock;
use meerkat_core::error::{AgentError, ToolError};
use meerkat_core::ops::ToolDispatchOutcome;
use meerkat_core::types::{ToolCallView, ToolDef, ToolResult};
use meerkat_core::{
AgentToolDispatcher, ToolCatalogCapabilities, ToolCatalogEntry, ToolDeadlineContributor,
ToolDeadlineOwner, ToolExecutionContract, ToolExecutionMode,
};
use meerkat_mob::{MobDefinition, MobStorage};
use serde::Deserialize;
use serde_json::{Value, json};
use tokio::sync::{Mutex, mpsc, oneshot};
#[derive(Clone, Default)]
struct GatewayGatingConfig {
action_risk_tiers: HashMap<String, String>,
}
struct GatewayRuntimeOptions {
runtime_options: RuntimeOptions,
identity_bootstrap_mode: Option<meerkat_mobkit::IdentityBootstrapMode>,
max_sessions: usize,
routing_routes: Vec<RuntimeRoute>,
schedules: Vec<ScheduleDefinition>,
gating: GatewayGatingConfig,
event_log: Option<EventLogConfig>,
decisions: Option<RuntimeDecisionState>,
console_ui: ConsoleUiConfig,
console_require_app_auth: Option<bool>,
console_read_only: Option<bool>,
console_fetch_timeout_ms: Option<u64>,
access: Option<meerkat_mobkit::AccessController>,
demo_llm: bool,
member_comms_address: Option<String>,
agent_memory: Option<GatewayAgentMemoryOptions>,
workgraph: GatewayWorkgraphOption,
live: GatewayLiveOption,
host_runnables: Vec<meerkat::HostRunnableName>,
runtime_store_ephemeral: bool,
compaction: Option<meerkat_core::config::CompactionRuntimeConfig>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
enum GatewayLiveOption {
#[default]
Disabled,
Enabled {
public_base_url: Option<String>,
seed_max_chars: Option<usize>,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
enum GatewayWorkgraphOption {
#[default]
Enabled,
Disabled,
DurableDir(std::path::PathBuf),
}
struct GatewayAgentMemoryOptions {
config: meerkat_mobkit::AgentMemoryConfig,
path: std::path::PathBuf,
store: GatewayAgentMemoryStoreKind,
selector: Option<meerkat_mobkit::memory::selector::SelectorSpec>,
distiller: meerkat_mobkit::memory::distiller::DistillerConfig,
steward: meerkat_mobkit::memory::steward::StewardConfig,
hygienist: meerkat_mobkit::memory::hygienist::HygienistConfig,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
enum GatewayAgentMemoryStoreKind {
Markdown,
#[default]
Sqlite,
}
fn operator_scope_recall_inert(
agent_memory: Option<&GatewayAgentMemoryOptions>,
resolver_installed: bool,
) -> bool {
!resolver_installed
&& agent_memory.is_some_and(|memory| {
memory.config.operator_scope == meerkat_mobkit::AgentMemoryOperatorScope::Provisional
})
}
impl Default for GatewayRuntimeOptions {
fn default() -> Self {
Self {
runtime_options: RuntimeOptions::default(),
identity_bootstrap_mode: None,
max_sessions: 16,
routing_routes: Vec::new(),
schedules: Vec::new(),
gating: GatewayGatingConfig::default(),
event_log: None,
decisions: None,
console_ui: ConsoleUiConfig::default(),
console_require_app_auth: None,
console_read_only: None,
console_fetch_timeout_ms: None,
access: None,
demo_llm: false,
member_comms_address: None,
agent_memory: None,
workgraph: GatewayWorkgraphOption::Enabled,
live: GatewayLiveOption::Disabled,
host_runnables: Vec::new(),
runtime_store_ephemeral: false,
compaction: None,
}
}
}
fn gateway_agent_config(options: &GatewayRuntimeOptions) -> Config {
let mut config = Config::default();
if let Some(address) = options.member_comms_address.as_ref() {
config.comms.mode = meerkat_core::CommsRuntimeMode::Tcp;
config.comms.address = Some(address.clone());
}
if let Some(compaction) = options.compaction.as_ref() {
if let Err(error) = meerkat_mobkit::apply_compaction_policy(&mut config, compaction) {
tracing::error!(
%error,
"refusing an invalid runtime_options.compaction declaration; \
inheriting meerkat's default compaction policy instead"
);
}
}
config
}
#[derive(Default)]
struct InMemoryEventLogStore {
events: std::sync::Mutex<Vec<PersistedEvent>>,
}
impl EventLogStore for InMemoryEventLogStore {
fn append_batch(
&self,
events: Vec<PersistedEvent>,
) -> Pin<Box<dyn std::future::Future<Output = Result<(), EventLogError>> + Send + '_>> {
Box::pin(async move {
self.events.lock().unwrap().extend(events);
Ok(())
})
}
fn query(
&self,
query: EventQuery,
) -> Pin<
Box<
dyn std::future::Future<Output = Result<Vec<PersistedEvent>, EventLogError>>
+ Send
+ '_,
>,
> {
Box::pin(async move {
let mut events = self.events.lock().unwrap().clone();
events.retain(|event| {
query
.since_ms
.is_none_or(|since| event.timestamp_ms >= since)
&& query
.until_ms
.is_none_or(|until| event.timestamp_ms < until)
&& query
.after_seq
.is_none_or(|after_seq| event.seq > after_seq)
&& query
.member_id
.as_ref()
.is_none_or(|member_id| event.member_id.as_ref() == Some(member_id))
&& (query.event_types.is_empty()
|| query.event_types.iter().any(|event_type| {
matches!(
&event.event,
meerkat_mobkit::UnifiedEvent::Module(module)
if &module.event_type == event_type
)
}))
});
events.sort_by_key(|event| event.seq);
if let Some(limit) = query.limit {
events.truncate(limit);
}
Ok(events)
})
}
}
fn minimal_decision_state() -> RuntimeDecisionState {
RuntimeDecisionState {
bigquery: BigQueryNaming {
dataset: "default_dataset".to_string(),
table: "default_table".to_string(),
},
modules: vec![],
auth: AuthPolicy::default(),
trusted_oidc: TrustedOidcRuntimeConfig {
discovery_json: r#"{"issuer":"https://noop.example.com","authorization_endpoint":"https://noop.example.com/auth","token_endpoint":"https://noop.example.com/token","jwks_uri":"https://noop.example.com/.well-known/jwks.json","response_types_supported":["code"],"subject_types_supported":["public"],"id_token_signing_alg_values_supported":["RS256"]}"#.to_string(),
jwks_json: r#"{"keys":[]}"#.to_string(),
audience: "persistent-gateway".to_string(),
},
console: ConsolePolicy::default(),
ops: RuntimeOpsPolicy::default(),
release_metadata: ReleaseMetadata {
targets: vec![
"crates.io".to_string(),
"npm".to_string(),
"pypi".to_string(),
"github-releases".to_string(),
],
support_matrix: "lts".to_string(),
},
}
}
fn shell_module(id: &str, script: &str) -> ModuleConfig {
ModuleConfig {
id: id.to_string(),
command: "sh".to_string(),
args: vec!["-c".to_string(), script.to_string()],
restart_policy: RestartPolicy::Never,
}
}
const MODULE_BOUNDARY_ENV_KEY: &str = "MOBKIT_MODULE_BOUNDARY";
const MODULE_BOUNDARY_MCP: &str = "mcp";
#[derive(Debug, Deserialize)]
struct GatewayModuleConfig {
id: String,
command: String,
#[serde(default)]
args: Vec<String>,
#[serde(default = "gateway_restart_policy_never")]
restart_policy: RestartPolicy,
#[serde(default)]
env: BTreeMap<String, String>,
#[serde(default)]
boundary: Option<String>,
}
fn gateway_restart_policy_never() -> RestartPolicy {
RestartPolicy::Never
}
impl GatewayModuleConfig {
fn into_module_and_pre_spawn(self) -> (ModuleConfig, Option<PreSpawnData>) {
let GatewayModuleConfig {
id,
command,
args,
restart_policy,
mut env,
boundary,
} = self;
if boundary
.as_deref()
.is_some_and(|value| value.eq_ignore_ascii_case(MODULE_BOUNDARY_MCP))
{
env.insert(
MODULE_BOUNDARY_ENV_KEY.to_string(),
MODULE_BOUNDARY_MCP.to_string(),
);
}
let pre_spawn = if env.is_empty() {
None
} else {
Some(PreSpawnData {
module_id: id.clone(),
env: env.into_iter().collect(),
})
};
(
ModuleConfig {
id,
command,
args,
restart_policy,
},
pre_spawn,
)
}
}
fn parse_gateway_modules(params: &Value) -> (Vec<ModuleConfig>, Vec<PreSpawnData>) {
let gateway_modules: Vec<GatewayModuleConfig> = params
.get("modules")
.and_then(|value| serde_json::from_value(value.clone()).ok())
.unwrap_or_default();
let mut modules = Vec::with_capacity(gateway_modules.len());
let mut pre_spawn = Vec::new();
for gateway_module in gateway_modules {
let (module, maybe_pre_spawn) = gateway_module.into_module_and_pre_spawn();
modules.push(module);
if let Some(pre_spawn_data) = maybe_pre_spawn {
pre_spawn.push(pre_spawn_data);
}
}
(modules, pre_spawn)
}
#[cfg(test)]
mod tests {
use super::*;
use meerkat_mobkit::mob_handle_runtime::MobRuntimeError;
use meerkat_mobkit::unified_runtime::types::IdentityAuthorityReleaseOutcome;
use meerkat_mobkit::{RuntimeShutdownReport, ShutdownDrainReport};
#[test]
fn default_tracing_filter_surfaces_own_info_keeps_deps_at_warn() {
use tracing_subscriber::layer::SubscriberExt;
let filter = tracing_subscriber::EnvFilter::try_new(DEFAULT_TRACING_FILTER)
.expect("default filter must parse");
let subscriber = tracing_subscriber::registry().with(filter);
tracing::subscriber::with_default(subscriber, || {
assert!(
tracing::enabled!(
target: "meerkat_mobkit::identity_first::local_store",
tracing::Level::INFO
),
"the crate's own INFO lines (conversion progress) must pass the default filter"
);
assert!(
tracing::enabled!(target: "rpc_gateway", tracing::Level::INFO),
"the gateway binary's own INFO lines must pass the default filter"
);
assert!(
!tracing::enabled!(target: "meerkat_runtime::ops_lifecycle", tracing::Level::INFO),
"dependency INFO noise must stay filtered"
);
assert!(
tracing::enabled!(target: "meerkat_runtime::ops_lifecycle", tracing::Level::WARN),
"dependency warnings must still pass"
);
});
}
#[test]
fn gateway_module_boundary_becomes_pre_spawn_data() {
let params = json!({
"modules": [{
"id": "router",
"command": "python3",
"args": ["router.py"],
"restart_policy": "on_failure",
"boundary": "mcp",
"env": {
"ROUTER_FIXTURE": "homecore"
}
}]
});
let (modules, pre_spawn) = parse_gateway_modules(¶ms);
assert_eq!(
modules,
vec![ModuleConfig {
id: "router".to_string(),
command: "python3".to_string(),
args: vec!["router.py".to_string()],
restart_policy: RestartPolicy::OnFailure,
}]
);
assert_eq!(
pre_spawn,
vec![PreSpawnData {
module_id: "router".to_string(),
env: vec![
("MOBKIT_MODULE_BOUNDARY".to_string(), "mcp".to_string()),
("ROUTER_FIXTURE".to_string(), "homecore".to_string()),
],
}]
);
}
#[test]
fn gateway_module_without_env_does_not_create_pre_spawn_data() {
let params = json!({
"modules": [{
"id": "delivery",
"command": "python3"
}]
});
let (modules, pre_spawn) = parse_gateway_modules(¶ms);
assert_eq!(modules.len(), 1);
assert_eq!(modules[0].id, "delivery");
assert_eq!(modules[0].args, Vec::<String>::new());
assert_eq!(modules[0].restart_policy, RestartPolicy::Never);
assert!(pre_spawn.is_empty());
}
#[test]
fn gateway_runtime_options_parse_host_runnables() {
let params = json!({
"runtime_options": {
"host_runnables": ["digest", "backup.rotate"]
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options
.host_runnables
.iter()
.map(meerkat::HostRunnableName::as_str)
.collect::<Vec<_>>(),
vec!["digest", "backup.rotate"]
);
}
#[test]
fn gateway_runtime_options_reject_invalid_host_runnables() {
for (runtime_options, needle) in [
(json!({"host_runnables": "digest"}), "must be an array"),
(json!({"host_runnables": [7]}), "must be strings"),
(json!({"host_runnables": [" "]}), "is invalid"),
(
json!({"host_runnables": ["digest", "digest"]}),
"duplicated",
),
(
json!({"host_runnables": [
meerkat_mobkit::schedule_wiring::STEWARD_DREAM_RUNNABLE
]}),
"reserved",
),
] {
let params = json!({ "runtime_options": runtime_options });
let err = match parse_gateway_runtime_options(¶ms, None) {
Err(err) => err,
Ok(_) => panic!("expected rejection for {params}"),
};
assert!(err.contains(needle), "{err} should mention '{needle}'");
}
}
#[test]
fn gateway_runtime_options_reject_unknown_fields() {
let params = json!({ "runtime_options": { "blob_storage": {} } });
let err = match parse_gateway_runtime_options(¶ms, None) {
Err(err) => err,
Ok(_) => panic!("unknown runtime_options keys must be rejected"),
};
assert!(
err.contains("unsupported runtime_options fields: blob_storage"),
"{err}"
);
}
#[test]
fn gateway_runtime_options_parse_compaction_policy() {
let defaulted = parse_gateway_runtime_options(&json!({ "runtime_options": {} }), None)
.expect("runtime options");
assert!(
defaulted.compaction.is_none(),
"an undeclared policy must inherit, not invent a threshold"
);
assert!(
!gateway_agent_config(&defaulted)
.compaction
.auto_compact_threshold_explicit,
"the inheriting form must stay un-pinned"
);
let params = json!({
"runtime_options": {
"compaction": {
"auto_compact_threshold": 120_000,
"recent_turn_budget": 6,
}
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
let config = gateway_agent_config(&options);
assert_eq!(config.compaction.auto_compact_threshold, 120_000);
assert!(
config.compaction.auto_compact_threshold_explicit,
"a declared gateway threshold must be pinned against model-aware scaling"
);
assert_eq!(config.compaction.recent_turn_budget, 6);
}
#[test]
fn gateway_runtime_options_reject_invalid_compaction_policy() {
for (compaction, needle) in [
(json!({"auto_compact_treshold": 100}), "unsupported fields"),
(json!({"auto_compact_threshold": 0}), "greater than 0"),
(json!({"auto_compact_threshold": "lots"}), "is invalid"),
(json!(120_000), "must be a JSON object"),
] {
let params = json!({ "runtime_options": { "compaction": compaction } });
let err = match parse_gateway_runtime_options(¶ms, None) {
Err(err) => err,
Ok(_) => panic!("expected rejection for {params}"),
};
assert!(
err.contains("runtime_options.compaction"),
"{err} should name the offending path"
);
assert!(err.contains(needle), "{err} should mention '{needle}'");
}
}
#[test]
fn gateway_runtime_options_parse_runtime_store_declaration() {
let params = json!({ "runtime_options": { "runtime_store": { "storage": "memory" } } });
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert!(options.runtime_store_ephemeral);
let defaulted = parse_gateway_runtime_options(&json!({ "runtime_options": {} }), None)
.expect("runtime options");
assert!(!defaulted.runtime_store_ephemeral);
for (value, needle) in [
(json!({"storage": "sqlite"}), "unsupported"),
(json!({"storage": 7}), "must be 'memory'"),
(json!("memory"), "must be a JSON object"),
] {
let params = json!({ "runtime_options": { "runtime_store": value } });
let err = match parse_gateway_runtime_options(¶ms, None) {
Err(err) => err,
Ok(_) => panic!("expected rejection for {params}"),
};
assert!(err.contains(needle), "{err} should mention '{needle}'");
}
}
#[test]
fn gateway_runtime_options_event_log_storage_declarations() {
for storage in ["memory", "in_memory", "null"] {
let params = json!({ "runtime_options": { "event_log": { "storage": storage } } });
let options = parse_gateway_runtime_options(¶ms, None)
.unwrap_or_else(|err| panic!("storage '{storage}' must parse: {err}"));
assert!(options.event_log.is_some());
}
let err = match parse_gateway_runtime_options(
&json!({ "runtime_options": { "event_log": { "storage": "sqlite" } } }),
None,
) {
Err(err) => err,
Ok(_) => panic!("undeclared event_log storage must be rejected"),
};
assert!(
err.contains("unsupported runtime_options.event_log.storage"),
"{err}"
);
}
#[test]
fn gateway_runtime_options_parse_live_seed_max_chars() {
let params = json!({
"runtime_options": {
"live": { "seed_max_chars": 200000 }
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options.live,
GatewayLiveOption::Enabled {
public_base_url: None,
seed_max_chars: Some(200_000),
}
);
let err = match parse_gateway_runtime_options(
&json!({"runtime_options": {"live": {"seed_max_chars": 0}}}),
None,
) {
Err(err) => err,
Ok(_) => panic!("zero seed_max_chars must be rejected"),
};
assert!(err.contains("seed_max_chars"), "{err}");
}
fn test_callback_bridge() -> StdioCallbackBridge {
let (stdout_tx, _stdout_rx) = mpsc::channel::<String>(4);
StdioCallbackBridge::new(stdout_tx)
}
#[test]
fn callback_tool_dispatcher_defs_carry_wire_schemas() {
let schema = json!({
"type": "object",
"properties": { "city": { "type": "string" } },
"required": ["city"]
});
let specs = vec![
CallbackToolSpec::parse(&json!("legacy_name")).expect("legacy string spec"),
CallbackToolSpec::parse(&json!({
"name": "weather",
"description": "Look up the weather",
"input_schema": schema,
"execution": {
"mode": "detached",
"runner": {"name": "weather.scan", "version": "1"},
"restart_class": "non_resumable",
"idempotency_scope": "interaction_and_arguments",
"submission_timeout_ms": 30000,
"credential_scopes": ["weather.read"]
}
}))
.expect("object spec"),
];
let dispatcher =
CallbackToolDispatcher::new(test_callback_bridge(), "build-1".to_string(), specs, None);
let defs = AgentToolDispatcher::tools(&dispatcher);
assert_eq!(defs.len(), 2);
assert_eq!(defs[0].name.as_ref(), "legacy_name");
assert_eq!(defs[0].description, "Python callback tool");
assert_eq!(defs[0].input_schema, json!({"type": "object"}));
assert_eq!(defs[1].name.as_ref(), "weather");
assert_eq!(defs[1].description, "Look up the weather");
assert_eq!(defs[1].input_schema, schema);
let catalog = AgentToolDispatcher::tool_catalog(&dispatcher);
assert_eq!(
catalog[1].execution.default_mode(),
meerkat_core::ToolExecutionMode::Detached
);
let detached = catalog[1]
.execution
.detached_policy()
.expect("detached policy");
assert_eq!(detached.runner().name(), "weather.scan");
assert_eq!(detached.runner().version(), "1");
assert_eq!(
detached.restart_class(),
meerkat_core::RestartClass::NonResumable
);
assert_eq!(
detached.idempotency_scope(),
meerkat_core::IdempotencyScope::InteractionAndArguments
);
assert_eq!(
detached.credential_scopes(),
&std::collections::BTreeSet::from(["weather.read".to_string()])
);
}
#[test]
fn callback_tool_spec_rejects_malformed_wire_entries() {
for value in [
json!(7),
json!({"description": "no name"}),
json!({"name": ""}),
json!({"name": "x", "input_schema": "not-an-object"}),
json!({"name": "x", "description": 3}),
json!({"name": "x", "execution": {"mode": "detached"}}),
] {
assert!(
CallbackToolSpec::parse(&value).is_err(),
"expected rejection for {value}"
);
}
}
#[test]
fn callback_event_delivery_builds_stable_runtime_owned_ingress() {
let job_id = meerkat::JobId::new("019f74fb-1907-7b21-932d-ab22c4d1f500").expect("job id");
let session_id = meerkat_core::SessionId::parse("019f74fb-1907-7b21-932d-ab22c4d1f501")
.expect("session id");
let subscription = meerkat::JobSubscription::new(
meerkat::JobSubscriptionId::new("review-agent").expect("subscription"),
session_id,
meerkat::JobDeliveryKind::Event {
handling_mode: meerkat_core::HandlingMode::Queue,
},
);
let lineage =
meerkat::InteractionLineageId::from_string("019f74fb-1907-7b21-932d-ab22c4d1f502")
.expect("lineage");
let content = meerkat::JobDeliveryContent::Terminal(meerkat::JobTerminalResult::WorkerLost);
let input = callback_job_event_input(
&job_id,
7,
&subscription,
&lineage,
meerkat_core::HandlingMode::Queue,
&content,
);
let meerkat_runtime::Input::ExternalEvent(event) = input else {
panic!("job event delivery must use canonical external-event ingress");
};
assert_eq!(event.event_type, "job.terminal");
assert_eq!(event.handling_mode, meerkat_core::HandlingMode::Queue);
assert_eq!(
event
.header
.idempotency_key
.as_ref()
.map(ToString::to_string),
Some(format!("job:{job_id}:7:review-agent"))
);
assert_eq!(
event.payload,
json!({
"job_id": job_id.to_string(),
"delivery_sequence": 7,
"content": {
"kind": "terminal",
"result": "WorkerLost"
}
})
);
assert_eq!(
event.header.correlation_id.map(|id| id.to_string()),
Some(lineage.as_str().to_string())
);
}
#[test]
fn callback_runner_ownership_requires_exact_tool_runner_and_version() {
let binary: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
let blobs: Arc<dyn meerkat_core::BlobStore> = Arc::new(Base64BlobStoreAdapter::new(binary));
let runtime = DetachedCallbackJobRuntime::new(
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
Arc::new(meerkat::MemoryDetachedJobStore::new()),
blobs,
);
let dispatcher = CallbackToolDispatcher::new(
test_callback_bridge(),
"build-1".to_string(),
vec![
CallbackToolSpec::parse(&json!({
"name": "security_scan",
"execution": {
"mode": "detached",
"runner": {"name": "homecore.security_scan", "version": "1"},
"restart_class": "adoptable",
"idempotency_scope": "interaction_and_arguments",
"submission_timeout_ms": 30000
}
}))
.expect("tool spec"),
],
Some(runtime.clone()),
);
assert!(
dispatcher.reconcile_registered_catalog,
"the first exact callback owner must trigger recovery"
);
let duplicate = CallbackToolDispatcher::new(
test_callback_bridge(),
"build-2".to_string(),
vec![
CallbackToolSpec::parse(&json!({
"name": "security_scan",
"execution": {
"mode": "detached",
"runner": {"name": "homecore.security_scan", "version": "1"},
"restart_class": "adoptable",
"idempotency_scope": "interaction_and_arguments",
"submission_timeout_ms": 30000
}
}))
.expect("tool spec"),
],
Some(runtime.clone()),
);
assert!(
!duplicate.reconcile_registered_catalog,
"rebuilding an already-known exact owner must not start a duplicate recovery census"
);
let later_runner = CallbackToolDispatcher::new(
test_callback_bridge(),
"build-3".to_string(),
vec![
CallbackToolSpec::parse(&json!({
"name": "report_export",
"execution": {
"mode": "detached",
"runner": {"name": "homecore.report_export", "version": "1"},
"restart_class": "adoptable",
"idempotency_scope": "interaction_and_arguments",
"submission_timeout_ms": 30000
}
}))
.expect("tool spec"),
],
Some(runtime.clone()),
);
assert!(
later_runner.reconcile_registered_catalog,
"a later exact owner must get its own recovery census"
);
let make_spec = |tool: &str, runner: &str, version: &str| {
meerkat::JobSpec::new(
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
meerkat_core::SessionId::parse("019f74fb-1907-7b21-932d-ab22c4d1f503")
.expect("session"),
meerkat::ExecutionIntentId::from_string(format!("intent:{tool}")).expect("intent"),
meerkat::InteractionLineageId::from_string(format!("lineage:{tool}"))
.expect("lineage"),
meerkat::ToolIdentity::new(tool, version).expect("tool"),
meerkat::RunnerIdentity::new(runner, version).expect("runner"),
meerkat::RestartClass::Adoptable,
meerkat::CanonicalArgumentsHash::new(format!("sha256:{}", "a".repeat(64)))
.expect("hash"),
meerkat::JobSubmissionKey::new(format!("submission:{tool}:{runner}:{version}"))
.expect("submission"),
)
};
assert!(runtime.owns_callback_job(&make_spec(
"security_scan",
"homecore.security_scan",
"1"
)));
assert!(!runtime.owns_callback_job(&make_spec(
"other_tool",
"homecore.security_scan",
"1"
)));
assert!(!runtime.owns_callback_job(&make_spec(
"security_scan",
"homecore.security_scan",
"2"
)));
}
#[tokio::test]
async fn callback_reconcile_rehydrates_exact_committed_authority_without_advancing_fence() {
let store = Arc::new(meerkat::MemoryDetachedJobStore::new());
let service = meerkat::DetachedJobService::new(store.clone());
let session_id = meerkat_core::SessionId::parse("019f74fb-1907-7b21-932d-ab22c4d1f532")
.expect("session id");
let spec = meerkat::JobSpec::new(
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
session_id,
meerkat::ExecutionIntentId::from_string("intent:reconcile").expect("intent"),
meerkat::InteractionLineageId::from_string("lineage:reconcile").expect("lineage"),
meerkat::ToolIdentity::new("security_scan", "1").expect("tool"),
meerkat::RunnerIdentity::new("homecore.security_scan", "1").expect("runner"),
meerkat::RestartClass::Adoptable,
meerkat::CanonicalArgumentsHash::new(format!("sha256:{}", "a".repeat(64)))
.expect("arguments hash"),
meerkat::JobSubmissionKey::new("callback:reconcile").expect("submission key"),
);
let receipt = service.submit(spec).await.expect("submit");
let claim = service
.claim_attempt(
&receipt.job_id,
meerkat::AttemptClaim::new(
meerkat::WorkerId::new("worker-1").expect("worker"),
100,
u64::MAX,
meerkat::RunnerHandleRef::new("external:scan-1").expect("handle"),
),
)
.await
.expect("claim");
let (stdout_tx, mut stdout_rx) = mpsc::channel::<String>(4);
let bridge = StdioCallbackBridge::new(stdout_tx);
let binary: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
let blobs: Arc<dyn meerkat_core::BlobStore> = Arc::new(Base64BlobStoreAdapter::new(binary));
let dispatcher = CallbackToolDispatcher::new(
bridge.clone(),
"build-1".to_string(),
vec![
CallbackToolSpec::parse(&json!({
"name": "security_scan",
"execution": {
"mode": "detached",
"runner": {"name": "homecore.security_scan", "version": "1"},
"restart_class": "adoptable",
"idempotency_scope": "interaction_and_arguments",
"submission_timeout_ms": 30000
}
}))
.expect("tool spec"),
],
Some(DetachedCallbackJobRuntime::new(
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
store.clone(),
blobs,
)),
);
let reconcile = tokio::spawn(async move { dispatcher.reconcile_detached_jobs().await });
let request_line = stdout_rx.recv().await.expect("reconcile callback");
let request: Value = serde_json::from_str(&request_line).expect("callback json");
assert_eq!(
request.pointer("/params/attempts/0/authority/fence"),
Some(&json!(claim.fence.get()))
);
assert_eq!(
request.pointer("/params/attempts/0/authority/attempt_id"),
Some(&json!(claim.attempt_id.to_string()))
);
assert_eq!(
request.pointer("/params/attempts/0/runner_handle"),
Some(&json!("external:scan-1"))
);
bridge
.route_callback_response(json!({
"jsonrpc": "2.0",
"id": request["id"].clone(),
"result": {
"live_attempts": [{
"job_id": receipt.job_id.to_string(),
"attempt_id": claim.attempt_id.to_string(),
"fence": claim.fence.get()
}]
}
}))
.await;
reconcile
.await
.expect("reconcile task")
.expect("reconcile succeeds");
let reopened = meerkat::DetachedJobStore::get(&*store, &receipt.job_id)
.await
.expect("read")
.expect("job");
assert_eq!(reopened.machine_state.attempt_count, 1);
assert_eq!(reopened.machine_state.current_fence, claim.fence.get());
assert_eq!(
reopened.machine_state.current_attempt_id.as_deref(),
Some(claim.attempt_id.as_str())
);
}
#[test]
fn gateway_runtime_options_parse_implicit_delegate_retirement() {
let params = json!({
"runtime_options": {
"implicit_delegate_idle_retire_secs": 42,
"implicit_delegate_idle_sweep_interval_ms": 2500
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options.runtime_options.implicit_delegate_idle_retire_secs,
Some(42)
);
assert_eq!(
options
.runtime_options
.implicit_delegate_idle_sweep_interval_ms,
2500
);
}
#[test]
fn gateway_runtime_options_null_disables_implicit_delegate_retirement() {
let params = json!({
"runtime_options": {
"implicit_delegate_idle_retire_secs": null
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options.runtime_options.implicit_delegate_idle_retire_secs,
None
);
}
#[test]
fn gateway_runtime_options_parse_console_config_path() {
let tmp = tempfile::tempdir().expect("temp dir");
let path = tmp.path().join("console.toml");
std::fs::write(
&path,
r#"
[sidebar]
visible_controls = ["topology", "roster"]
[agent_list]
group_by = ["labels.console_group", "group"]
"#,
)
.expect("write console config");
let params = json!({
"runtime_options": {
"console_config_path": path.to_string_lossy()
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options.console_ui.sidebar.visible_controls,
Some(vec!["topology".to_string(), "roster".to_string()])
);
assert_eq!(
options.console_ui.agent_list.group_by,
vec!["labels.console_group".to_string(), "group".to_string()]
);
}
#[test]
fn gateway_runtime_options_parse_access_config_path() {
let tmp = tempfile::tempdir().expect("temp dir");
let path = tmp.path().join("access.toml");
std::fs::write(
&path,
r#"
enabled = true
admins = ["root@example.test"]
[[rules]]
id = "everyone-views"
actions = ["agent.view"]
"#,
)
.expect("write access config");
let params = json!({
"runtime_options": {
"access_config_path": path.to_string_lossy()
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
let access = options.access.expect("access controller");
assert!(access.enabled());
let (config, _) = access.snapshot();
assert_eq!(config.admins, vec!["root@example.test".to_string()]);
assert_eq!(config.rules.len(), 1);
}
#[test]
fn gateway_runtime_options_reject_invalid_access_config() {
let tmp = tempfile::tempdir().expect("temp dir");
let path = tmp.path().join("access.toml");
std::fs::write(&path, "enabled = true").expect("write access config");
let params = json!({
"runtime_options": {
"access_config_path": path.to_string_lossy()
}
});
let err = match parse_gateway_runtime_options(¶ms, None) {
Ok(_) => panic!("invalid access config must be rejected"),
Err(err) => err,
};
assert!(err.contains("access_config_path"), "{err}");
}
#[test]
fn gateway_runtime_options_can_disable_console_auth_for_local_console() {
let params = json!({
"runtime_options": {
"console_require_app_auth": false
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert!(
!options
.decisions
.expect("decisions")
.console
.require_app_auth
);
}
#[test]
fn gateway_runtime_options_can_enable_read_only_console() {
let params = json!({
"runtime_options": {
"console_read_only": true
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert!(options.decisions.expect("decisions").console.read_only);
}
#[test]
fn gateway_runtime_options_parse_console_fetch_timeout_ms() {
let params = json!({
"runtime_options": {
"console_fetch_timeout_ms": 120_000
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(
options
.decisions
.expect("decisions")
.console
.fetch_timeout_ms,
Some(120_000)
);
}
#[test]
fn gateway_runtime_options_reject_zero_console_fetch_timeout_ms() {
let params = json!({
"runtime_options": {
"console_fetch_timeout_ms": 0
}
});
let err = match parse_gateway_runtime_options(¶ms, None) {
Ok(_) => panic!("zero should fail"),
Err(err) => err,
};
assert!(err.contains("runtime_options.console_fetch_timeout_ms"));
}
#[test]
fn gateway_runtime_options_parse_demo_llm() {
let params = json!({
"runtime_options": {
"demo_llm": true
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert!(options.demo_llm);
}
#[test]
fn gateway_runtime_options_parse_member_comms_address() {
let params = json!({
"runtime_options": {
"member_comms_address": "127.0.0.1:0"
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(options.member_comms_address.as_deref(), Some("127.0.0.1:0"));
let config = gateway_agent_config(&options);
assert_eq!(config.comms.mode, meerkat_core::CommsRuntimeMode::Tcp);
assert_eq!(config.comms.address.as_deref(), Some("127.0.0.1:0"));
}
#[test]
fn gateway_runtime_options_parse_max_sessions() {
let params = json!({
"runtime_options": {
"max_sessions": 320
}
});
let options = parse_gateway_runtime_options(¶ms, None).expect("runtime options");
assert_eq!(options.max_sessions, 320);
}
#[test]
fn callback_build_agent_options_include_profile_name_from_spawn_labels() {
let build = meerkat_core::service::SessionBuildOptions {
peer_meta: Some(
meerkat_core::PeerMeta::default()
.with_label("role", "security")
.with_label("agent_identity", "domain:security"),
),
..Default::default()
};
let req = CreateSessionRequest {
model: "security-model".to_string(),
prompt: meerkat_core::ContentInput::Text("noop".to_string()),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
build: Some(build),
labels: None,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
injected_context: Vec::new(),
};
let options = callback_build_agent_options(&req, "build-test");
assert_eq!(options["scope_id"], "build-test");
assert_eq!(
options["profile_name"], "security",
"SDK build_agent must see the adopted roster profile instead of \
falling back to checkpoint-local defaults"
);
assert_eq!(options["labels"]["agent_identity"], "domain:security");
}
#[test]
fn callback_build_agent_options_merge_request_and_spawn_labels() {
let build = meerkat_core::service::SessionBuildOptions {
peer_meta: Some(
meerkat_core::PeerMeta::default()
.with_label("role", "security")
.with_label("agent_identity", "domain:security"),
),
..Default::default()
};
let mut request_labels = BTreeMap::new();
request_labels.insert("session_id".to_string(), "session-123".to_string());
let req = CreateSessionRequest {
model: "security-model".to_string(),
prompt: meerkat_core::ContentInput::Text("noop".to_string()),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
build: Some(build),
labels: Some(request_labels),
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
injected_context: Vec::new(),
};
let options = callback_build_agent_options(&req, "build-test");
assert_eq!(options["profile_name"], "security");
assert_eq!(options["session_id"], "session-123");
assert_eq!(options["labels"]["session_id"], "session-123");
assert_eq!(options["labels"]["agent_identity"], "domain:security");
}
#[test]
fn gateway_runtime_options_parse_agent_memory_defaults() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": true
}
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.config.realm, "default");
assert_eq!(
agent_memory.config.selection,
meerkat_mobkit::AgentMemorySelection::Contextual
);
assert_eq!(agent_memory.config.recall_timeout_ms, 500);
assert_eq!(
agent_memory.config.recall_failure_policy,
meerkat_mobkit::AgentMemoryRecallFailurePolicy::Skip
);
assert!(agent_memory.config.defang_inbound);
assert_eq!(agent_memory.path, tmp.path().join("agent-memory"));
}
#[test]
fn agent_memory_census_slot_covers_both_store_kinds() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({ "runtime_options": { "agent_memory": true } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let slot = agent_memory_census_slot(&options.agent_memory.expect("agent memory options"));
assert_eq!(slot.declaration.domain, "agent-memory");
assert_eq!(
slot.declaration.resolution,
meerkat_core::DurabilityResolution::Persistent
);
assert_eq!(slot.backend, "SqliteAgentMemoryStore");
assert!(!slot.degraded);
let params = json!({ "runtime_options": { "agent_memory": { "store": "markdown" } } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let slot = agent_memory_census_slot(&options.agent_memory.expect("agent memory options"));
assert_eq!(slot.declaration.domain, "agent-memory");
assert_eq!(slot.backend, "MarkdownAgentMemoryStore");
}
#[test]
fn gateway_runtime_options_parse_agent_memory_defang_inbound() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": { "defang_inbound": false }
}
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let agent_memory = options.agent_memory.expect("agent memory options");
assert!(!agent_memory.config.defang_inbound);
let params = json!({
"runtime_options": {
"agent_memory": { "defang_inbound": "yes" }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("non-boolean defang_inbound should fail loudly"),
Err(err) => err,
};
assert!(err.contains("defang_inbound"), "{err}");
}
#[test]
fn gateway_runtime_options_parse_agent_memory_recall_policy() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": {
"recall_timeout_ms": 1200,
"recall_failure_policy": "fail"
}
}
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.config.recall_timeout_ms, 1200);
assert_eq!(
agent_memory.config.recall_failure_policy,
meerkat_mobkit::AgentMemoryRecallFailurePolicy::Fail
);
}
#[test]
fn gateway_runtime_options_reject_non_boolean_agent_memory_enabled() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": {
"enabled": "false"
}
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("non-boolean agent memory enabled should fail"),
Err(err) => err,
};
assert!(err.contains("enabled"), "{err}");
}
#[test]
fn gateway_runtime_options_reject_agent_memory_without_state_path() {
let params = json!({
"runtime_options": {
"agent_memory": true
}
});
let err = match parse_gateway_runtime_options(¶ms, None) {
Ok(_) => panic!("agent memory without path should fail"),
Err(err) => err,
};
assert!(err.contains("agent_memory"), "{err}");
}
#[test]
fn gateway_runtime_options_reject_agent_memory_path_override() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": {
"path": "/tmp/other-agent-memory"
}
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("agent memory path override should fail"),
Err(err) => err,
};
assert!(err.contains("path"), "{err}");
}
#[test]
fn gateway_runtime_options_agent_memory_store_defaults_to_sqlite() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": true
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("boolean agent memory config should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.store, GatewayAgentMemoryStoreKind::Sqlite);
let params = json!({
"runtime_options": {
"agent_memory": { "enabled": true }
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("object agent memory config should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.store, GatewayAgentMemoryStoreKind::Sqlite);
}
#[test]
fn gateway_runtime_options_agent_memory_store_accepts_markdown() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown" }
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("markdown store config should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.store, GatewayAgentMemoryStoreKind::Markdown);
assert_eq!(agent_memory.path, tmp.path().join("agent-memory"));
}
#[test]
fn gateway_runtime_options_agent_memory_store_rejects_unknown_backend() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": { "store": "postgres" }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("unknown store backend should fail loudly"),
Err(err) => err,
};
assert!(err.contains("'markdown' or 'sqlite'"), "{err}");
}
#[test]
fn gateway_runtime_options_agent_memory_budgeted_injection_requires_sqlite() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "per_turn_injection": "budgeted" }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("markdown + budgeted injection should fail loudly"),
Err(err) => err,
};
assert!(
err.contains("per_turn_injection='budgeted' requires store='sqlite'"),
"{err}"
);
let params = json!({
"runtime_options": {
"agent_memory": { "per_turn_injection": "budgeted" }
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("sqlite + budgeted should parse");
assert_eq!(
options
.agent_memory
.expect("agent memory")
.config
.per_turn_injection,
meerkat_mobkit::AgentMemoryPerTurnInjection::Budgeted
);
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "per_turn_injection": "off" }
}
});
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("markdown + off should parse");
}
#[test]
fn operator_scope_warning_fires_exactly_when_provisional_without_resolver() {
let tmp = tempfile::tempdir().expect("temp dir");
let parse = |scope: &str| {
let params = json!({
"runtime_options": {
"agent_memory": { "operator_scope": scope }
}
});
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("parse")
.agent_memory
.expect("agent memory options")
};
let provisional = parse("provisional");
let off = parse("off");
assert!(operator_scope_recall_inert(Some(&provisional), false));
assert!(!operator_scope_recall_inert(Some(&provisional), true));
assert!(!operator_scope_recall_inert(Some(&off), false));
assert!(!operator_scope_recall_inert(Some(&off), true));
assert!(!operator_scope_recall_inert(None, false));
assert!(!operator_scope_recall_inert(None, true));
}
#[test]
fn gateway_runtime_options_agent_memory_distiller_parse_matrix() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({ "runtime_options": { "agent_memory": true } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let agent_memory = options.agent_memory.expect("agent memory options");
assert!(!agent_memory.distiller.enabled);
assert_eq!(agent_memory.distiller.runs_per_hour, 12);
assert_eq!(agent_memory.distiller.min_interactions, 3);
assert_eq!(agent_memory.distiller.model, None);
let params = json!({
"runtime_options": {
"agent_memory": {
"distiller": {
"enabled": true,
"runs_per_hour": 6,
"min_interactions": 5,
"model": "claude-haiku-4-5"
}
}
}
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
let distiller = options.agent_memory.expect("agent memory").distiller;
assert!(distiller.enabled);
assert_eq!(distiller.runs_per_hour, 6);
assert_eq!(distiller.min_interactions, 5);
assert_eq!(distiller.model.as_deref(), Some("claude-haiku-4-5"));
let params = json!({
"runtime_options": { "agent_memory": { "distiller": true } }
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
assert!(
options
.agent_memory
.expect("agent memory")
.distiller
.enabled
);
let params = json!({
"runtime_options": { "agent_memory": { "distiller": { "runs_per_hour": 2 } } }
});
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("runtime options");
assert!(
options
.agent_memory
.expect("agent memory")
.distiller
.enabled
);
for (params, needle) in [
(
json!({ "runtime_options": { "agent_memory": { "distiller": { "cadence": 5 } } } }),
"unsupported runtime_options.agent_memory.distiller fields",
),
(
json!({ "runtime_options": { "agent_memory": { "distiller": "on" } } }),
"must be a boolean or object",
),
(
json!({ "runtime_options": { "agent_memory": { "distiller": { "runs_per_hour": 0 } } } }),
"runs_per_hour must be between 1 and 240",
),
(
json!({ "runtime_options": { "agent_memory": { "distiller": { "min_interactions": 0 } } } }),
"min_interactions must be between 1 and 100",
),
(
json!({ "runtime_options": { "agent_memory": { "distiller": { "model": "" } } } }),
"model must be a non-empty string",
),
] {
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("expected fail-loud parse for {params}"),
Err(err) => err,
};
assert!(err.contains(needle), "{err}");
}
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "distiller": { "enabled": true } }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("markdown + distiller should fail loudly"),
Err(err) => err,
};
assert!(err.contains("distiller requires store='sqlite'"), "{err}");
}
#[test]
fn gateway_runtime_options_agent_memory_steward_parse_matrix() {
let tmp = tempfile::tempdir().expect("tempdir");
let params = json!({ "runtime_options": { "agent_memory": true } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("defaults parse");
let steward = options.agent_memory.expect("agent memory").steward;
assert!(!steward.enabled);
assert_eq!(steward.cadence, "*/6h");
assert!(!steward.per_mob);
assert_eq!(steward.runs_per_day, 4);
assert_eq!(steward.min_signals, 3);
assert_eq!(steward.model, None);
let params = json!({
"runtime_options": {
"agent_memory": {
"steward": {
"enabled": true,
"cadence": "*/30m",
"model": "claude-sonnet-4-6",
"per_mob": true,
"runs_per_day": 8,
"min_signals": 5
}
}
}
});
let steward = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("object form parses")
.agent_memory
.expect("agent memory")
.steward;
assert!(steward.enabled);
assert_eq!(steward.cadence, "*/30m");
assert_eq!(steward.model.as_deref(), Some("claude-sonnet-4-6"));
assert!(steward.per_mob);
assert_eq!(steward.runs_per_day, 8);
assert_eq!(steward.min_signals, 5);
let params = json!({
"runtime_options": { "agent_memory": { "steward": true } }
});
assert!(
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("bool form parses")
.agent_memory
.expect("agent memory")
.steward
.enabled
);
let params = json!({
"runtime_options": { "agent_memory": { "steward": { "runs_per_day": 2 } } }
});
assert!(
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("object without enabled parses")
.agent_memory
.expect("agent memory")
.steward
.enabled
);
for (params, needle) in [
(
json!({ "runtime_options": { "agent_memory": { "steward": { "tempo": 5 } } } }),
"unsupported runtime_options.agent_memory.steward fields",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": "on" } } }),
"must be a boolean or object",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "cadence": "0 9 * * *" } } } }),
"not an interval marker",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "cadence": "" } } } }),
"cadence must be a non-empty string",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "runs_per_day": 0 } } } }),
"runs_per_day must be between 1 and 96",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "min_signals": 0 } } } }),
"min_signals must be between 1 and 1000",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "model": "" } } } }),
"model must be a non-empty string",
),
(
json!({ "runtime_options": { "agent_memory": { "steward": { "per_mob": "yes" } } } }),
"per_mob must be a boolean",
),
] {
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("expected fail-loud parse for {params}"),
Err(err) => err,
};
assert!(err.contains(needle), "{err}");
}
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "steward": { "enabled": true } }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("markdown + steward should fail loudly"),
Err(err) => err,
};
assert!(err.contains("steward requires store='sqlite'"), "{err}");
}
#[test]
fn gateway_runtime_options_agent_memory_hygienist_parse_matrix() {
let tmp = tempfile::tempdir().expect("tempdir");
let params = json!({ "runtime_options": { "agent_memory": true } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("defaults parse");
let hygienist = options.agent_memory.expect("agent memory").hygienist;
assert!(!hygienist.enabled);
assert_eq!(hygienist.runs_per_day, 2);
assert_eq!(hygienist.model, None);
let params = json!({
"runtime_options": {
"agent_memory": {
"hygienist": {
"enabled": true,
"runs_per_day": 4,
"model": "claude-sonnet-4-6"
}
}
}
});
let hygienist = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("object form parses")
.agent_memory
.expect("agent memory")
.hygienist;
assert!(hygienist.enabled);
assert_eq!(hygienist.runs_per_day, 4);
assert_eq!(hygienist.model.as_deref(), Some("claude-sonnet-4-6"));
let params = json!({
"runtime_options": { "agent_memory": { "hygienist": true } }
});
assert!(
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("bool form parses")
.agent_memory
.expect("agent memory")
.hygienist
.enabled
);
let params = json!({
"runtime_options": { "agent_memory": { "hygienist": { "runs_per_day": 1 } } }
});
assert!(
parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("object without enabled parses")
.agent_memory
.expect("agent memory")
.hygienist
.enabled
);
for (params, needle) in [
(
json!({ "runtime_options": { "agent_memory": { "hygienist": { "cadence": "*/6h" } } } }),
"unsupported runtime_options.agent_memory.hygienist fields",
),
(
json!({ "runtime_options": { "agent_memory": { "hygienist": "on" } } }),
"must be a boolean or object",
),
(
json!({ "runtime_options": { "agent_memory": { "hygienist": { "runs_per_day": 0 } } } }),
"runs_per_day must be between 1 and 24",
),
(
json!({ "runtime_options": { "agent_memory": { "hygienist": { "runs_per_day": 48 } } } }),
"runs_per_day must be between 1 and 24",
),
(
json!({ "runtime_options": { "agent_memory": { "hygienist": { "model": "" } } } }),
"model must be a non-empty string",
),
(
json!({ "runtime_options": { "agent_memory": { "hygienist": { "enabled": "yes" } } } }),
"enabled must be a boolean",
),
] {
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("expected fail-loud parse for {params}"),
Err(err) => err,
};
assert!(err.contains(needle), "{err}");
}
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "hygienist": { "enabled": true } }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("markdown + hygienist should fail loudly"),
Err(err) => err,
};
assert!(err.contains("hygienist requires store='sqlite'"), "{err}");
}
#[test]
fn gateway_runtime_options_agent_memory_operator_scope_parse_matrix() {
let tmp = tempfile::tempdir().expect("tempdir");
let params = json!({ "runtime_options": { "agent_memory": true } });
let options =
parse_gateway_runtime_options(¶ms, Some(tmp.path())).expect("defaults parse");
assert_eq!(
options
.agent_memory
.expect("agent memory")
.config
.operator_scope,
meerkat_mobkit::AgentMemoryOperatorScope::Off
);
for (value, expected) in [
("off", meerkat_mobkit::AgentMemoryOperatorScope::Off),
(
"provisional",
meerkat_mobkit::AgentMemoryOperatorScope::Provisional,
),
] {
let params = json!({
"runtime_options": { "agent_memory": { "operator_scope": value } }
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("operator_scope parses");
assert_eq!(
options
.agent_memory
.expect("agent memory")
.config
.operator_scope,
expected
);
}
for params in [
json!({ "runtime_options": { "agent_memory": { "operator_scope": "final" } } }),
json!({ "runtime_options": { "agent_memory": { "operator_scope": true } } }),
] {
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("expected fail-loud parse for {params}"),
Err(err) => err,
};
assert!(err.contains("operator_scope must be 'off' or"), "{err}");
}
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "operator_scope": "provisional" }
}
});
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("markdown + operator_scope should fail loudly"),
Err(err) => err,
};
assert!(
err.contains("operator_scope requires store='sqlite'"),
"{err}"
);
}
#[test]
fn gateway_mob_binding_resolves_hosting_mob_only_for_matching_realm() {
use meerkat_mobkit::memory::coordinator::{MobScopeResolver, StaticMobBinding};
let binding = StaticMobBinding {
realm: "default".to_string(),
mob: "mob:example".to_string(),
};
assert_eq!(
binding.active_mobs("default", "identity:frontdesk"),
vec!["mob:example".to_string()],
"matching realm must compose the hosting mob"
);
assert!(
binding
.active_mobs("other-realm", "identity:frontdesk")
.is_empty(),
"foreign realm must compose no mob scope"
);
}
#[test]
fn compaction_reset_sink_fires_on_attributed_compaction_only() {
use meerkat_core::event::AgentEvent;
let seen: Arc<std::sync::Mutex<Vec<String>>> = Arc::new(std::sync::Mutex::new(Vec::new()));
let sink = meerkat_mobkit::CompactionResetSink::new({
let seen = seen.clone();
Arc::new(move |session: &str| {
seen.lock()
.unwrap_or_else(|err| err.into_inner())
.push(session.to_string());
})
});
let session = meerkat_core::types::SessionId::new();
let envelope = |source, payload| meerkat_core::event::EventEnvelope {
event_id: Default::default(),
source,
seq: 0,
mob_id: None,
timestamp_ms: 0,
payload,
};
let compaction = || AgentEvent::CompactionCompleted {
summary_tokens: 10,
messages_before: 20,
messages_after: 2,
};
use meerkat_mobkit::MemberAgentEventSink as _;
sink.observe(
"identity:a",
&envelope(
meerkat_core::event::EventSourceIdentity::Session {
session_id: session.clone(),
},
compaction(),
),
);
assert_eq!(
seen.lock().unwrap_or_else(|err| err.into_inner()).clone(),
vec![session.to_string()],
"reset must fire with the compacted session's key"
);
sink.observe(
"identity:a",
&envelope(
meerkat_core::event::EventSourceIdentity::Session {
session_id: session.clone(),
},
AgentEvent::RunCompleted {
session_id: session.clone(),
result: "done".to_string(),
structured_output: None,
extraction_required: false,
usage: Default::default(),
terminal_cause_kind: None,
},
),
);
sink.observe(
"identity:a",
&envelope(
meerkat_core::event::EventSourceIdentity::Callback,
compaction(),
),
);
assert_eq!(
seen.lock().unwrap_or_else(|err| err.into_inner()).len(),
1,
"only attributed compaction events reset"
);
}
#[test]
fn gateway_runtime_options_parse_agent_memory_taint_knobs() {
let tmp = tempfile::tempdir().expect("temp dir");
let params = json!({
"runtime_options": {
"agent_memory": {
"llm_writes": "quarantined",
"recorder_tool": false,
"content_trust": {
"trusted_mcp_servers": ["knowledge_graph"],
"untrusted_tools": ["scrape_page"],
"trusted_tools": ["safe_calc"],
}
}
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("taint knobs should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(
agent_memory.config.llm_writes,
meerkat_mobkit::AgentMemoryLlmWrites::Quarantined
);
assert!(!agent_memory.config.recorder_tool);
assert_eq!(
agent_memory.config.content_trust.trusted_mcp_servers,
vec!["knowledge_graph"]
);
assert_eq!(
agent_memory.config.content_trust.untrusted_tools,
vec!["scrape_page"]
);
let params = json!({
"runtime_options": { "agent_memory": true }
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("boolean agent memory config should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(
agent_memory.config.llm_writes,
meerkat_mobkit::AgentMemoryLlmWrites::Observed
);
assert!(agent_memory.config.recorder_tool);
assert_eq!(
agent_memory.config.content_trust,
meerkat_mobkit::ContentTrustConfig::default()
);
}
#[test]
fn gateway_runtime_options_reject_bad_taint_knobs() {
let tmp = tempfile::tempdir().expect("temp dir");
for (block, needle) in [
(
json!({ "llm_writes": "yolo" }),
"'observed' or 'quarantined'",
),
(json!({ "llm_writes": true }), "'observed' or 'quarantined'"),
(json!({ "recorder_tool": "yes" }), "must be a boolean"),
(
json!({ "content_trust": { "servers": [] } }),
"unsupported content_trust fields",
),
(
json!({ "content_trust": { "trusted_mcp_servers": "kg" } }),
"must be an array",
),
(
json!({ "store": "markdown", "llm_writes": "quarantined" }),
"require store='sqlite'",
),
(
json!({ "store": "markdown", "content_trust": {} }),
"require store='sqlite'",
),
] {
let params = json!({ "runtime_options": { "agent_memory": block } });
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("config must fail loudly: {params}"),
Err(err) => err,
};
assert!(err.contains(needle), "{err} (wanted '{needle}')");
}
}
#[test]
fn gateway_runtime_options_parse_agent_memory_selector() {
use meerkat_mobkit::memory::selector::SelectorSpec;
let tmp = tempfile::tempdir().expect("temp dir");
for params in [
json!({ "runtime_options": { "agent_memory": true } }),
json!({ "runtime_options": { "agent_memory": {} } }),
] {
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("agent memory config should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.selector, None);
}
for (value, want) in [
("off", SelectorSpec::Off),
("default", SelectorSpec::Default),
(
"profile:/etc/mobkit/selector.toml",
SelectorSpec::Profile(std::path::PathBuf::from("/etc/mobkit/selector.toml")),
),
] {
let params = json!({
"runtime_options": { "agent_memory": { "selector": value } }
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("selector value should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(agent_memory.selector, Some(want), "selector '{value}'");
}
}
#[test]
fn gateway_runtime_options_reject_bad_agent_memory_selector() {
let tmp = tempfile::tempdir().expect("temp dir");
for (block, needle) in [
(json!({ "selector": "on" }), "'off', 'default'"),
(json!({ "selector": "profile:" }), "'off', 'default'"),
(json!({ "selector": true }), "must be a string"),
(
json!({ "store": "markdown", "selector": "default" }),
"selector requires store='sqlite'",
),
] {
let params = json!({ "runtime_options": { "agent_memory": block } });
let err = match parse_gateway_runtime_options(¶ms, Some(tmp.path())) {
Ok(_) => panic!("selector config must fail loudly: {params}"),
Err(err) => err,
};
assert!(err.contains(needle), "{err} (wanted '{needle}')");
}
let params = json!({
"runtime_options": {
"agent_memory": { "store": "markdown", "selector": "off" }
}
});
let options = parse_gateway_runtime_options(¶ms, Some(tmp.path()))
.expect("selector=off with markdown should parse");
let agent_memory = options.agent_memory.expect("agent memory options");
assert_eq!(
agent_memory.selector,
Some(meerkat_mobkit::memory::selector::SelectorSpec::Off)
);
}
#[test]
fn selector_config_takes_precedence_over_env_fallback() {
use meerkat_mobkit::memory::selector::SelectorSpec;
assert_eq!(
resolve_selector_spec(Some(&SelectorSpec::Default)).expect("configured spec resolves"),
SelectorSpec::Default,
"agent_memory.selector must be used ahead of the env fallback"
);
let profile = SelectorSpec::Profile(std::path::PathBuf::from("/etc/mobkit/selector.toml"));
assert_eq!(
resolve_selector_spec(Some(&profile)).expect("configured profile resolves"),
profile
);
assert_eq!(
resolve_selector_spec(None).expect("env fallback resolves"),
SelectorSpec::Off,
"unset config must fall back to MOBKIT_AGENT_MEMORY_SELECTOR (unset ⇒ off)"
);
}
#[test]
fn gateway_runtime_options_reject_zero_max_sessions() {
let params = json!({
"runtime_options": {
"max_sessions": 0
}
});
let err = match parse_gateway_runtime_options(¶ms, None) {
Ok(_) => panic!("zero should fail"),
Err(err) => err,
};
assert!(err.contains("runtime_options.max_sessions"));
}
#[test]
fn gateway_runtime_options_workgraph_parses_bool_and_path_forms() {
let defaults = parse_gateway_runtime_options(&json!({}), None).expect("defaults parse");
assert_eq!(
defaults.workgraph,
GatewayWorkgraphOption::Enabled,
"workgraph defaults on"
);
let enabled = parse_gateway_runtime_options(
&json!({ "runtime_options": { "workgraph": true } }),
None,
)
.expect("explicit true parses");
assert_eq!(enabled.workgraph, GatewayWorkgraphOption::Enabled);
let disabled = parse_gateway_runtime_options(
&json!({ "runtime_options": { "workgraph": false } }),
None,
)
.expect("explicit false parses");
assert_eq!(disabled.workgraph, GatewayWorkgraphOption::Disabled);
let durable = parse_gateway_runtime_options(
&json!({ "runtime_options": { "workgraph": "/var/lib/mobkit/workgraph" } }),
None,
)
.expect("string path parses");
assert_eq!(
durable.workgraph,
GatewayWorkgraphOption::DurableDir("/var/lib/mobkit/workgraph".into())
);
for invalid in [json!(42), json!(""), json!(" "), json!({ "dir": "x" })] {
let err = match parse_gateway_runtime_options(
&json!({ "runtime_options": { "workgraph": invalid } }),
None,
) {
Ok(_) => panic!("invalid workgraph value should fail: {invalid}"),
Err(err) => err,
};
assert!(err.contains("runtime_options.workgraph"), "{err}");
}
}
#[test]
fn gateway_identity_bootstrap_mode_parse_matrix_is_strict() {
use meerkat_mobkit::IdentityBootstrapMode;
let defaults = parse_gateway_runtime_options(&json!({}), None).expect("defaults");
assert_eq!(defaults.identity_bootstrap_mode, None);
for (wire, expected) in [
(
json!({"mode": "eager_materialize"}),
IdentityBootstrapMode::EagerMaterialize,
),
(
json!({"mode": "lazy_materialize"}),
IdentityBootstrapMode::LazyMaterialize,
),
(
json!({"mode": "lazy_with_background_warm", "concurrency": 2}),
IdentityBootstrapMode::LazyWithBackgroundWarm { concurrency: 2 },
),
] {
let options = parse_gateway_runtime_options(
&json!({"runtime_options": {"identity_bootstrap_mode": wire}}),
None,
)
.expect("valid bootstrap mode");
assert_eq!(options.identity_bootstrap_mode, Some(expected));
}
for invalid in [
json!("lazy_materialize"),
json!({}),
json!({"mode": "unknown"}),
json!({"mode": "lazy_materialize", "concurrency": 2}),
json!({"mode": "lazy_with_background_warm"}),
json!({"mode": "lazy_with_background_warm", "concurrency": 0}),
json!({"mode": "lazy_with_background_warm", "concurrency": 17}),
json!({"mode": "eager_materialize", "extra": true}),
] {
let error = parse_gateway_runtime_options(
&json!({"runtime_options": {"identity_bootstrap_mode": invalid}}),
None,
)
.err()
.expect("invalid bootstrap mode must fail");
assert!(error.contains("identity_bootstrap_mode"), "{error}");
}
}
#[tokio::test]
async fn callback_close_wakes_pending_call_and_rejects_late_admission() {
let (stdout_tx, mut stdout_rx) = mpsc::channel(4);
let bridge = StdioCallbackBridge::new(stdout_tx);
let pending = tokio::spawn({
let bridge = bridge.clone();
async move { bridge.call("callback/build_agent", json!({})).await }
});
tokio::time::timeout(Duration::from_secs(1), stdout_rx.recv())
.await
.expect("pending callback should be written")
.expect("stdout callback channel should remain open");
bridge.close().await;
let error = tokio::time::timeout(Duration::from_secs(1), pending)
.await
.expect("close should wake the pending callback")
.expect("callback task should not panic")
.expect_err("closed callback must fail");
assert!(error.contains("channel dropped"), "{error}");
let late_error = bridge
.call("callback/build_agent", json!({}))
.await
.expect_err("close must reject callbacks admitted after EOF");
assert!(late_error.contains("transport closed"), "{late_error}");
assert!(bridge.state.lock().await.pending.is_empty());
}
#[tokio::test]
async fn callback_builder_delegates_absent_session_compaction_reconciliation() {
let tmp = tempfile::tempdir().expect("temp dir");
let inner = FactoryAgentBuilder::new(AgentFactory::new(tmp.path()), Config::default());
let (stdout_tx, _stdout_rx) = mpsc::channel(1);
let builder = StdioCallbackAgentBuilder {
inner,
bridge: StdioCallbackBridge::new(stdout_tx),
has_session_builder: false,
session_store: None,
detached_jobs: None,
};
let session_id = meerkat_core::SessionId::parse("019f74fb-1907-7b21-932d-ab22c4d1f532")
.expect("valid session id");
builder
.abort_absent_session_compaction_stages(&session_id)
.await
.expect("callback wrapper must preserve the inner factory's durable-memory seam");
}
#[test]
fn advertised_shutdown_horizon_covers_every_bounded_gateway_phase() {
fn completed_report() -> UnifiedRuntimeShutdownReport {
UnifiedRuntimeShutdownReport {
drain: ShutdownDrainReport {
drained_count: 1,
timed_out: false,
drain_duration_ms: 2,
},
module_shutdown: RuntimeShutdownReport {
terminated_modules: vec!["router".to_string()],
orphan_processes: 0,
},
mob_stop: Ok(()),
identity_authority_release: IdentityAuthorityReleaseOutcome::Released {
grant_count: 1,
},
}
}
fn assert_not_attested(report: Option<&UnifiedRuntimeShutdownReport>) {
let response = gateway_shutdown_response(json!("shutdown"), report);
assert_eq!(response["result"]["shutdown"], false);
assert_eq!(response["result"]["runtime_cleanup_completed"], false);
}
assert_eq!(PROVIDER_CALLBACK_TIMEOUT, Duration::from_secs(130));
let mob_quiesce_window = Duration::from_secs(10);
let scheduler_overhead = Duration::from_secs(10);
assert_eq!(
GATEWAY_RUNTIME_SHUTDOWN_TIMEOUT,
PROVIDER_CALLBACK_TIMEOUT
+ PROVIDER_CALLBACK_TIMEOUT
+ GATEWAY_RUNTIME_EVENT_DRAIN_TIMEOUT
+ mob_quiesce_window
+ scheduler_overhead,
"runtime budget must exactly cover both provider callbacks, event drain, mob quiesce, and scheduler overhead"
);
let gateway_phase_budget = GATEWAY_RPC_DRAIN_TIMEOUT
+ GATEWAY_HTTP_DRAIN_TIMEOUT
+ GATEWAY_RUNTIME_SHUTDOWN_TIMEOUT
+ GATEWAY_STDOUT_DRAIN_TIMEOUT;
assert_eq!(gateway_phase_budget, Duration::from_secs(325));
assert_eq!(
Duration::from_millis(GATEWAY_SHUTDOWN_HORIZON_MS),
Duration::from_secs(335)
);
assert_eq!(
Duration::from_millis(GATEWAY_SHUTDOWN_HORIZON_MS).saturating_sub(gateway_phase_budget),
Duration::from_secs(10),
"advertised horizon must leave response/reaping margin"
);
assert_not_attested(None);
let successful_report = completed_report();
let completed = gateway_shutdown_response(json!("shutdown"), Some(&successful_report));
assert_eq!(completed["result"]["shutdown"], true);
assert_eq!(completed["result"]["runtime_cleanup_completed"], true);
let mut report = completed_report();
report.drain.timed_out = true;
assert_not_attested(Some(&report));
let mut report = completed_report();
report.mob_stop = Err(MobRuntimeError::InvalidConfig(
"mob stop failed".to_string(),
));
report.identity_authority_release = IdentityAuthorityReleaseOutcome::SkippedMobStopFailed;
assert_not_attested(Some(&report));
let mut report = completed_report();
report.identity_authority_release =
IdentityAuthorityReleaseOutcome::SkippedResetCleanupFailed {
error: "superseded generation cleanup remains pending".to_string(),
};
assert_not_attested(Some(&report));
let mut report = completed_report();
report.identity_authority_release = IdentityAuthorityReleaseOutcome::Failed {
error: "provider release failed".to_string(),
};
assert_not_attested(Some(&report));
let mut report = completed_report();
report.module_shutdown.orphan_processes = 1;
assert_not_attested(Some(&report));
}
#[test]
fn gateway_explicit_identity_bootstrap_mode_always_requires_roster() {
use meerkat_mobkit::IdentityBootstrapMode;
validate_gateway_identity_bootstrap_intent(None, false)
.expect("omitted mode preserves the classic gateway");
validate_gateway_identity_bootstrap_intent(None, true)
.expect("a roster may use the eager compatibility default");
for mode in [
IdentityBootstrapMode::EagerMaterialize,
IdentityBootstrapMode::LazyMaterialize,
IdentityBootstrapMode::LazyWithBackgroundWarm { concurrency: 2 },
] {
let error = validate_gateway_identity_bootstrap_intent(Some(&mode), false)
.expect_err("every explicit identity mode requires a roster");
assert!(error.contains("roster provider"), "{error}");
validate_gateway_identity_bootstrap_intent(Some(&mode), true)
.expect("an explicit mode with a roster is valid");
}
}
}
fn parse_gateway_identity_bootstrap_mode(
value: &Value,
) -> Result<meerkat_mobkit::IdentityBootstrapMode, String> {
use meerkat_mobkit::identity_first::MAX_IDENTITY_BACKGROUND_WARM_CONCURRENCY;
let object = value
.as_object()
.ok_or_else(|| "runtime_options.identity_bootstrap_mode must be an object".to_string())?;
let unsupported = object
.keys()
.filter(|key| !matches!(key.as_str(), "mode" | "concurrency"))
.cloned()
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.identity_bootstrap_mode fields: {}",
unsupported.join(", ")
));
}
let mode = object.get("mode").and_then(Value::as_str).ok_or_else(|| {
"runtime_options.identity_bootstrap_mode.mode must be a string".to_string()
})?;
let parsed = match mode {
"eager_materialize" => {
if object.contains_key("concurrency") {
return Err(
"runtime_options.identity_bootstrap_mode.concurrency is only valid for lazy_with_background_warm"
.to_string(),
);
}
meerkat_mobkit::IdentityBootstrapMode::EagerMaterialize
}
"lazy_materialize" => {
if object.contains_key("concurrency") {
return Err(
"runtime_options.identity_bootstrap_mode.concurrency is only valid for lazy_with_background_warm"
.to_string(),
);
}
meerkat_mobkit::IdentityBootstrapMode::LazyMaterialize
}
"lazy_with_background_warm" => {
let concurrency = object
.get("concurrency")
.and_then(Value::as_u64)
.ok_or_else(|| {
"runtime_options.identity_bootstrap_mode.concurrency must be a positive integer"
.to_string()
})?;
let concurrency = usize::try_from(concurrency).map_err(|_| {
"runtime_options.identity_bootstrap_mode.concurrency is too large".to_string()
})?;
if !(1..=MAX_IDENTITY_BACKGROUND_WARM_CONCURRENCY).contains(&concurrency) {
return Err(format!(
"runtime_options.identity_bootstrap_mode.concurrency must be between 1 and {MAX_IDENTITY_BACKGROUND_WARM_CONCURRENCY}"
));
}
meerkat_mobkit::IdentityBootstrapMode::LazyWithBackgroundWarm { concurrency }
}
_ => {
return Err(format!(
"runtime_options.identity_bootstrap_mode.mode '{mode}' is unsupported"
));
}
};
Ok(parsed)
}
fn validate_gateway_identity_bootstrap_intent(
configured_mode: Option<&meerkat_mobkit::IdentityBootstrapMode>,
has_roster_provider: bool,
) -> Result<(), String> {
if configured_mode.is_some() && !has_roster_provider {
return Err(
"runtime_options.identity_bootstrap_mode requires an identity-first roster provider"
.to_string(),
);
}
Ok(())
}
fn parse_gateway_runtime_options(
params: &Value,
persistent_state: Option<&std::path::Path>,
) -> Result<GatewayRuntimeOptions, String> {
let Some(runtime_options) = params.get("runtime_options") else {
return Ok(GatewayRuntimeOptions::default());
};
let runtime_options = runtime_options
.as_object()
.ok_or_else(|| "runtime_options must be a JSON object".to_string())?;
let supported = [
"memory_config",
"identity_bootstrap_mode",
"routing_config_path",
"scheduling_files",
"gating_config_path",
"auth_config",
"access_config_path",
"console_config_path",
"console_require_app_auth",
"console_read_only",
"console_fetch_timeout_ms",
"demo_llm",
"member_comms_address",
"max_sessions",
"event_log",
"agent_memory",
"implicit_delegate_idle_retire_secs",
"implicit_delegate_idle_sweep_interval_ms",
"workgraph",
"live",
"host_runnables",
"runtime_store",
"compaction",
];
let unsupported = runtime_options
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options fields: {}",
unsupported.join(", ")
));
}
let mut parsed = GatewayRuntimeOptions::default();
if let Some(value) = runtime_options.get("identity_bootstrap_mode") {
parsed.identity_bootstrap_mode = Some(parse_gateway_identity_bootstrap_mode(value)?);
}
if let Some(memory_config) = runtime_options.get("memory_config") {
parsed.runtime_options.memory_backend = Some(parse_gateway_memory_config(
memory_config,
persistent_state,
)?);
}
if let Some(agent_memory) = runtime_options.get("agent_memory") {
parsed.agent_memory = parse_gateway_agent_memory_config(agent_memory, persistent_state)?;
}
if let Some(path) = runtime_options
.get("routing_config_path")
.and_then(Value::as_str)
{
parsed.routing_routes = parse_gateway_routing_config_path(path)?;
}
if let Some(files) = runtime_options.get("scheduling_files") {
parsed.schedules = parse_gateway_scheduling_files(files)?;
}
if let Some(path) = runtime_options
.get("gating_config_path")
.and_then(Value::as_str)
{
parsed.gating = parse_gateway_gating_config_path(path)?;
}
if let Some(auth_config) = runtime_options.get("auth_config") {
parsed.decisions = Some(parse_gateway_auth_config(auth_config)?);
}
if let Some(path) = runtime_options
.get("console_config_path")
.and_then(Value::as_str)
{
parsed.console_ui = load_console_ui_config_from_path_for_realm(path, None)
.map_err(|err| format!("runtime_options.console_config_path is invalid: {err}"))?;
}
if let Some(path) = runtime_options
.get("access_config_path")
.and_then(Value::as_str)
{
parsed.access = Some(
meerkat_mobkit::AccessController::load_or_default(path)
.map_err(|err| format!("runtime_options.access_config_path is invalid: {err}"))?,
);
}
if let Some(value) = runtime_options.get("console_require_app_auth") {
parsed.console_require_app_auth = Some(value.as_bool().ok_or_else(|| {
"runtime_options.console_require_app_auth must be a boolean".to_string()
})?);
}
if let Some(value) = runtime_options.get("console_read_only") {
parsed.console_read_only =
Some(value.as_bool().ok_or_else(|| {
"runtime_options.console_read_only must be a boolean".to_string()
})?);
}
if let Some(value) = runtime_options.get("console_fetch_timeout_ms") {
if !value.is_null() {
let timeout_ms = value.as_u64().ok_or_else(|| {
"runtime_options.console_fetch_timeout_ms must be a positive integer or null"
.to_string()
})?;
if timeout_ms == 0 {
return Err(
"runtime_options.console_fetch_timeout_ms must be greater than zero"
.to_string(),
);
}
parsed.console_fetch_timeout_ms = Some(timeout_ms);
}
}
if let Some(value) = runtime_options.get("workgraph") {
parsed.workgraph = match value {
Value::Bool(true) => GatewayWorkgraphOption::Enabled,
Value::Bool(false) => GatewayWorkgraphOption::Disabled,
Value::String(path) if !path.trim().is_empty() => {
GatewayWorkgraphOption::DurableDir(std::path::PathBuf::from(path))
}
_ => {
return Err(
"runtime_options.workgraph must be a boolean or a non-empty string \
(the directory for the durable workgraph store)"
.to_string(),
);
}
};
}
if let Some(value) = runtime_options.get("live") {
parsed.live = match value {
Value::Bool(true) => GatewayLiveOption::Enabled {
public_base_url: None,
seed_max_chars: None,
},
Value::Bool(false) => GatewayLiveOption::Disabled,
Value::Object(map) => {
let enabled = map.get("enabled").and_then(Value::as_bool).unwrap_or(true);
let seed_max_chars = match map.get("seed_max_chars") {
None | Some(Value::Null) => None,
Some(value) => {
let chars = value.as_u64().filter(|chars| *chars > 0).ok_or_else(|| {
"runtime_options.live.seed_max_chars must be a positive integer"
.to_string()
})?;
Some(usize::try_from(chars).map_err(|_| {
"runtime_options.live.seed_max_chars is too large".to_string()
})?)
}
};
if enabled {
GatewayLiveOption::Enabled {
public_base_url: map
.get("public_base_url")
.and_then(Value::as_str)
.map(str::to_string),
seed_max_chars,
}
} else {
GatewayLiveOption::Disabled
}
}
_ => {
return Err(
"runtime_options.live must be a boolean or an object ({enabled?, public_base_url?, seed_max_chars?})"
.to_string(),
);
}
};
}
if let Some(value) = runtime_options.get("host_runnables") {
let entries = value.as_array().ok_or_else(|| {
"runtime_options.host_runnables must be an array of runnable names".to_string()
})?;
let mut names = Vec::with_capacity(entries.len());
for entry in entries {
let text = entry.as_str().ok_or_else(|| {
"runtime_options.host_runnables entries must be strings".to_string()
})?;
let name = meerkat::HostRunnableName::parse(text).map_err(|err| {
format!("runtime_options.host_runnables entry '{text}' is invalid: {err}")
})?;
if name.as_str() == meerkat_mobkit::schedule_wiring::STEWARD_DREAM_RUNNABLE {
return Err(format!(
"runtime_options.host_runnables entry '{text}' collides with the \
reserved steward dream runnable"
));
}
if names.contains(&name) {
return Err(format!(
"runtime_options.host_runnables entry '{text}' is duplicated"
));
}
names.push(name);
}
parsed.host_runnables = names;
}
if let Some(value) = runtime_options.get("demo_llm") {
parsed.demo_llm = value
.as_bool()
.ok_or_else(|| "runtime_options.demo_llm must be a boolean".to_string())?;
}
if let Some(value) = runtime_options.get("member_comms_address") {
let address = value.as_str().ok_or_else(|| {
"runtime_options.member_comms_address must be a socket address string".to_string()
})?;
let socket = address
.parse::<std::net::SocketAddr>()
.map_err(|error| format!("runtime_options.member_comms_address is invalid: {error}"))?;
if socket.ip().is_unspecified() {
return Err(
"runtime_options.member_comms_address must name a concrete interface; wildcard binds cannot be advertised to external peers"
.to_string(),
);
}
parsed.member_comms_address = Some(address.to_string());
}
if let Some(value) = runtime_options.get("max_sessions") {
let max_sessions = value
.as_u64()
.ok_or_else(|| "runtime_options.max_sessions must be a positive integer".to_string())?;
if max_sessions == 0 {
return Err("runtime_options.max_sessions must be greater than zero".to_string());
}
parsed.max_sessions = usize::try_from(max_sessions)
.map_err(|_| "runtime_options.max_sessions is too large".to_string())?;
}
if let Some(event_log) = runtime_options.get("event_log") {
parsed.event_log = Some(parse_gateway_event_log_config(event_log)?);
}
if let Some(runtime_store) = runtime_options.get("runtime_store") {
parsed.runtime_store_ephemeral = parse_gateway_runtime_store_config(runtime_store)?;
}
if let Some(compaction) = runtime_options.get("compaction") {
parsed.compaction = Some(
meerkat_mobkit::parse_compaction_policy(compaction)
.map_err(|error| format!("runtime_options.compaction {error}"))?,
);
}
if let Some(value) = runtime_options.get("implicit_delegate_idle_retire_secs") {
parsed.runtime_options.implicit_delegate_idle_retire_secs = if value.is_null() {
None
} else {
Some(value.as_u64().ok_or_else(|| {
"runtime_options.implicit_delegate_idle_retire_secs must be a non-negative integer or null".to_string()
})?)
};
}
if let Some(value) = runtime_options.get("implicit_delegate_idle_sweep_interval_ms") {
let interval = value.as_u64().ok_or_else(|| {
"runtime_options.implicit_delegate_idle_sweep_interval_ms must be a positive integer"
.to_string()
})?;
if interval == 0 {
return Err(
"runtime_options.implicit_delegate_idle_sweep_interval_ms must be greater than zero"
.to_string(),
);
}
parsed
.runtime_options
.implicit_delegate_idle_sweep_interval_ms = interval;
}
if let Some(require_app_auth) = parsed.console_require_app_auth {
parsed
.decisions
.get_or_insert_with(minimal_decision_state)
.console
.require_app_auth = require_app_auth;
}
if let Some(read_only) = parsed.console_read_only {
parsed
.decisions
.get_or_insert_with(minimal_decision_state)
.console
.read_only = read_only;
}
if let Some(fetch_timeout_ms) = parsed.console_fetch_timeout_ms {
parsed
.decisions
.get_or_insert_with(minimal_decision_state)
.console
.fetch_timeout_ms = Some(fetch_timeout_ms);
}
if let Some(decisions) = parsed.decisions.as_mut() {
decisions.console.ui = parsed.console_ui.clone();
}
Ok(parsed)
}
fn read_gateway_config_file(path: &str, option_name: &str) -> Result<Value, String> {
let text = std::fs::read_to_string(path)
.map_err(|err| format!("failed to read runtime_options.{option_name} '{path}': {err}"))?;
if path.ends_with(".json") {
return serde_json::from_str(&text)
.map_err(|err| format!("invalid JSON in runtime_options.{option_name}: {err}"));
}
let toml_value: toml::Value = toml::from_str(&text)
.map_err(|err| format!("invalid TOML in runtime_options.{option_name}: {err}"))?;
serde_json::to_value(toml_value)
.map_err(|err| format!("failed to normalize runtime_options.{option_name}: {err}"))
}
fn parse_gateway_routing_config_path(path: &str) -> Result<Vec<RuntimeRoute>, String> {
let value = read_gateway_config_file(path, "routing_config_path")?;
let routes_value = value
.get("routes")
.cloned()
.unwrap_or_else(|| value.clone());
serde_json::from_value(routes_value)
.map_err(|err| format!("runtime_options.routing_config_path routes are invalid: {err}"))
}
fn parse_gateway_scheduling_files(files: &Value) -> Result<Vec<ScheduleDefinition>, String> {
let files = files
.as_array()
.ok_or_else(|| "runtime_options.scheduling_files must be an array".to_string())?;
let mut schedules = Vec::new();
for file in files {
let path = file.as_str().ok_or_else(|| {
"runtime_options.scheduling_files entries must be strings".to_string()
})?;
let value = read_gateway_config_file(path, "scheduling_files")?;
let schedules_value = value
.get("schedules")
.cloned()
.unwrap_or_else(|| value.clone());
let mut parsed: Vec<ScheduleDefinition> =
serde_json::from_value(schedules_value).map_err(|err| {
format!("runtime_options.scheduling_files schedule definitions are invalid: {err}")
})?;
schedules.append(&mut parsed);
}
meerkat_mobkit::evaluate_schedules_at_tick(&schedules, 0)
.map_err(|err| format!("runtime_options.scheduling_files are invalid: {err:?}"))?;
Ok(schedules)
}
fn parse_gateway_gating_config_path(path: &str) -> Result<GatewayGatingConfig, String> {
let value = read_gateway_config_file(path, "gating_config_path")?;
let actions = value
.get("actions")
.and_then(Value::as_object)
.ok_or_else(|| {
"runtime_options.gating_config_path must define an actions object".to_string()
})?;
let mut action_risk_tiers = HashMap::new();
for (action, config) in actions {
let risk_tier = config.as_str().or_else(|| {
config
.as_object()
.and_then(|object| object.get("risk_tier"))
.and_then(Value::as_str)
});
let risk_tier = risk_tier.ok_or_else(|| {
format!("runtime_options.gating_config_path action '{action}' must define risk_tier")
})?;
let normalized = risk_tier.trim().to_ascii_lowercase();
if !matches!(normalized.as_str(), "r0" | "r1" | "r2" | "r3") {
return Err(format!(
"runtime_options.gating_config_path action '{action}' has unsupported risk_tier '{risk_tier}'"
));
}
action_risk_tiers.insert(action.trim().to_string(), normalized);
}
Ok(GatewayGatingConfig { action_risk_tiers })
}
fn parse_gateway_event_log_config(value: &Value) -> Result<EventLogConfig, String> {
let object = value
.as_object()
.ok_or_else(|| "runtime_options.event_log must be a JSON object".to_string())?;
let storage = object
.get("storage")
.and_then(Value::as_str)
.ok_or_else(|| {
"runtime_options.event_log.storage must be 'memory' or 'null'".to_string()
})?;
let store: Box<dyn EventLogStore> = match storage {
"memory" | "in_memory" => Box::new(InMemoryEventLogStore::default()),
"null" => Box::new(meerkat_mobkit::unified_runtime::NullEventLogStore),
other => {
return Err(format!(
"unsupported runtime_options.event_log.storage '{other}'"
));
}
};
let batch_size = object
.get("batch_size")
.and_then(Value::as_u64)
.and_then(|value| usize::try_from(value).ok())
.unwrap_or(64);
let flush_interval_ms = match object.get("flush_interval_ms") {
Some(value) => {
let flush_interval_ms = value.as_u64().ok_or_else(|| {
"runtime_options.event_log.flush_interval_ms must be a positive integer".to_string()
})?;
if flush_interval_ms == 0 {
return Err(
"runtime_options.event_log.flush_interval_ms must be greater than zero"
.to_string(),
);
}
flush_interval_ms
}
None => 1_000,
};
Ok(EventLogConfig {
store,
filter: None,
batch_size,
flush_interval: Duration::from_millis(flush_interval_ms),
})
}
fn gateway_event_log_slot(
options: &GatewayRuntimeOptions,
) -> meerkat_mobkit::storage_health::StorageSlotSummary {
if options.event_log.is_some() {
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"event_log",
"in-process store",
"declared via runtime_options.event_log ('memory' retains a bounded queryable \
buffer; 'null' drops events explicitly)",
)
} else {
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"event_log",
"not configured",
"operational events are not ingested; set runtime_options.event_log to declare a \
store explicitly",
)
}
}
fn agent_memory_census_slot(
agent_memory: &GatewayAgentMemoryOptions,
) -> meerkat_mobkit::storage_health::StorageSlotSummary {
match agent_memory.store {
GatewayAgentMemoryStoreKind::Sqlite => {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"agent-memory",
"SqliteAgentMemoryStore",
)
}
GatewayAgentMemoryStoreKind::Markdown => {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"agent-memory",
"MarkdownAgentMemoryStore",
)
}
}
}
fn parse_gateway_runtime_store_config(value: &Value) -> Result<bool, String> {
let object = value
.as_object()
.ok_or_else(|| "runtime_options.runtime_store must be a JSON object".to_string())?;
let storage = object
.get("storage")
.and_then(Value::as_str)
.ok_or_else(|| "runtime_options.runtime_store.storage must be 'memory'".to_string())?;
if !matches!(storage, "memory" | "in_memory") {
return Err(format!(
"unsupported runtime_options.runtime_store.storage '{storage}' (persistent SQLite \
is the default; the only declaration is 'memory')"
));
}
Ok(true)
}
fn parse_gateway_auth_config(value: &Value) -> Result<RuntimeDecisionState, String> {
let object = value
.as_object()
.ok_or_else(|| "runtime_options.auth_config must be a JSON object".to_string())?;
let provider = object
.get("provider")
.and_then(Value::as_str)
.or_else(|| {
if object.contains_key("sharedSecret") || object.contains_key("shared_secret") {
Some("jwt")
} else {
None
}
})
.ok_or_else(|| "runtime_options.auth_config.provider is required".to_string())?;
if provider != "jwt" {
return Err(format!(
"unsupported runtime_options.auth_config.provider '{provider}'"
));
}
let shared_secret = object
.get("shared_secret")
.or_else(|| object.get("sharedSecret"))
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"runtime_options.auth_config.shared_secret must be a non-empty string".to_string()
})?;
let issuer = object
.get("issuer")
.and_then(Value::as_str)
.unwrap_or("http://127.0.0.1/mobkit-gateway");
let audience = object
.get("audience")
.and_then(Value::as_str)
.unwrap_or("persistent-gateway");
let email_allowlist = object
.get("email_allowlist")
.or_else(|| object.get("emailAllowlist"))
.and_then(Value::as_array)
.map(|values| {
values
.iter()
.filter_map(Value::as_str)
.map(ToString::to_string)
.collect::<Vec<_>>()
})
.unwrap_or_default();
let key = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(shared_secret.as_bytes());
let discovery_json = serde_json::to_string(&json!({
"issuer": issuer,
"jwks_uri": "http://127.0.0.1/mobkit-gateway/jwks.json"
}))
.map_err(|err| format!("failed to build trusted OIDC discovery: {err}"))?;
let jwks_json = serde_json::to_string(&json!({
"keys": [{
"kty": "oct",
"alg": "HS256",
"k": key
}]
}))
.map_err(|err| format!("failed to build trusted JWKS: {err}"))?;
Ok(RuntimeDecisionState {
bigquery: BigQueryNaming {
dataset: "default_dataset".to_string(),
table: "default_table".to_string(),
},
modules: vec![],
auth: AuthPolicy {
default_provider: AuthProvider::GenericOidc,
email_allowlist,
},
trusted_oidc: TrustedOidcRuntimeConfig {
discovery_json,
jwks_json,
audience: audience.to_string(),
},
console: ConsolePolicy {
require_app_auth: true,
..ConsolePolicy::default()
},
ops: RuntimeOpsPolicy::default(),
release_metadata: ReleaseMetadata {
targets: vec![
"crates.io".to_string(),
"npm".to_string(),
"pypi".to_string(),
"github-releases".to_string(),
],
support_matrix: "lts".to_string(),
},
})
}
fn apply_gateway_runtime_config_to_request(
request_line: &str,
schedules: &[ScheduleDefinition],
gating: &GatewayGatingConfig,
) -> String {
let Ok(mut request) = serde_json::from_str::<Value>(request_line) else {
return request_line.to_string();
};
let method = request.get("method").and_then(Value::as_str).unwrap_or("");
match method {
"mobkit/scheduling/evaluate" | "mobkit/scheduling/dispatch" if !schedules.is_empty() => {
let params = request.get_mut("params").and_then(Value::as_object_mut);
if let Some(params) = params
&& !params.contains_key("schedules")
{
params.insert(
"schedules".to_string(),
serde_json::to_value(schedules).unwrap_or(Value::Null),
);
}
}
"mobkit/gating/evaluate" => {
let params = request.get_mut("params").and_then(Value::as_object_mut);
if let Some(params) = params
&& !params.contains_key("risk_tier")
&& let Some(action) = params.get("action").and_then(Value::as_str)
&& let Some(risk_tier) = gating.action_risk_tiers.get(action.trim())
{
params.insert("risk_tier".to_string(), Value::String(risk_tier.clone()));
}
}
_ => {}
}
serde_json::to_string(&request).unwrap_or_else(|_| request_line.to_string())
}
const LEGACY_MEMORY_LEDGER_STATE_FILE: &str = "elephant-memory-state.json";
const MEMORY_LEDGER_STATE_FILE: &str = "memory-ledger-state.json";
fn parse_gateway_memory_config(
memory_config: &Value,
persistent_state: Option<&std::path::Path>,
) -> Result<MemoryBackendConfig, String> {
let object = memory_config
.as_object()
.ok_or_else(|| "runtime_options.memory_config must be a JSON object".to_string())?;
let backend = object.get("backend").and_then(Value::as_str).ok_or_else(|| {
"runtime_options.memory_config.backend must be 'local_json' (or the deprecated 'elephant')"
.to_string()
})?;
let persistent_state = persistent_state.ok_or_else(|| {
"runtime_options.memory_config requires persistent_state so the memory ledger has a stable path"
.to_string()
})?;
match backend {
"local_json" => {
let health_check_endpoint = match object.get("health_check_endpoint") {
None => None,
Some(value) => Some(
value
.as_str()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"runtime_options.memory_config.health_check_endpoint must be a non-empty string when provided"
.to_string()
})?
.to_string(),
),
};
let unsupported = object
.keys()
.filter(|key| key.as_str() != "backend" && key.as_str() != "health_check_endpoint")
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.memory_config fields: {}",
unsupported.join(", ")
));
}
let new_path = persistent_state.join(MEMORY_LEDGER_STATE_FILE);
let legacy_path = persistent_state.join(LEGACY_MEMORY_LEDGER_STATE_FILE);
let state_path = if !new_path.exists() && legacy_path.exists() {
legacy_path
} else {
new_path
};
Ok(MemoryBackendConfig::LocalJson(
LocalJsonMemoryBackendConfig {
state_path: state_path.to_string_lossy().to_string(),
health_check_endpoint,
},
))
}
"elephant" => {
let endpoint = object
.get("endpoint")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"runtime_options.memory_config.endpoint must be a non-empty string".to_string()
})?;
let unsupported = object
.keys()
.filter(|key| key.as_str() != "backend" && key.as_str() != "endpoint")
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.memory_config fields: {}",
unsupported.join(", ")
));
}
eprintln!(
"[mobkit-gateway] runtime_options.memory_config.backend 'elephant' is deprecated: \
it only health-checks the endpoint and persists the ledger as local JSON; use \
backend 'local_json' with an optional health_check_endpoint"
);
let state_path = persistent_state.join(LEGACY_MEMORY_LEDGER_STATE_FILE);
Ok(MemoryBackendConfig::LocalJson(
LocalJsonMemoryBackendConfig {
state_path: state_path.to_string_lossy().to_string(),
health_check_endpoint: Some(endpoint.to_string()),
},
))
}
other => Err(format!(
"unsupported runtime_options.memory_config.backend '{other}'"
)),
}
}
fn resolve_agent_memory_root(
persistent_state: &std::path::Path,
) -> Result<std::path::PathBuf, String> {
meerkat_mobkit::MobKitStorageLayout::with_injected_roots(persistent_state.to_path_buf(), None)
.agent_memory_root()
.map(|resolved| resolved.path)
.map_err(|e| e.to_string())
}
fn parse_gateway_agent_memory_config(
agent_memory: &Value,
persistent_state: Option<&std::path::Path>,
) -> Result<Option<GatewayAgentMemoryOptions>, String> {
if let Some(enabled) = agent_memory.as_bool() {
if !enabled {
return Ok(None);
}
let path = resolve_agent_memory_root(persistent_state.ok_or_else(|| {
"runtime_options.agent_memory=true requires persistent_state".to_string()
})?)?;
return Ok(Some(GatewayAgentMemoryOptions {
config: meerkat_mobkit::AgentMemoryConfig::default(),
path,
store: GatewayAgentMemoryStoreKind::default(),
selector: None,
distiller: meerkat_mobkit::memory::distiller::DistillerConfig::default(),
steward: meerkat_mobkit::memory::steward::StewardConfig::default(),
hygienist: meerkat_mobkit::memory::hygienist::HygienistConfig::default(),
}));
}
let object = agent_memory
.as_object()
.ok_or_else(|| "runtime_options.agent_memory must be a boolean or object".to_string())?;
let supported = [
"enabled",
"realm",
"selection",
"max_entries",
"recall_timeout_ms",
"recall_failure_policy",
"instruction_header",
"per_turn_injection",
"defang_inbound",
"store",
"llm_writes",
"recorder_tool",
"content_trust",
"selector",
"distiller",
"steward",
"operator_scope",
"hygienist",
];
let unsupported = object
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.agent_memory fields: {}",
unsupported.join(", ")
));
}
if let Some(enabled) = object.get("enabled") {
let enabled = enabled
.as_bool()
.ok_or_else(|| "runtime_options.agent_memory.enabled must be a boolean".to_string())?;
if !enabled {
return Ok(None);
}
}
let realm = object
.get("realm")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("default")
.to_string();
let selection = match object
.get("selection")
.and_then(Value::as_str)
.map(str::trim)
.unwrap_or("contextual")
{
"always" => meerkat_mobkit::AgentMemorySelection::Always,
"contextual" => meerkat_mobkit::AgentMemorySelection::Contextual,
other => {
return Err(format!(
"runtime_options.agent_memory.selection must be 'always' or 'contextual' (got '{other}')"
));
}
};
let max_entries = match object.get("max_entries") {
None => 8,
Some(value) => {
let Some(value) = value.as_u64() else {
return Err(
"runtime_options.agent_memory.max_entries must be a positive integer"
.to_string(),
);
};
if value == 0 || value > 64 {
return Err(
"runtime_options.agent_memory.max_entries must be between 1 and 64".to_string(),
);
}
value as usize
}
};
let recall_timeout_ms = match object.get("recall_timeout_ms") {
None => 500,
Some(value) => {
let Some(value) = value.as_u64() else {
return Err(
"runtime_options.agent_memory.recall_timeout_ms must be a positive integer"
.to_string(),
);
};
if value == 0 || value > 30_000 {
return Err(
"runtime_options.agent_memory.recall_timeout_ms must be between 1 and 30000"
.to_string(),
);
}
value
}
};
let recall_failure_policy = match object
.get("recall_failure_policy")
.and_then(Value::as_str)
.map(str::trim)
.unwrap_or("skip")
{
"skip" => meerkat_mobkit::AgentMemoryRecallFailurePolicy::Skip,
"fail" => meerkat_mobkit::AgentMemoryRecallFailurePolicy::Fail,
other => {
return Err(format!(
"runtime_options.agent_memory.recall_failure_policy must be 'skip' or 'fail' (got '{other}')"
));
}
};
let instruction_header = match object.get("instruction_header") {
None => None,
Some(value) => Some(
value
.as_str()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"runtime_options.agent_memory.instruction_header must be a non-empty string"
.to_string()
})?
.to_string(),
),
};
let per_turn_injection = match object
.get("per_turn_injection")
.and_then(Value::as_str)
.map(str::trim)
.unwrap_or("off")
{
"off" => meerkat_mobkit::AgentMemoryPerTurnInjection::Off,
"budgeted" => meerkat_mobkit::AgentMemoryPerTurnInjection::Budgeted,
other => {
return Err(format!(
"runtime_options.agent_memory.per_turn_injection must be 'off' or 'budgeted' (got '{other}')"
));
}
};
let defang_inbound = match object.get("defang_inbound") {
None => true,
Some(value) => value.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.defang_inbound must be a boolean".to_string()
})?,
};
let store = match object
.get("store")
.and_then(Value::as_str)
.map(str::trim)
.unwrap_or("sqlite")
{
"markdown" => GatewayAgentMemoryStoreKind::Markdown,
"sqlite" => GatewayAgentMemoryStoreKind::Sqlite,
other => {
return Err(format!(
"runtime_options.agent_memory.store must be 'markdown' or 'sqlite' (got '{other}')"
));
}
};
let llm_writes = match object.get("llm_writes") {
None => meerkat_mobkit::AgentMemoryLlmWrites::Observed,
Some(value) => match value.as_str().map(str::trim) {
Some("observed") => meerkat_mobkit::AgentMemoryLlmWrites::Observed,
Some("quarantined") => meerkat_mobkit::AgentMemoryLlmWrites::Quarantined,
_ => {
return Err(format!(
"runtime_options.agent_memory.llm_writes must be 'observed' or 'quarantined' \
(got '{value}')"
));
}
},
};
let recorder_tool = match object.get("recorder_tool") {
None => true,
Some(value) => value.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.recorder_tool must be a boolean".to_string()
})?,
};
let content_trust = match object.get("content_trust") {
None => meerkat_mobkit::ContentTrustConfig::default(),
Some(value) => meerkat_mobkit::ContentTrustConfig::from_json_value(value)
.map_err(|err| format!("runtime_options.agent_memory.{err}"))?,
};
let operator_scope = match object.get("operator_scope") {
None => meerkat_mobkit::AgentMemoryOperatorScope::Off,
Some(value) => match value.as_str().map(str::trim) {
Some("off") => meerkat_mobkit::AgentMemoryOperatorScope::Off,
Some("provisional") => meerkat_mobkit::AgentMemoryOperatorScope::Provisional,
_ => {
return Err(format!(
"runtime_options.agent_memory.operator_scope must be 'off' or \
'provisional' (got '{value}')"
));
}
},
};
let selector = match object.get("selector") {
None => None,
Some(value) => {
let value = value.as_str().map(str::trim).ok_or_else(|| {
"runtime_options.agent_memory.selector must be a string \
('off', 'default', or 'profile:<path>')"
.to_string()
})?;
match value {
"off" => Some(meerkat_mobkit::memory::selector::SelectorSpec::Off),
"default" => Some(meerkat_mobkit::memory::selector::SelectorSpec::Default),
other => match other.strip_prefix("profile:") {
Some(path) if !path.trim().is_empty() => {
Some(meerkat_mobkit::memory::selector::SelectorSpec::Profile(
std::path::PathBuf::from(path.trim()),
))
}
_ => {
return Err(format!(
"runtime_options.agent_memory.selector must be 'off', 'default', \
or 'profile:<path>' (got '{other}'); this option overrides the \
MOBKIT_AGENT_MEMORY_SELECTOR environment variable"
));
}
},
}
}
};
let distiller = match object.get("distiller") {
None => meerkat_mobkit::memory::distiller::DistillerConfig::default(),
Some(value) => parse_gateway_distiller_config(value)?,
};
let steward = match object.get("steward") {
None => meerkat_mobkit::memory::steward::StewardConfig::default(),
Some(value) => parse_gateway_steward_config(value)?,
};
let hygienist = match object.get("hygienist") {
None => meerkat_mobkit::memory::hygienist::HygienistConfig::default(),
Some(value) => parse_gateway_hygienist_config(value)?,
};
if store == GatewayAgentMemoryStoreKind::Markdown
&& (llm_writes != meerkat_mobkit::AgentMemoryLlmWrites::Observed
|| object.contains_key("content_trust"))
{
return Err(
"runtime_options.agent_memory.llm_writes/content_trust require store='sqlite'"
.to_string(),
);
}
if store == GatewayAgentMemoryStoreKind::Markdown
&& selector
.as_ref()
.is_some_and(|spec| *spec != meerkat_mobkit::memory::selector::SelectorSpec::Off)
{
return Err("runtime_options.agent_memory.selector requires store='sqlite'".to_string());
}
if store == GatewayAgentMemoryStoreKind::Markdown && distiller.enabled {
return Err("runtime_options.agent_memory.distiller requires store='sqlite'".to_string());
}
if store == GatewayAgentMemoryStoreKind::Markdown && steward.enabled {
return Err("runtime_options.agent_memory.steward requires store='sqlite'".to_string());
}
if store == GatewayAgentMemoryStoreKind::Markdown && hygienist.enabled {
return Err("runtime_options.agent_memory.hygienist requires store='sqlite'".to_string());
}
if store == GatewayAgentMemoryStoreKind::Markdown
&& operator_scope != meerkat_mobkit::AgentMemoryOperatorScope::Off
{
return Err(
"runtime_options.agent_memory.operator_scope requires store='sqlite'".to_string(),
);
}
if store == GatewayAgentMemoryStoreKind::Markdown
&& per_turn_injection == meerkat_mobkit::AgentMemoryPerTurnInjection::Budgeted
{
return Err(
"runtime_options.agent_memory.per_turn_injection='budgeted' requires store='sqlite'"
.to_string(),
);
}
let path =
resolve_agent_memory_root(persistent_state.ok_or_else(|| {
"runtime_options.agent_memory requires persistent_state".to_string()
})?)?;
Ok(Some(GatewayAgentMemoryOptions {
config: meerkat_mobkit::AgentMemoryConfig {
realm,
selection,
max_entries,
recall_timeout_ms,
recall_failure_policy,
instruction_header,
per_turn_injection,
defang_inbound,
llm_writes,
recorder_tool,
content_trust,
operator_scope,
},
path,
store,
selector,
distiller,
steward,
hygienist,
}))
}
fn parse_gateway_distiller_config(
value: &Value,
) -> Result<meerkat_mobkit::memory::distiller::DistillerConfig, String> {
let mut config = meerkat_mobkit::memory::distiller::DistillerConfig::default();
if let Some(enabled) = value.as_bool() {
config.enabled = enabled;
return Ok(config);
}
let object = value.as_object().ok_or_else(|| {
"runtime_options.agent_memory.distiller must be a boolean or object".to_string()
})?;
let supported = ["enabled", "runs_per_hour", "min_interactions", "model"];
let unsupported = object
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.agent_memory.distiller fields: {}",
unsupported.join(", ")
));
}
if let Some(enabled) = object.get("enabled") {
config.enabled = enabled.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.distiller.enabled must be a boolean".to_string()
})?;
} else {
config.enabled = true;
}
if let Some(value) = object.get("runs_per_hour") {
let runs = value.as_u64().ok_or_else(|| {
"runtime_options.agent_memory.distiller.runs_per_hour must be a positive integer"
.to_string()
})?;
if runs == 0 || runs > 240 {
return Err(
"runtime_options.agent_memory.distiller.runs_per_hour must be between 1 and 240"
.to_string(),
);
}
config.runs_per_hour = runs as u32;
}
if let Some(value) = object.get("min_interactions") {
let min = value.as_u64().ok_or_else(|| {
"runtime_options.agent_memory.distiller.min_interactions must be a positive integer"
.to_string()
})?;
if min == 0 || min > 100 {
return Err(
"runtime_options.agent_memory.distiller.min_interactions must be between 1 and 100"
.to_string(),
);
}
config.min_interactions = min as u32;
}
if let Some(value) = object.get("model") {
let model = value
.as_str()
.map(str::trim)
.filter(|model| !model.is_empty())
.ok_or_else(|| {
"runtime_options.agent_memory.distiller.model must be a non-empty string"
.to_string()
})?;
config.model = Some(model.to_string());
}
Ok(config)
}
fn parse_gateway_steward_config(
value: &Value,
) -> Result<meerkat_mobkit::memory::steward::StewardConfig, String> {
let mut config = meerkat_mobkit::memory::steward::StewardConfig::default();
if let Some(enabled) = value.as_bool() {
config.enabled = enabled;
return Ok(config);
}
let object = value.as_object().ok_or_else(|| {
"runtime_options.agent_memory.steward must be a boolean or object".to_string()
})?;
let supported = [
"enabled",
"cadence",
"model",
"per_mob",
"runs_per_day",
"min_signals",
];
let unsupported = object
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.agent_memory.steward fields: {}",
unsupported.join(", ")
));
}
if let Some(enabled) = object.get("enabled") {
config.enabled = enabled.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.steward.enabled must be a boolean".to_string()
})?;
} else {
config.enabled = true;
}
if let Some(value) = object.get("cadence") {
let cadence = value
.as_str()
.map(str::trim)
.filter(|cadence| !cadence.is_empty())
.ok_or_else(|| {
"runtime_options.agent_memory.steward.cadence must be a non-empty string"
.to_string()
})?;
meerkat_mobkit::memory::steward::StewardConfig::parse_cadence(cadence)
.map_err(|err| format!("runtime_options.agent_memory.steward.cadence: {err}"))?;
config.cadence = cadence.to_string();
}
if let Some(value) = object.get("model") {
let model = value
.as_str()
.map(str::trim)
.filter(|model| !model.is_empty())
.ok_or_else(|| {
"runtime_options.agent_memory.steward.model must be a non-empty string".to_string()
})?;
config.model = Some(model.to_string());
}
if let Some(value) = object.get("per_mob") {
config.per_mob = value.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.steward.per_mob must be a boolean".to_string()
})?;
}
if let Some(value) = object.get("runs_per_day") {
let runs = value.as_u64().ok_or_else(|| {
"runtime_options.agent_memory.steward.runs_per_day must be a positive integer"
.to_string()
})?;
if runs == 0 || runs > 96 {
return Err(
"runtime_options.agent_memory.steward.runs_per_day must be between 1 and 96"
.to_string(),
);
}
config.runs_per_day = runs as u32;
}
if let Some(value) = object.get("min_signals") {
let min = value.as_u64().ok_or_else(|| {
"runtime_options.agent_memory.steward.min_signals must be a positive integer"
.to_string()
})?;
if min == 0 || min > 1000 {
return Err(
"runtime_options.agent_memory.steward.min_signals must be between 1 and 1000"
.to_string(),
);
}
config.min_signals = min as u32;
}
Ok(config)
}
fn parse_gateway_hygienist_config(
value: &Value,
) -> Result<meerkat_mobkit::memory::hygienist::HygienistConfig, String> {
let mut config = meerkat_mobkit::memory::hygienist::HygienistConfig::default();
if let Some(enabled) = value.as_bool() {
config.enabled = enabled;
return Ok(config);
}
let object = value.as_object().ok_or_else(|| {
"runtime_options.agent_memory.hygienist must be a boolean or object".to_string()
})?;
let supported = ["enabled", "runs_per_day", "model"];
let unsupported = object
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported runtime_options.agent_memory.hygienist fields: {}",
unsupported.join(", ")
));
}
if let Some(enabled) = object.get("enabled") {
config.enabled = enabled.as_bool().ok_or_else(|| {
"runtime_options.agent_memory.hygienist.enabled must be a boolean".to_string()
})?;
} else {
config.enabled = true;
}
if let Some(value) = object.get("runs_per_day") {
let runs = value.as_u64().ok_or_else(|| {
"runtime_options.agent_memory.hygienist.runs_per_day must be a positive integer"
.to_string()
})?;
if runs == 0 || runs > 24 {
return Err(
"runtime_options.agent_memory.hygienist.runs_per_day must be between 1 and 24"
.to_string(),
);
}
config.runs_per_day = runs as u32;
}
if let Some(value) = object.get("model") {
let model = value
.as_str()
.map(str::trim)
.filter(|model| !model.is_empty())
.ok_or_else(|| {
"runtime_options.agent_memory.hygienist.model must be a non-empty string"
.to_string()
})?;
config.model = Some(model.to_string());
}
Ok(config)
}
fn resolve_selector_spec(
configured: Option<&meerkat_mobkit::memory::selector::SelectorSpec>,
) -> Result<
meerkat_mobkit::memory::selector::SelectorSpec,
meerkat_mobkit::memory::selector::SelectorError,
> {
match configured {
Some(spec) => Ok(spec.clone()),
None => meerkat_mobkit::memory::selector::spec_from_env(),
}
}
fn run_single_shot() {
let request = std::env::var("MOBKIT_RPC_REQUEST")
.expect("MOBKIT_RPC_REQUEST must be set for rpc_gateway");
let config = MobKitConfig {
modules: vec![shell_module(
"routing",
r#"printf '%s\n' '{"event_id":"evt-routing","source":"module","timestamp_ms":101,"event":{"kind":"module","module":"routing","event_type":"ready","payload":{"family":"routing","health":{"state":"healthy"},"tools":{"list_method":"routing/tools.list","representative_call":{"method":"routing/tool.call","params_schema":{"tool":"string","input":"json"}}}}}}'"#,
)],
discovery: DiscoverySpec {
namespace: "mobkit-rpc".to_string(),
modules: vec!["routing".to_string()],
},
pre_spawn: vec![],
};
let mut runtime =
start_mobkit_runtime(config, vec![], Duration::from_secs(1)).expect("runtime starts");
let response = handle_mobkit_rpc_json(&mut runtime, &request, Duration::from_secs(1));
print!("{response}");
let _ = runtime.shutdown();
}
#[derive(Clone)]
struct StdioCallbackBridge {
stdout_tx: mpsc::Sender<String>,
state: Arc<Mutex<StdioCallbackState>>,
counter: Arc<std::sync::atomic::AtomicU64>,
}
#[derive(Default)]
struct StdioCallbackState {
closed: bool,
pending: HashMap<String, oneshot::Sender<Value>>,
}
const GATEWAY_SHUTDOWN_METHOD: &str = "mobkit/shutdown";
const PROVIDER_CALLBACK_TIMEOUT: Duration = Duration::from_secs(130);
const GATEWAY_RPC_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
const GATEWAY_HTTP_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
const GATEWAY_RUNTIME_EVENT_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
const GATEWAY_RUNTIME_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(310);
const GATEWAY_STDOUT_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
const GATEWAY_SHUTDOWN_HORIZON_MS: u64 = 335_000;
#[derive(Debug)]
struct GatewayShutdownRequest {
response_id: Value,
}
fn gateway_shutdown_request(message: &Value) -> Option<GatewayShutdownRequest> {
if message.get("method").and_then(Value::as_str) != Some(GATEWAY_SHUTDOWN_METHOD) {
return None;
}
Some(GatewayShutdownRequest {
response_id: message.get("id")?.clone(),
})
}
fn gateway_shutdown_response(
response_id: Value,
runtime_shutdown: Option<&UnifiedRuntimeShutdownReport>,
) -> Value {
let runtime_cleanup_completed =
runtime_shutdown.is_some_and(UnifiedRuntimeShutdownReport::cleanup_completed);
json!({
"jsonrpc": "2.0",
"id": response_id,
"result": {
"shutdown": runtime_cleanup_completed,
"runtime_cleanup_completed": runtime_cleanup_completed
}
})
}
impl StdioCallbackBridge {
fn new(stdout_tx: mpsc::Sender<String>) -> Self {
Self {
stdout_tx,
state: Arc::new(Mutex::new(StdioCallbackState::default())),
counter: Arc::new(std::sync::atomic::AtomicU64::new(1)),
}
}
fn notify(&self, method: &str, params: Value) {
let notification = json!({
"jsonrpc": "2.0",
"method": method,
"params": params,
});
if let Ok(line) = serde_json::to_string(¬ification) {
let _ = self.stdout_tx.try_send(line);
}
}
async fn notify_reliable(&self, method: &str, params: Value) {
let notification = json!({
"jsonrpc": "2.0",
"method": method,
"params": params,
});
if let Ok(line) = serde_json::to_string(¬ification) {
if let Err(e) = self.stdout_tx.send(line).await {
eprintln!("[mobkit-gateway] failed to deliver {method}: {e}");
}
}
}
async fn call(&self, method: &str, params: Value) -> Result<Value, String> {
let id = self
.counter
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let id_str = format!("cb-{id}");
let (tx, rx) = oneshot::channel();
{
let mut state = self.state.lock().await;
if state.closed {
return Err("callback transport closed".to_string());
}
state.pending.insert(id_str.clone(), tx);
}
let request = json!({
"jsonrpc": "2.0",
"id": id_str,
"method": method,
"params": params,
});
let line = match serde_json::to_string(&request) {
Ok(l) => l,
Err(e) => {
self.state.lock().await.pending.remove(&id_str);
return Err(e.to_string());
}
};
if let Err(_) = self.stdout_tx.send(line).await {
self.state.lock().await.pending.remove(&id_str);
return Err("stdout channel closed".to_string());
}
match tokio::time::timeout(PROVIDER_CALLBACK_TIMEOUT, rx).await {
Ok(Ok(value)) => {
if let Some(error) = value.get("error") {
Err(format!(
"callback error: {}",
error
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("unknown")
))
} else {
Ok(value.get("result").cloned().unwrap_or(Value::Null))
}
}
Ok(Err(_)) => Err("callback response channel dropped".to_string()),
Err(_) => {
self.state.lock().await.pending.remove(&id_str);
Err(format!(
"callback timed out after {}s",
PROVIDER_CALLBACK_TIMEOUT.as_secs()
))
}
}
}
async fn route_callback_response(&self, msg: Value) {
let id = msg
.get("id")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if let Some(tx) = self.state.lock().await.pending.remove(&id) {
let _ = tx.send(msg);
}
}
async fn close(&self) {
let pending = {
let mut state = self.state.lock().await;
state.closed = true;
std::mem::take(&mut state.pending)
};
drop(pending);
}
}
#[async_trait]
impl meerkat_mobkit::identity_first::gateway_bridges::CallbackBridge for StdioCallbackBridge {
async fn call(&self, method: &str, params: Value) -> Result<Value, String> {
self.call(method, params).await
}
}
struct CallbackToolSpec {
name: String,
description: Option<String>,
input_schema: Option<Value>,
execution: ToolExecutionContract,
}
#[derive(Clone)]
struct DetachedCallbackJobRuntime {
realm_id: String,
store: Arc<dyn meerkat::DetachedJobStore>,
service: meerkat::DetachedJobService,
blob_store: Arc<dyn meerkat_core::BlobStore>,
runtime_inbox: Option<meerkat_runtime::RuntimeDeliveryInbox>,
delivery_service: Arc<std::sync::RwLock<Option<Arc<dyn meerkat_mob::MobSessionService>>>>,
delivery_driver_started: Arc<AtomicBool>,
monitor_recovery_completed: Arc<AtomicBool>,
monitor_shell_config: Option<meerkat_tools::builtin::shell::ShellConfig>,
monitor_managers: Arc<Mutex<HashMap<String, Arc<meerkat_tools::builtin::shell::JobManager>>>>,
callback_runners: Arc<std::sync::RwLock<BTreeSet<(String, String, String)>>>,
}
impl DetachedCallbackJobRuntime {
fn new(
realm_id: impl Into<String>,
store: Arc<dyn meerkat::DetachedJobStore>,
blob_store: Arc<dyn meerkat_core::BlobStore>,
) -> Self {
Self {
realm_id: realm_id.into(),
service: meerkat::DetachedJobService::new(Arc::clone(&store)),
store,
blob_store,
runtime_inbox: None,
delivery_service: Arc::new(std::sync::RwLock::new(None)),
delivery_driver_started: Arc::new(AtomicBool::new(false)),
monitor_recovery_completed: Arc::new(AtomicBool::new(false)),
monitor_shell_config: None,
monitor_managers: Arc::new(Mutex::new(HashMap::new())),
callback_runners: Arc::new(std::sync::RwLock::new(BTreeSet::new())),
}
}
fn with_runtime_delivery_store(
mut self,
runtime_store: Arc<dyn meerkat_runtime::RuntimeStore>,
) -> Self {
self.runtime_inbox = Some(meerkat_runtime::RuntimeDeliveryInbox::new(runtime_store));
self
}
fn with_monitor_shell(mut self, project_root: PathBuf, enabled: bool) -> Self {
if enabled {
self.monitor_shell_config =
Some(meerkat_tools::builtin::shell::ShellConfig::with_project_root(project_root));
}
self
}
fn shell_delivery_projector(
&self,
) -> Result<Arc<dyn meerkat_tools::builtin::shell::ShellJobDeliveryProjector>, String> {
let runtime_inbox = self.runtime_inbox.clone().ok_or_else(|| {
"durable monitor execution requires the persistent runtime inbox".to_string()
})?;
Ok(Arc::new(meerkat::JobOutboxProjector::new_for_realm(
Arc::clone(&self.store),
runtime_inbox,
self.realm_id.clone(),
)))
}
async fn monitor_manager(
&self,
session_id: &meerkat_core::SessionId,
) -> Result<Arc<meerkat_tools::builtin::shell::JobManager>, String> {
let key = session_id.to_string();
let mut managers = self.monitor_managers.lock().await;
if let Some(manager) = managers.get(&key) {
return Ok(Arc::clone(manager));
}
let config = self.monitor_shell_config.clone().ok_or_else(|| {
"monitors/start requires shell tooling to be enabled for this MobKit runtime"
.to_string()
})?;
let durable = meerkat_tools::builtin::shell::DurableShellJobRuntime::new(
self.realm_id.clone(),
session_id.clone(),
Arc::clone(&self.store),
Arc::clone(&self.blob_store),
self.shell_delivery_projector()?,
)
.map_err(|error| error.to_string())?;
let manager = Arc::new(
meerkat_tools::builtin::shell::JobManager::new(config)
.with_durable_job_runtime(durable)
.bind_canonical_async_ops(
session_id.clone(),
Arc::new(meerkat_runtime::RuntimeOpsLifecycleRegistry::new()),
),
);
managers.insert(key, Arc::clone(&manager));
Ok(manager)
}
fn register_callback_catalog(&self, catalog: &[ToolCatalogEntry]) -> bool {
let mut runners = self
.callback_runners
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut registered_new_runner = false;
for entry in catalog {
if let Some(policy) = entry.execution.detached_policy() {
registered_new_runner |= runners.insert((
entry.tool.name.to_string(),
policy.runner().name().to_string(),
policy.runner().version().to_string(),
));
}
}
registered_new_runner
}
fn owns_callback_job(&self, spec: &meerkat::JobSpec) -> bool {
self.callback_runners
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(&(
spec.tool.name().to_string(),
spec.runner.name().to_string(),
spec.runner.version().to_string(),
))
}
async fn recover_monitor_jobs(&self) -> Result<(), String> {
if self.monitor_shell_config.is_none()
|| self.monitor_recovery_completed.load(Ordering::Acquire)
{
return Ok(());
}
let sessions = self
.store
.list_all(usize::MAX)
.await
.map_err(|error| error.to_string())?
.into_iter()
.filter(|job| {
job.spec.realm_id == self.realm_id
&& matches!(
job.spec.runner.name(),
"meerkat.shell" | "meerkat.monitor_script"
)
})
.map(|job| {
let session_id = job.spec.origin_session_id;
(session_id.to_string(), session_id)
})
.collect::<BTreeMap<_, _>>();
for session_id in sessions.into_values() {
self.monitor_manager(&session_id)
.await?
.list_jobs()
.await
.map_err(|error| error.to_string())?;
}
self.monitor_recovery_completed
.store(true, Ordering::Release);
Ok(())
}
fn attach_delivery_service(&self, service: Arc<dyn meerkat_mob::MobSessionService>) {
*self
.delivery_service
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(service);
}
fn arm_delivery_driver(&self, unified_runtime: Arc<UnifiedRuntime>) {
if self
.delivery_driver_started
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}
let runtime = self.clone();
tokio::spawn(async move {
loop {
if let Err(error) = runtime.recover_monitor_jobs().await {
tracing::warn!(%error, "durable monitor recovery failed");
}
if let Err(error) = runtime.drain_deliveries().await {
tracing::warn!(%error, "durable callback delivery drain failed");
}
match runtime.health_projection().await {
Ok(projection) => unified_runtime.set_job_health_projection(Some(projection)),
Err(error) => {
tracing::warn!(%error, "durable callback health projection failed");
unified_runtime.set_job_health_projection(Some(json!({
"status": "degraded",
"detached_jobs": {
"status": "degraded",
"reason": "job_health_projection_failed"
}
})));
}
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
});
}
async fn drain_deliveries(&self) -> Result<(), String> {
let Some(runtime_inbox) = self.runtime_inbox.clone() else {
return Ok(());
};
let delivery_service = self
.delivery_service
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let Some(delivery_service) = delivery_service else {
return Ok(());
};
let projector = meerkat::JobOutboxProjector::new_for_realm(
Arc::clone(&self.store),
runtime_inbox.clone(),
self.realm_id.clone(),
);
projector
.project_pending(256)
.await
.map_err(|error| error.to_string())?;
let sink: Arc<dyn meerkat::JobDeliverySink> =
Arc::new(CallbackJobDeliverySink { delivery_service });
let applier = meerkat::JobRuntimeDeliveryApplier::new(runtime_inbox, sink);
for session_id in self.delivery_origin_sessions().await? {
applier
.apply_pending(
&meerkat_runtime::LogicalRuntimeId::for_session(&session_id),
256,
)
.await
.map_err(|error| error.to_string())?;
}
Ok(())
}
async fn delivery_origin_sessions(&self) -> Result<Vec<meerkat_core::SessionId>, String> {
Ok(self
.store
.list_all(usize::MAX)
.await
.map_err(|error| error.to_string())?
.into_iter()
.filter(|job| job.spec.realm_id == self.realm_id)
.map(|job| {
let origin = job.spec.origin_session_id;
(origin.to_string(), origin)
})
.collect::<BTreeMap<_, _>>()
.into_values()
.collect())
}
async fn runtime_delivery_backlog(&self) -> Result<u64, String> {
let Some(runtime_inbox) = self.runtime_inbox.clone() else {
return Ok(0);
};
let mut total = 0_u64;
for session_id in self.delivery_origin_sessions().await? {
let pending = runtime_inbox
.list_pending(
&meerkat_runtime::LogicalRuntimeId::for_session(&session_id),
usize::MAX,
)
.await
.map_err(|error| error.to_string())?;
total = total.saturating_add(u64::try_from(pending.len()).unwrap_or(u64::MAX));
}
Ok(total)
}
async fn health_projection(&self) -> Result<Value, String> {
let now_ms = callback_unix_time_ms()?;
let health = self
.service
.health_snapshot_for_realm(&self.realm_id, now_ms, usize::MAX)
.await
.map_err(|error| error.to_string())?;
let runtime_delivery_backlog = self.runtime_delivery_backlog().await?;
let delivery_backlog = health
.delivery_backlog
.saturating_add(runtime_delivery_backlog);
let degraded = health.is_degraded() || runtime_delivery_backlog > 0;
let mut by_session = serde_json::Map::new();
for job in self
.store
.list_all(usize::MAX)
.await
.map_err(|error| error.to_string())?
.into_iter()
.filter(|job| job.spec.realm_id == self.realm_id)
{
let key = job.spec.origin_session_id.to_string();
let entry = by_session.entry(key).or_insert_with(|| {
json!({
"active": 0_u64,
"awaiting_detached": false,
"queued": 0_u64,
"running": 0_u64,
"needs_attention": 0_u64
})
});
let active = matches!(
job.machine_state.lifecycle_phase,
meerkat::JobPhase::Unsubmitted
| meerkat::JobPhase::Queued
| meerkat::JobPhase::Claimed
| meerkat::JobPhase::Running
| meerkat::JobPhase::WaitingExternal
| meerkat::JobPhase::LossObserved
| meerkat::JobPhase::RetryScheduled
);
if active {
entry["active"] = json!(
entry["active"]
.as_u64()
.unwrap_or_default()
.saturating_add(1)
);
entry["awaiting_detached"] = Value::Bool(true);
}
let field = match job.machine_state.lifecycle_phase {
meerkat::JobPhase::Queued | meerkat::JobPhase::RetryScheduled => Some("queued"),
meerkat::JobPhase::Running | meerkat::JobPhase::WaitingExternal => Some("running"),
meerkat::JobPhase::NeedsAttention => Some("needs_attention"),
_ => None,
};
if let Some(field) = field {
entry[field] = json!(entry[field].as_u64().unwrap_or_default().saturating_add(1));
}
}
Ok(json!({
"status": if degraded { "degraded" } else { "ok" },
"monitors_available": self.monitor_shell_config.is_some(),
"detached_jobs": {
"status": if degraded { "degraded" } else { "ok" },
"queued": health.queued,
"running": health.running,
"awaiting_members": health.awaiting_members,
"stale_leases": health.stale_leases,
"needs_attention": health.needs_attention,
"delivery_backlog": delivery_backlog
},
"by_session": by_session
}))
}
}
struct CallbackJobDeliverySink {
delivery_service: Arc<dyn meerkat_mob::MobSessionService>,
}
#[async_trait]
impl meerkat::JobDeliverySink for CallbackJobDeliverySink {
async fn apply(&self, application: meerkat::JobDeliveryApplication) -> Result<(), String> {
match application {
meerkat::JobDeliveryApplication::Record { .. } => Ok(()),
meerkat::JobDeliveryApplication::Notification {
job_id,
delivery_sequence,
subscription,
content,
} => {
let mut request = meerkat_core::service::AppendSystemContextRequest::from_text(
callback_job_delivery_text(&job_id, &content),
);
request.source = Some(format!("detached_job:{job_id}"));
request.idempotency_key = Some(format!(
"job:{job_id}:{delivery_sequence}:{}",
subscription.subscription_id()
));
self.delivery_service
.append_system_context(subscription.session_id(), request)
.await
.map(|_| ())
.map_err(|error| error.to_string())
}
meerkat::JobDeliveryApplication::Event {
job_id,
delivery_sequence,
subscription,
interaction_lineage_id,
handling_mode,
content,
} => {
let runtime_adapter = self
.delivery_service
.runtime_adapter()
.ok_or_else(|| {
format!(
"session {} has no runtime-owned durable event ingress for detached job {job_id}",
subscription.session_id()
)
})?;
let input = callback_job_event_input(
&job_id,
delivery_sequence,
&subscription,
&interaction_lineage_id,
handling_mode,
&content,
);
meerkat_runtime::SessionServiceRuntimeExt::accept_input(
runtime_adapter.as_ref(),
subscription.session_id(),
input,
)
.await
.map(|_| ())
.map_err(|error| error.to_string())
}
}
}
}
fn callback_job_event_input(
job_id: &meerkat::JobId,
delivery_sequence: u64,
subscription: &meerkat::JobSubscription,
interaction_lineage_id: &meerkat::InteractionLineageId,
handling_mode: meerkat_core::HandlingMode,
content: &meerkat::JobDeliveryContent,
) -> meerkat_runtime::Input {
let event_type = match content {
meerkat::JobDeliveryContent::Notification(_) => "job.notification",
meerkat::JobDeliveryContent::Terminal(_) => "job.terminal",
};
let content_value = match content {
meerkat::JobDeliveryContent::Notification(notification) => json!({
"kind": "notification",
"notification": notification,
}),
meerkat::JobDeliveryContent::Terminal(result) => json!({
"kind": "terminal",
"result": result,
}),
};
let correlation_id = uuid::Uuid::parse_str(interaction_lineage_id.as_str())
.ok()
.map(meerkat_runtime::CorrelationId::from_uuid);
let idempotency_key = format!(
"job:{job_id}:{delivery_sequence}:{}",
subscription.subscription_id()
);
meerkat_runtime::Input::ExternalEvent(meerkat_runtime::ExternalEventInput {
objective_id: None,
header: meerkat_runtime::InputHeader {
id: meerkat_core::lifecycle::InputId::new(),
timestamp: chrono::Utc::now(),
source: meerkat_runtime::InputOrigin::External {
source_name: event_type.to_string(),
},
durability: meerkat_runtime::InputDurability::Durable,
visibility: meerkat_runtime::InputVisibility::default(),
idempotency_key: Some(meerkat_runtime::IdempotencyKey::new(idempotency_key)),
supersession_key: None,
correlation_id,
},
event_type: event_type.to_string(),
payload: json!({
"job_id": job_id.to_string(),
"delivery_sequence": delivery_sequence,
"content": content_value,
}),
blocks: None,
handling_mode,
render_metadata: None,
})
}
fn callback_job_delivery_text(
job_id: &meerkat::JobId,
content: &meerkat::JobDeliveryContent,
) -> String {
match content {
meerkat::JobDeliveryContent::Notification(notification) => format!(
"Detached job {job_id}: {}\n\n{}",
notification.title(),
notification.body()
),
meerkat::JobDeliveryContent::Terminal(result) => {
format!("Detached job {job_id} reached terminal state: {result:?}")
}
}
}
impl CallbackToolSpec {
fn parse(value: &Value) -> Result<Self, String> {
if let Some(name) = value.as_str() {
return Ok(Self {
name: name.to_string(),
description: None,
input_schema: None,
execution: ToolExecutionContract::default(),
});
}
let Some(object) = value.as_object() else {
return Err(format!(
"tools entries must be strings or {{name, description?, input_schema?}} \
objects, got: {value}"
));
};
let name = object
.get("name")
.and_then(Value::as_str)
.filter(|name| !name.is_empty())
.ok_or_else(|| format!("tool object requires a non-empty string name, got: {value}"))?
.to_string();
let description = match object.get("description") {
None | Some(Value::Null) => None,
Some(Value::String(text)) => Some(text.clone()),
Some(other) => {
return Err(format!(
"tool '{name}' description must be a string, got: {other}"
));
}
};
let input_schema = match object.get("input_schema") {
None | Some(Value::Null) => None,
Some(schema @ Value::Object(_)) => Some(schema.clone()),
Some(other) => {
return Err(format!(
"tool '{name}' input_schema must be a JSON object, got: {other}"
));
}
};
let execution = object
.get("execution")
.cloned()
.map(serde_json::from_value::<meerkat_contracts::CallbackToolExecution>)
.transpose()
.map_err(|error| format!("tool '{name}' execution is invalid: {error}"))?
.map_or_else(
|| Ok(ToolExecutionContract::default()),
callback_execution_contract,
)?;
Ok(Self {
name,
description,
input_schema,
execution,
})
}
}
fn callback_execution_contract(
execution: meerkat_contracts::CallbackToolExecution,
) -> Result<ToolExecutionContract, String> {
match execution {
meerkat_contracts::CallbackToolExecution::Fast => Ok(ToolExecutionContract::default()),
meerkat_contracts::CallbackToolExecution::Detached {
runner,
restart_class,
idempotency_scope,
submission_timeout_ms,
credential_scopes,
} => {
let runner = meerkat_core::RunnerIdentity::new(runner.name, runner.version)
.map_err(|error| error.to_string())?;
let restart_class = match restart_class {
meerkat_contracts::JobRestartClass::Adoptable => {
meerkat_core::RestartClass::Adoptable
}
meerkat_contracts::JobRestartClass::CheckpointResumable => {
meerkat_core::RestartClass::CheckpointResumable
}
meerkat_contracts::JobRestartClass::Replayable => {
meerkat_core::RestartClass::Replayable
}
meerkat_contracts::JobRestartClass::NonResumable => {
meerkat_core::RestartClass::NonResumable
}
};
let idempotency_scope = match idempotency_scope {
meerkat_contracts::JobIdempotencyScope::ToolCall => {
meerkat_core::IdempotencyScope::ToolCall
}
meerkat_contracts::JobIdempotencyScope::InteractionAndArguments => {
meerkat_core::IdempotencyScope::InteractionAndArguments
}
meerkat_contracts::JobIdempotencyScope::HostSemanticKey => {
meerkat_core::IdempotencyScope::HostSemanticKey
}
};
let policy = meerkat_core::DetachedToolExecutionPolicy::new(
runner,
restart_class,
idempotency_scope,
Duration::from_millis(submission_timeout_ms),
)
.map_err(|error| error.to_string())?
.with_credential_scopes(credential_scopes);
ToolExecutionContract::new(
std::collections::BTreeSet::from([ToolExecutionMode::Detached]),
ToolExecutionMode::Detached,
None,
Some(policy),
)
.map_err(|error| error.to_string())
}
}
}
#[derive(Clone)]
struct CallbackToolDispatcher {
bridge: StdioCallbackBridge,
scope_id: String,
tool_defs: Arc<[Arc<ToolDef>]>,
tool_catalog: Arc<[ToolCatalogEntry]>,
detached_jobs: Option<DetachedCallbackJobRuntime>,
reconcile_registered_catalog: bool,
}
impl CallbackToolDispatcher {
fn new(
bridge: StdioCallbackBridge,
scope_id: String,
tools: Vec<CallbackToolSpec>,
detached_jobs: Option<DetachedCallbackJobRuntime>,
) -> Self {
let entries: Vec<(Arc<ToolDef>, ToolExecutionContract)> = tools
.into_iter()
.map(|tool| {
(
Arc::new(ToolDef {
name: tool.name.into(),
description: tool
.description
.unwrap_or_else(|| "Python callback tool".to_string()),
input_schema: tool
.input_schema
.unwrap_or_else(|| json!({"type": "object"})),
provenance: None,
}),
tool.execution,
)
})
.collect();
let tool_defs = entries
.iter()
.map(|(tool, _)| Arc::clone(tool))
.collect::<Vec<_>>();
let tool_catalog = entries
.into_iter()
.map(|(tool, execution)| {
ToolCatalogEntry::session_inline(tool, true).with_execution_contract(execution)
})
.collect::<Vec<_>>();
let reconcile_registered_catalog = detached_jobs
.as_ref()
.is_some_and(|runtime| runtime.register_callback_catalog(&tool_catalog));
Self {
bridge,
scope_id,
tool_defs: tool_defs.into(),
tool_catalog: tool_catalog.into(),
detached_jobs,
reconcile_registered_catalog,
}
}
async fn submit_detached(
&self,
call: ToolCallView<'_>,
context: &meerkat_core::ToolDispatchContext,
plan: &meerkat_core::ResolvedToolExecutionPlan,
) -> Result<ToolDispatchOutcome, ToolError> {
let runtime = self.detached_jobs.clone().ok_or_else(|| {
ToolError::unavailable(
call.name,
meerkat_core::ToolUnavailableReason::ExecutionModeOwnerUnavailable,
)
})?;
let origin_session_id = context.origin_session_id().cloned().ok_or_else(|| {
ToolError::execution_failed(
"detached callback dispatch requires runtime-owned session identity".to_string(),
)
})?;
let interaction_lineage = context.interaction_lineage_id().ok_or_else(|| {
ToolError::execution_failed(
"detached callback dispatch requires runtime-owned interaction lineage".to_string(),
)
})?;
let arguments_sha256 = plan.canonical_arguments_sha256().ok_or_else(|| {
ToolError::execution_failed(
"detached callback dispatch requires a root-fenced canonical argument digest"
.to_string(),
)
})?;
let policy = match plan.kind() {
meerkat_core::ResolvedExecutionKind::Detached(policy) => policy,
_ => {
return Err(ToolError::execution_failed(
"detached callback owner received a non-detached execution plan".to_string(),
));
}
};
let arguments: Value = serde_json::from_str(call.args.get())
.map_err(|error| ToolError::invalid_arguments(call.name, error.to_string()))?;
let arguments_hash = callback_sha256(arguments_sha256);
let lineage = interaction_lineage.to_string();
let submission_key = callback_submission_key(
&runtime.realm_id,
&origin_session_id,
&lineage,
call,
policy,
&arguments_hash,
&arguments,
)?;
let specification = runtime
.blob_store
.put_artifact(
"application/vnd.meerkat.callback-arguments+json",
call.args.get(),
)
.await
.map_err(|error| {
ToolError::execution_failed(format!(
"failed to persist detached callback specification: {error}"
))
})?;
let spec = meerkat::JobSpec::new(
runtime.realm_id.clone(),
origin_session_id,
meerkat::ExecutionIntentId::from_string(format!(
"intent:{lineage}:{}:{arguments_hash}",
call.name
))
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
meerkat::InteractionLineageId::from_string(lineage)
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
meerkat::ToolIdentity::new(call.name, policy.runner().version())
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
meerkat::RunnerIdentity::new(policy.runner().name(), policy.runner().version())
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
callback_job_restart_class(policy.restart_class()),
meerkat::CanonicalArgumentsHash::new(arguments_hash)
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
meerkat::JobSubmissionKey::new(submission_key)
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
)
.with_runner_specification_ref(
meerkat::RunnerSpecificationRef::new(specification.blob_id.to_string())
.map_err(|error| ToolError::execution_failed(error.to_string()))?,
)
.with_credential_context_refs(match plan.credential_context_refs() {
meerkat_core::ToolExecutionApplicability::Applicable(references) => references.clone(),
meerkat_core::ToolExecutionApplicability::NotApplicable => Vec::new(),
});
let receipt = runtime
.service
.submit(spec)
.await
.map_err(|error| ToolError::execution_failed(error.to_string()))?;
let projected = meerkat::project_job_receipt(receipt.clone());
let bridge = self.bridge.clone();
tokio::spawn(async move {
if let Err(error) =
start_detached_callback_attempt(runtime, bridge, receipt.job_id).await
{
tracing::warn!(%error, "detached callback attempt start did not complete");
}
});
let content = serde_json::to_string(&projected).map_err(|error| {
ToolError::execution_failed(format!("failed to encode detached job receipt: {error}"))
})?;
Ok(ToolResult::new(call.id.to_string(), content, false).into())
}
fn owns_detached_callback_spec(&self, spec: &meerkat::JobSpec) -> bool {
self.tool_catalog.iter().any(|entry| {
entry.tool.name == spec.tool.name()
&& entry.execution.detached_policy().is_some_and(|policy| {
policy.runner().name() == spec.runner.name()
&& policy.runner().version() == spec.runner.version()
})
})
}
async fn reconcile_detached_jobs(&self) -> Result<(), String> {
let Some(runtime) = self.detached_jobs.clone() else {
return Ok(());
};
let jobs = runtime.store.list_all(usize::MAX).await.map_err(|error| {
format!("failed to enumerate detached callbacks for reconciliation: {error}")
})?;
let mut attempts = Vec::new();
for job in jobs.into_iter().filter(|job| {
job.spec.realm_id == runtime.realm_id && self.owns_detached_callback_spec(&job.spec)
}) {
if matches!(
job.machine_state.lifecycle_phase,
meerkat::JobPhase::Running | meerkat::JobPhase::WaitingExternal
) {
let attempt_id =
job.machine_state
.current_attempt_id
.as_deref()
.ok_or_else(|| {
format!(
"active detached callback {} has no committed attempt id",
job.job_id
)
})?;
let runner_handle =
job.machine_state.runner_handle.as_deref().ok_or_else(|| {
format!(
"active detached callback {} has no committed runner handle",
job.job_id
)
})?;
let lease_expires_at_ms =
job.machine_state.lease_expires_at_ms.ok_or_else(|| {
format!(
"active detached callback {} has no committed lease",
job.job_id
)
})?;
attempts.push(meerkat_contracts::CallbackJobReconcileAttempt {
authority: meerkat_contracts::JobAttemptAuthority {
job_id: job.job_id.to_string(),
attempt_id: attempt_id.to_string(),
fence: job.machine_state.current_fence,
},
runner: meerkat_contracts::JobRunner {
name: job.spec.runner.name().to_string(),
version: job.spec.runner.version().to_string(),
},
restart_class: callback_wire_restart_class(job.spec.restart_class),
runner_handle: runner_handle.to_string(),
checkpoint_ref: job
.machine_state
.checkpoint_ref
.as_ref()
.map(ToString::to_string),
lease_expires_at_ms,
});
}
}
if !attempts.is_empty() {
let offered_attempts = attempts.clone();
let offered = attempts
.iter()
.map(|attempt| attempt.authority.clone())
.collect::<Vec<_>>();
let result: meerkat_contracts::CallbackJobReconcileResult = serde_json::from_value(
self.bridge
.call(
"callback/job/reconcile",
serde_json::to_value(meerkat_contracts::CallbackJobReconcileParams {
attempts,
})
.map_err(|error| error.to_string())?,
)
.await?,
)
.map_err(|error| error.to_string())?;
if result
.live_attempts
.iter()
.any(|authority| !offered.contains(authority))
{
return Err(
"callback/job/reconcile returned authority that was not offered".to_string(),
);
}
let now_ms = callback_unix_time_ms()?;
for attempt in offered_attempts.iter().filter(|attempt| {
attempt.lease_expires_at_ms <= now_ms
&& !result.live_attempts.contains(&attempt.authority)
}) {
let job_id =
meerkat::JobId::new(&attempt.authority.job_id).map_err(|e| e.to_string())?;
let write = meerkat::AttemptWriteAuthority {
attempt_id: meerkat::AttemptId::new(&attempt.authority.attempt_id)
.map_err(|e| e.to_string())?,
fence: meerkat::FenceToken::new(attempt.authority.fence),
};
match runtime
.service
.observe_lease_expired(&job_id, write, now_ms)
.await
{
Ok(_) => {}
Err(
meerkat::DetachedJobError::StaleRevision { .. }
| meerkat::DetachedJobError::InvalidTransition { .. },
) => continue,
Err(error) => return Err(error.to_string()),
}
match attempt.restart_class {
meerkat_contracts::JobRestartClass::NonResumable => {
runtime
.service
.classify_worker_loss(&job_id, now_ms)
.await
.map_err(|error| error.to_string())?;
}
meerkat_contracts::JobRestartClass::CheckpointResumable
if attempt.checkpoint_ref.is_none() =>
{
runtime
.service
.mark_needs_attention(
&job_id,
now_ms,
meerkat::JobFailureCode::new(
"checkpoint_resume_missing_checkpoint",
)
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
}
meerkat_contracts::JobRestartClass::Adoptable
| meerkat_contracts::JobRestartClass::CheckpointResumable
| meerkat_contracts::JobRestartClass::Replayable => {
runtime
.service
.schedule_retry(&job_id, now_ms)
.await
.map_err(|error| error.to_string())?;
}
}
}
for attempt in offered_attempts.iter().filter(|attempt| {
attempt.lease_expires_at_ms > now_ms
|| result.live_attempts.contains(&attempt.authority)
}) {
let runtime = runtime.clone();
let bridge = self.bridge.clone();
let authority = attempt.authority.clone();
spawn_callback_lease_tracker(runtime, bridge, authority);
}
}
self.start_due_detached_jobs().await
}
async fn start_due_detached_jobs(&self) -> Result<(), String> {
let Some(runtime) = self.detached_jobs.clone() else {
return Ok(());
};
let now_ms = callback_unix_time_ms()?;
let jobs = runtime.store.list_all(usize::MAX).await.map_err(|error| {
format!("failed to enumerate detached callbacks for runnable work: {error}")
})?;
for job in jobs.into_iter().filter(|job| {
job.spec.realm_id == runtime.realm_id
&& self.owns_detached_callback_spec(&job.spec)
&& (job.machine_state.lifecycle_phase == meerkat::JobPhase::Queued
|| job.machine_state.lifecycle_phase == meerkat::JobPhase::RetryScheduled)
}) {
let runtime = runtime.clone();
let bridge = self.bridge.clone();
let delay_ms = job
.machine_state
.retry_due_at_ms
.unwrap_or(now_ms)
.saturating_sub(now_ms);
tokio::spawn(async move {
if delay_ms > 0 {
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
}
if let Err(error) =
start_detached_callback_attempt(runtime, bridge, job.job_id).await
{
tracing::warn!(%error, "runnable detached callback start did not complete");
}
});
}
Ok(())
}
}
fn spawn_callback_lease_tracker(
runtime: DetachedCallbackJobRuntime,
bridge: StdioCallbackBridge,
authority: meerkat_contracts::JobAttemptAuthority,
) {
tokio::spawn(async move {
if let Err(error) = reconcile_missing_callback_after_lease(runtime, bridge, authority).await
{
tracing::warn!(%error, "detached callback lease tracking failed");
}
});
}
async fn reconcile_missing_callback_after_lease(
runtime: DetachedCallbackJobRuntime,
bridge: StdioCallbackBridge,
authority: meerkat_contracts::JobAttemptAuthority,
) -> Result<(), String> {
let job_id = meerkat::JobId::new(&authority.job_id).map_err(|error| error.to_string())?;
loop {
let stored = runtime
.store
.get(&job_id)
.await
.map_err(|error| error.to_string())?
.ok_or_else(|| format!("detached callback {job_id} disappeared before lease expiry"))?;
let exact_attempt = stored.machine_state.current_attempt_id.as_deref()
== Some(authority.attempt_id.as_str())
&& stored.machine_state.current_fence == authority.fence
&& matches!(
stored.machine_state.lifecycle_phase,
meerkat::JobPhase::Running | meerkat::JobPhase::WaitingExternal
);
if !exact_attempt {
return Ok(());
}
let lease_expires_at_ms = stored
.machine_state
.lease_expires_at_ms
.ok_or_else(|| format!("active detached callback {job_id} has no committed lease"))?;
let now_ms = callback_unix_time_ms()?;
if lease_expires_at_ms > now_ms {
tokio::time::sleep(Duration::from_millis(
lease_expires_at_ms.saturating_sub(now_ms).saturating_add(1),
))
.await;
continue;
}
let offered = meerkat_contracts::CallbackJobReconcileAttempt {
authority: authority.clone(),
runner: meerkat_contracts::JobRunner {
name: stored.spec.runner.name().to_string(),
version: stored.spec.runner.version().to_string(),
},
restart_class: callback_wire_restart_class(stored.spec.restart_class),
runner_handle: stored
.machine_state
.runner_handle
.as_ref()
.ok_or_else(|| {
format!("active detached callback {job_id} has no committed runner handle")
})?
.clone(),
checkpoint_ref: stored
.machine_state
.checkpoint_ref
.as_ref()
.map(ToString::to_string),
lease_expires_at_ms,
};
let _live_attempts = match bridge
.call(
"callback/job/reconcile",
serde_json::to_value(meerkat_contracts::CallbackJobReconcileParams {
attempts: vec![offered],
})
.map_err(|error| error.to_string())?,
)
.await
{
Ok(value) => {
let result: meerkat_contracts::CallbackJobReconcileResult =
serde_json::from_value(value).map_err(|error| error.to_string())?;
if result
.live_attempts
.iter()
.any(|candidate| candidate != &authority)
{
return Err(
"callback/job/reconcile returned authority that was not offered"
.to_string(),
);
}
result.live_attempts
}
Err(error) => {
tracing::warn!(%error, %job_id, "host unavailable at detached callback lease boundary");
Vec::new()
}
};
let current = runtime
.store
.get(&job_id)
.await
.map_err(|error| error.to_string())?
.ok_or_else(|| format!("detached callback {job_id} disappeared during reconcile"))?;
let still_exact = current.machine_state.current_attempt_id.as_deref()
== Some(authority.attempt_id.as_str())
&& current.machine_state.current_fence == authority.fence
&& matches!(
current.machine_state.lifecycle_phase,
meerkat::JobPhase::Running | meerkat::JobPhase::WaitingExternal
);
if !still_exact {
return Ok(());
}
let current_lease = current
.machine_state
.lease_expires_at_ms
.ok_or_else(|| format!("active detached callback {job_id} has no committed lease"))?;
let now_ms = callback_unix_time_ms()?;
if current_lease > now_ms {
continue;
}
let write = meerkat::AttemptWriteAuthority {
attempt_id: meerkat::AttemptId::new(&authority.attempt_id)
.map_err(|error| error.to_string())?,
fence: meerkat::FenceToken::new(authority.fence),
};
match runtime
.service
.observe_lease_expired(&job_id, write, now_ms)
.await
{
Ok(_) => {}
Err(
meerkat::DetachedJobError::StaleRevision { .. }
| meerkat::DetachedJobError::InvalidTransition { .. },
) => return Ok(()),
Err(error) => return Err(error.to_string()),
}
let retry = match current.spec.restart_class {
meerkat::RestartClass::NonResumable => {
runtime
.service
.classify_worker_loss(&job_id, now_ms)
.await
.map_err(|error| error.to_string())?;
false
}
meerkat::RestartClass::CheckpointResumable
if current.machine_state.checkpoint_ref.is_none() =>
{
runtime
.service
.mark_needs_attention(
&job_id,
now_ms,
meerkat::JobFailureCode::new("checkpoint_resume_missing_checkpoint")
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
false
}
meerkat::RestartClass::Adoptable
| meerkat::RestartClass::CheckpointResumable
| meerkat::RestartClass::Replayable => {
runtime
.service
.schedule_retry(&job_id, now_ms)
.await
.map_err(|error| error.to_string())?;
true
}
};
if retry {
start_detached_callback_attempt(runtime, bridge, job_id).await?;
}
return Ok(());
}
}
#[async_trait]
impl AgentToolDispatcher for CallbackToolDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::clone(&self.tool_defs)
}
fn tool_catalog_capabilities(&self) -> ToolCatalogCapabilities {
ToolCatalogCapabilities {
exact_catalog: true,
may_require_catalog_control_plane: false,
}
}
fn tool_catalog(&self) -> Arc<[ToolCatalogEntry]> {
Arc::clone(&self.tool_catalog)
}
fn resolve_execution_plan(
&self,
call: ToolCallView<'_>,
_dispatch_context: &meerkat_core::ToolDispatchContext,
resolution_context: &meerkat_core::ToolExecutionResolutionContext,
) -> Result<meerkat_core::ResolvedToolExecutionPlan, meerkat_core::ToolExecutionResolutionError>
{
let entry = self
.tool_catalog
.iter()
.find(|entry| entry.tool.name == call.name)
.ok_or_else(|| meerkat_core::ToolExecutionResolutionError::NotFound {
tool_name: call.name.to_string(),
})?;
let resolution_context =
resolution_context.with_deadline(ToolDeadlineContributor::finite(
ToolDeadlineOwner::ToolInternal,
Duration::from_mins(2),
))?;
entry
.execution
.resolve_default(resolution_context.deadlines().clone())
.map_err(Into::into)
}
async fn dispatch(&self, call: ToolCallView<'_>) -> Result<ToolDispatchOutcome, ToolError> {
let args: Value =
serde_json::from_str(call.args.get()).map_err(|e| ToolError::InvalidArguments {
name: call.name.to_string(),
reason: e.to_string(),
})?;
let params = json!({
"scope_id": self.scope_id,
"tool": call.name,
"arguments": args,
});
match self.bridge.call("callback/call_tool", params).await {
Ok(result) => Ok(ToolResult {
tool_use_id: call.id.to_string(),
content:
meerkat_mobkit::identity_first::gateway_bridges::callback_result_to_content(
&result,
),
is_error: false,
}
.into()),
Err(err) => Ok(ToolResult {
tool_use_id: call.id.to_string(),
content: vec![ContentBlock::Text {
text: format!("Tool execution failed: {err}"),
}],
is_error: true,
}
.into()),
}
}
async fn dispatch_resolved_with_context(
&self,
call: ToolCallView<'_>,
context: &meerkat_core::ToolDispatchContext,
plan: &meerkat_core::ResolvedToolExecutionPlan,
) -> Result<ToolDispatchOutcome, ToolError> {
match plan.mode() {
ToolExecutionMode::Fast => self.dispatch(call).await,
ToolExecutionMode::Detached => self.submit_detached(call, context, plan).await,
ToolExecutionMode::Streaming => Err(ToolError::unavailable(
call.name,
meerkat_core::ToolUnavailableReason::ExecutionModeOwnerUnavailable,
)),
}
}
}
fn callback_sha256(digest: [u8; 32]) -> String {
let mut encoded = String::with_capacity("sha256:".len() + digest.len() * 2);
encoded.push_str("sha256:");
for byte in digest {
use std::fmt::Write as _;
let _ = write!(encoded, "{byte:02x}");
}
encoded
}
fn callback_submission_key(
realm_id: &str,
origin_session_id: &meerkat_core::SessionId,
interaction_lineage: &str,
call: ToolCallView<'_>,
policy: &meerkat_core::DetachedToolExecutionPolicy,
arguments_hash: &str,
arguments: &Value,
) -> Result<String, ToolError> {
let scope = match policy.idempotency_scope() {
meerkat_core::IdempotencyScope::ToolCall => format!("tool-call:{}", call.id),
meerkat_core::IdempotencyScope::InteractionAndArguments => {
format!("interaction:{interaction_lineage}:arguments:{arguments_hash}")
}
meerkat_core::IdempotencyScope::HostSemanticKey => {
let semantic_key = arguments
.get("idempotency_key")
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| {
ToolError::invalid_arguments(
call.name,
"host_semantic_key execution requires a non-empty string idempotency_key",
)
})?;
format!(
"host-semantic:sha256:{:x}",
Sha256::digest(semantic_key.as_bytes())
)
}
};
Ok(format!(
"callback:{realm_id}:{origin_session_id}:{}:{}:{}:{scope}",
call.name,
policy.runner().name(),
policy.runner().version(),
))
}
fn callback_job_restart_class(class: meerkat_core::RestartClass) -> meerkat::RestartClass {
match class {
meerkat_core::RestartClass::Adoptable => meerkat::RestartClass::Adoptable,
meerkat_core::RestartClass::CheckpointResumable => {
meerkat::RestartClass::CheckpointResumable
}
meerkat_core::RestartClass::Replayable => meerkat::RestartClass::Replayable,
meerkat_core::RestartClass::NonResumable => meerkat::RestartClass::NonResumable,
}
}
fn callback_wire_restart_class(class: meerkat::RestartClass) -> meerkat_contracts::JobRestartClass {
match class {
meerkat::RestartClass::Adoptable => meerkat_contracts::JobRestartClass::Adoptable,
meerkat::RestartClass::CheckpointResumable => {
meerkat_contracts::JobRestartClass::CheckpointResumable
}
meerkat::RestartClass::Replayable => meerkat_contracts::JobRestartClass::Replayable,
meerkat::RestartClass::NonResumable => meerkat_contracts::JobRestartClass::NonResumable,
}
}
async fn start_detached_callback_attempt(
runtime: DetachedCallbackJobRuntime,
bridge: StdioCallbackBridge,
job_id: meerkat::JobId,
) -> Result<(), String> {
let stored = runtime
.store
.get(&job_id)
.await
.map_err(|error| error.to_string())?
.ok_or_else(|| format!("detached callback job {job_id} disappeared after submission"))?;
let claimed_at_ms = callback_unix_time_ms()?;
let runnable = stored.machine_state.lifecycle_phase == meerkat::JobPhase::Queued
|| (stored.machine_state.lifecycle_phase == meerkat::JobPhase::RetryScheduled
&& stored
.machine_state
.retry_due_at_ms
.is_some_and(|due_at_ms| claimed_at_ms >= due_at_ms));
if !runnable {
return Ok(());
}
let lease_expires_at_ms = claimed_at_ms
.checked_add(120_000)
.ok_or_else(|| "detached callback lease timestamp overflowed".to_string())?;
let runner_handle = format!(
"callback:{job_id}:attempt:{}",
stored.machine_state.attempt_count.saturating_add(1)
);
let claim = match runtime
.service
.claim_attempt(
&job_id,
meerkat::AttemptClaim::new(
meerkat::WorkerId::new(format!("mobkit-callback:{}", std::process::id()))
.map_err(|error| error.to_string())?,
claimed_at_ms,
lease_expires_at_ms,
meerkat::RunnerHandleRef::new(runner_handle.clone())
.map_err(|error| error.to_string())?,
),
)
.await
{
Ok(claim) => claim,
Err(
meerkat::DetachedJobError::StaleRevision { .. }
| meerkat::DetachedJobError::InvalidTransition { .. },
) => return Ok(()),
Err(error) => return Err(error.to_string()),
};
let specification_ref = stored
.spec
.runner_specification_ref
.as_ref()
.ok_or_else(|| format!("detached callback job {job_id} has no runner specification"))?;
let specification = runtime
.blob_store
.get(&meerkat_core::BlobId::new(specification_ref.as_str()))
.await
.map_err(|error| error.to_string())?;
let arguments: Value =
serde_json::from_str(&specification.data).map_err(|error| error.to_string())?;
let credential_scopes = stored
.spec
.credential_context_refs
.iter()
.flat_map(|reference| match reference {
meerkat_core::ToolCredentialContextRef::OwningProfile { required_scopes }
| meerkat_core::ToolCredentialContextRef::AuthBinding {
required_scopes, ..
} => required_scopes.iter().cloned().collect::<Vec<_>>(),
})
.collect();
let params = meerkat_contracts::CallbackJobStartParams {
authority: meerkat_contracts::JobAttemptAuthority {
job_id: job_id.to_string(),
attempt_id: claim.attempt_id.to_string(),
fence: claim.fence.get(),
},
runner: meerkat_contracts::JobRunner {
name: stored.spec.runner.name().to_string(),
version: stored.spec.runner.version().to_string(),
},
restart_class: callback_wire_restart_class(stored.spec.restart_class),
runner_handle: runner_handle.clone(),
runner_specification_ref: Some(specification_ref.to_string()),
arguments,
credential_scopes,
resume_checkpoint: claim.resume_checkpoint.as_ref().map(ToString::to_string),
};
spawn_callback_lease_tracker(runtime.clone(), bridge.clone(), params.authority.clone());
let result: meerkat_contracts::CallbackJobStartResult = serde_json::from_value(
bridge
.call(
"callback/job/start",
serde_json::to_value(params).map_err(|error| error.to_string())?,
)
.await?,
)
.map_err(|error| error.to_string())?;
if !result.accepted || result.runner_handle != runner_handle {
runtime
.service
.mark_needs_attention(
&job_id,
callback_unix_time_ms().unwrap_or(claimed_at_ms),
meerkat::JobFailureCode::new(if result.accepted {
"callback_runner_handle_mismatch"
} else {
"callback_start_rejected"
})
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
}
Ok(())
}
fn callback_unix_time_ms() -> Result<u64, String> {
let millis = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_err(|error| error.to_string())?
.as_millis();
u64::try_from(millis).map_err(|_| "wall clock exceeds u64 milliseconds".to_string())
}
async fn callback_job_description(
runtime: &DetachedCallbackJobRuntime,
job_id: &meerkat::JobId,
) -> Result<meerkat::JobDescription, String> {
runtime
.service
.describe_for_realm(&runtime.realm_id, job_id)
.await
.map_err(|error| error.to_string())?
.ok_or_else(|| format!("detached job {job_id} does not exist in this realm"))
}
fn callback_write_authority(
authority: meerkat_contracts::JobAttemptAuthority,
) -> Result<(meerkat::JobId, meerkat::AttemptWriteAuthority), String> {
Ok((
meerkat::JobId::new(authority.job_id).map_err(|error| error.to_string())?,
meerkat::AttemptWriteAuthority {
attempt_id: meerkat::AttemptId::new(authority.attempt_id)
.map_err(|error| error.to_string())?,
fence: meerkat::FenceToken::new(authority.fence),
},
))
}
fn callback_job_rpc_response(id: Value, result: Result<Value, String>) -> Result<String, String> {
serde_json::to_string(&match result {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(message) => {
json!({"jsonrpc": "2.0", "id": id, "error": {"code": -32602, "message": message}})
}
})
.map_err(|error| error.to_string())
}
async fn handle_callback_job_rpc(
request_line: &str,
runtime: Option<&DetachedCallbackJobRuntime>,
bridge: &StdioCallbackBridge,
) -> Option<String> {
const METHODS: &[&str] = &[
"jobs/get",
"jobs/list",
"jobs/cancel",
"jobs/progress",
"jobs/result",
"jobs/artifacts",
"jobs/retry",
"jobs/health",
"jobs/subscribe",
"jobs/unsubscribe",
"monitors/start",
"mobkit/jobs/heartbeat",
"mobkit/jobs/progress",
"mobkit/jobs/checkpoint",
"mobkit/jobs/complete",
"mobkit/jobs/fail",
"mobkit/jobs/cancel_ack",
];
let request: Value = serde_json::from_str(request_line).ok()?;
let method = request.get("method").and_then(Value::as_str)?;
if !METHODS.contains(&method) {
return None;
}
let id = request.get("id").cloned().unwrap_or(Value::Null);
let Some(runtime) = runtime else {
return Some(
callback_job_rpc_response(
id,
Err(
"durable jobs require persistent_state; semantic detached admission is disabled"
.to_string(),
),
)
.unwrap_or_default(),
);
};
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
let result: Result<Value, String> = async {
match method {
"jobs/get" => {
let params: meerkat_contracts::JobsGetParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::JobsGetResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"jobs/list" => {
let params: meerkat_contracts::JobsListParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let session = params
.session_id
.ok_or_else(|| "jobs/list requires session_id".to_string())
.and_then(|raw| {
meerkat_core::SessionId::parse(&raw).map_err(|error| error.to_string())
})?;
let limit =
usize::try_from(params.limit.unwrap_or(100).min(1_000)).unwrap_or(1_000);
let jobs = runtime
.service
.list_descriptions_for_origin(&runtime.realm_id, &session, limit)
.await
.map_err(|error| error.to_string())?
.into_iter()
.map(meerkat::project_job_description)
.collect();
serde_json::to_value(meerkat_contracts::JobsListResult { jobs })
.map_err(|error| error.to_string())
}
"jobs/cancel" => {
let params: meerkat_contracts::JobsCancelParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let stored = runtime
.store
.get(&job_id)
.await
.map_err(|error| error.to_string())?
.filter(|job| job.spec.realm_id == runtime.realm_id)
.ok_or_else(|| format!("detached job {job_id} does not exist in this realm"))?;
if matches!(
stored.spec.runner.name(),
"meerkat.shell" | "meerkat.monitor_script"
) {
let manager = runtime
.monitor_manager(&stored.spec.origin_session_id)
.await?;
manager
.cancel_job(&meerkat_tools::builtin::shell::JobId::from_string(
job_id.to_string(),
))
.await
.map_err(|error| error.to_string())?;
} else {
let snapshot = runtime
.service
.request_cancel(&job_id)
.await
.map_err(|error| error.to_string())?;
if runtime.owns_callback_job(&stored.spec)
&& let Some(attempt_id) = snapshot.current_attempt_id.as_ref()
&& matches!(
snapshot.phase,
meerkat::JobPhase::Running | meerkat::JobPhase::WaitingExternal
)
{
let cancel: meerkat_contracts::CallbackJobCancelResult =
serde_json::from_value(
bridge
.call(
"callback/job/cancel",
serde_json::to_value(
meerkat_contracts::CallbackJobCancelParams {
authority: meerkat_contracts::JobAttemptAuthority {
job_id: job_id.to_string(),
attempt_id: attempt_id.to_string(),
fence: snapshot.current_fence.get(),
},
},
)
.map_err(|error| error.to_string())?,
)
.await?,
)
.map_err(|error| error.to_string())?;
if !cancel.accepted {
return Err(
"callback/job/cancel rejected the committed active attempt"
.to_string(),
);
}
}
}
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::JobsCancelResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"jobs/progress" => {
let params: meerkat_contracts::JobsProgressParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let job = meerkat::project_job_description(
callback_job_description(runtime, &job_id).await?,
);
serde_json::to_value(meerkat_contracts::JobsProgressResult {
job_id: job.job_id,
phase: job.phase,
progress: job.progress,
})
.map_err(|error| error.to_string())
}
"jobs/result" => {
let params: meerkat_contracts::JobsResultParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let job = meerkat::project_job_description(
callback_job_description(runtime, &job_id).await?,
);
serde_json::to_value(meerkat_contracts::JobsResultResult {
job_id: job.job_id,
phase: job.phase,
result: job.terminal_result,
})
.map_err(|error| error.to_string())
}
"jobs/artifacts" => {
let params: meerkat_contracts::JobsArtifactsParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let job = callback_job_description(runtime, &job_id).await?;
let reference = match job.terminal_result {
Some(
meerkat::JobTerminalResult::Succeeded {
result_ref: Some(reference),
}
| meerkat::JobTerminalResult::Failed {
detail_ref: Some(reference),
..
},
) => Some(reference.to_string()),
_ => None,
};
serde_json::to_value(meerkat_contracts::JobsArtifactsResult {
job_id: job_id.to_string(),
artifacts: reference
.into_iter()
.map(|reference| meerkat_contracts::JobArtifactRef { reference })
.collect(),
})
.map_err(|error| error.to_string())
}
"jobs/retry" => {
let params: meerkat_contracts::JobsRetryParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
let stored = runtime
.store
.get(&job_id)
.await
.map_err(|error| error.to_string())?
.filter(|job| job.spec.realm_id == runtime.realm_id)
.ok_or_else(|| format!("detached job {job_id} does not exist in this realm"))?;
let shell_owned = matches!(
stored.spec.runner.name(),
"meerkat.shell" | "meerkat.monitor_script"
);
runtime
.service
.schedule_retry(&job_id, params.retry_due_at_ms)
.await
.map_err(|error| error.to_string())?;
let now = callback_unix_time_ms()?;
let delay_ms = params.retry_due_at_ms.saturating_sub(now);
if shell_owned {
let manager = runtime
.monitor_manager(&stored.spec.origin_session_id)
.await?;
let public_job_id =
meerkat_tools::builtin::shell::JobId::from_string(job_id.to_string());
tokio::spawn(async move {
if delay_ms > 0 {
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
}
if let Err(error) = manager.get_status(&public_job_id).await {
tracing::warn!(%error, "durable shell/monitor retry start failed");
}
});
} else if runtime.owns_callback_job(&stored.spec) {
let runtime = (*runtime).clone();
let bridge = bridge.clone();
let retry_job_id = job_id.clone();
tokio::spawn(async move {
if delay_ms > 0 {
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
}
if let Err(error) =
start_detached_callback_attempt(runtime, bridge, retry_job_id).await
{
tracing::warn!(%error, "detached callback retry start failed");
}
});
} else {
}
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::JobsRetryResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"jobs/health" => {
let health = runtime
.service
.health_snapshot_for_realm(
&runtime.realm_id,
callback_unix_time_ms()?,
usize::MAX,
)
.await
.map_err(|error| error.to_string())?;
let runtime_delivery_backlog = runtime.runtime_delivery_backlog().await?;
let delivery_backlog = health
.delivery_backlog
.saturating_add(runtime_delivery_backlog);
serde_json::to_value(meerkat_contracts::JobsHealthResult {
detached_jobs: meerkat_contracts::JobHealthSummary {
status: if health.is_degraded() || runtime_delivery_backlog > 0 {
meerkat_contracts::JobHealthStatus::Degraded
} else {
meerkat_contracts::JobHealthStatus::Ok
},
queued: health.queued,
running: health.running,
awaiting_members: health.awaiting_members,
stale_leases: health.stale_leases,
needs_attention: health.needs_attention,
delivery_backlog,
},
})
.map_err(|error| error.to_string())
}
"jobs/subscribe" => {
let params: meerkat_contracts::JobsSubscribeParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
callback_job_description(runtime, &job_id).await?;
let delivery = match params.delivery {
meerkat_contracts::JobDeliveryKind::Record => meerkat::JobDeliveryKind::Record,
meerkat_contracts::JobDeliveryKind::Notification => {
meerkat::JobDeliveryKind::Notification
}
meerkat_contracts::JobDeliveryKind::Event { handling_mode } => {
meerkat::JobDeliveryKind::Event {
handling_mode: handling_mode.into(),
}
}
};
runtime
.service
.subscribe(
&job_id,
meerkat::JobSubscription::new(
meerkat::JobSubscriptionId::new(params.subscription_id)
.map_err(|error| error.to_string())?,
meerkat_core::SessionId::parse(¶ms.session_id)
.map_err(|error| error.to_string())?,
delivery,
),
)
.await
.map_err(|error| error.to_string())?;
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::JobsSubscribeResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"jobs/unsubscribe" => {
let params: meerkat_contracts::JobsUnsubscribeParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(params.job_id).map_err(|error| error.to_string())?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.unsubscribe(
&job_id,
&meerkat::JobSubscriptionId::new(params.subscription_id)
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::JobsUnsubscribeResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"monitors/start" => {
let params: meerkat_contracts::MonitorsStartParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let session_id = meerkat_core::SessionId::parse(¶ms.session_id)
.map_err(|error| error.to_string())?;
let mut limits = meerkat_tools::builtin::shell::MonitorProtocolLimits::default();
if let Some(value) = params.max_line_bytes {
limits.max_line_bytes = usize::try_from(value)
.map_err(|_| "max_line_bytes exceeds this host's size limit".to_string())?;
}
if let Some(value) = params.max_notifications_per_window {
limits.max_notifications_per_window = usize::try_from(value).map_err(|_| {
"max_notifications_per_window exceeds this host's size limit".to_string()
})?;
}
if let Some(value) = params.notification_window_ms {
limits.notification_window_ms = value;
}
if let Some(value) = params.max_retained_diagnostic_bytes {
limits.max_retained_diagnostic_bytes =
usize::try_from(value).map_err(|_| {
"max_retained_diagnostic_bytes exceeds this host's size limit"
.to_string()
})?;
}
let protocol = match params.protocol {
meerkat_contracts::MonitorOutputProtocol::FramedJsonl => {
meerkat_tools::builtin::shell::MonitorOutputProtocol::FramedJsonl
}
meerkat_contracts::MonitorOutputProtocol::Lines => {
meerkat_tools::builtin::shell::MonitorOutputProtocol::Lines
}
};
let restart_class = match params.restart_class {
meerkat_contracts::JobRestartClass::Adoptable => {
meerkat::RestartClass::Adoptable
}
meerkat_contracts::JobRestartClass::CheckpointResumable => {
meerkat::RestartClass::CheckpointResumable
}
meerkat_contracts::JobRestartClass::Replayable => {
meerkat::RestartClass::Replayable
}
meerkat_contracts::JobRestartClass::NonResumable => {
meerkat::RestartClass::NonResumable
}
};
let delivery = match params.delivery {
meerkat_contracts::JobDeliveryKind::Record => meerkat::JobDeliveryKind::Record,
meerkat_contracts::JobDeliveryKind::Notification => {
meerkat::JobDeliveryKind::Notification
}
meerkat_contracts::JobDeliveryKind::Event { handling_mode } => {
meerkat::JobDeliveryKind::Event {
handling_mode: handling_mode.into(),
}
}
};
let manager = runtime.monitor_manager(&session_id).await?;
let job_id = manager
.spawn_monitor_for_call(
¶ms.command,
params.working_dir.as_deref().map(std::path::Path::new),
params.timeout_secs,
¶ms.submission_key,
meerkat_tools::builtin::shell::MonitorStartOptions {
protocol,
restart_class,
limits,
delivery,
},
)
.await
.map_err(|error| error.to_string())?;
let job_id =
meerkat::JobId::new(job_id.to_string()).map_err(|error| error.to_string())?;
let job = callback_job_description(runtime, &job_id).await?;
serde_json::to_value(meerkat_contracts::MonitorsStartResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
"mobkit/jobs/heartbeat" => {
let params: meerkat_contracts::MobkitJobHeartbeatParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.renew_lease(
&job_id,
write,
params.heartbeat_at_ms,
params.lease_expires_at_ms,
)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
"mobkit/jobs/progress" => {
let params: meerkat_contracts::MobkitJobProgressParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.report_progress(
&job_id,
write,
meerkat::JobProgress::new(params.cursor, params.detail)
.map_err(|error| error.to_string())?,
params.observed_at_ms,
)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
"mobkit/jobs/checkpoint" => {
let params: meerkat_contracts::MobkitJobCheckpointParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.record_checkpoint(
&job_id,
write,
meerkat::CheckpointRef::new(params.checkpoint_ref)
.map_err(|error| error.to_string())?,
params.observed_at_ms,
)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
"mobkit/jobs/complete" => {
let params: meerkat_contracts::MobkitJobCompleteParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.complete_attempt(
&job_id,
write,
params.completed_at_ms,
params
.result_ref
.map(meerkat::JobResultRef::new)
.transpose()
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
"mobkit/jobs/fail" => {
let params: meerkat_contracts::MobkitJobFailParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.fail_attempt(
&job_id,
write,
params.failed_at_ms,
meerkat::JobFailureCode::new(params.code)
.map_err(|error| error.to_string())?,
params
.detail_ref
.map(meerkat::JobResultRef::new)
.transpose()
.map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
"mobkit/jobs/cancel_ack" => {
let params: meerkat_contracts::MobkitJobCancelAckParams =
serde_json::from_value(params).map_err(|error| error.to_string())?;
let (job_id, write) = callback_write_authority(params.authority)?;
callback_job_description(runtime, &job_id).await?;
runtime
.service
.acknowledge_cancel(&job_id, write, params.acknowledged_at_ms)
.await
.map_err(|error| error.to_string())?;
callback_mutation_projection(runtime, &job_id).await
}
_ => unreachable!("method allowlist and match must stay in sync"),
}
}
.await;
Some(callback_job_rpc_response(id, result).unwrap_or_default())
}
async fn callback_mutation_projection(
runtime: &DetachedCallbackJobRuntime,
job_id: &meerkat::JobId,
) -> Result<Value, String> {
let job = callback_job_description(runtime, job_id).await?;
serde_json::to_value(meerkat_contracts::MobkitJobMutationResult {
job: meerkat::project_job_description(job),
})
.map_err(|error| error.to_string())
}
struct StdioCallbackAgentBuilder {
inner: FactoryAgentBuilder,
bridge: StdioCallbackBridge,
has_session_builder: bool,
session_store: Option<Arc<dyn meerkat::SessionStore>>,
detached_jobs: Option<DetachedCallbackJobRuntime>,
}
fn callback_build_agent_options(req: &CreateSessionRequest, scope_id: &str) -> Value {
let request_labels = req.labels.as_ref();
let build_labels = req
.build
.as_ref()
.and_then(|build| build.peer_meta.as_ref())
.map(|meta| &meta.labels);
let mut labels = BTreeMap::new();
if let Some(build_labels) = build_labels {
labels.extend(
build_labels
.iter()
.map(|(key, value)| (key.clone(), value.clone())),
);
}
if let Some(request_labels) = request_labels {
labels.extend(
request_labels
.iter()
.map(|(key, value)| (key.clone(), value.clone())),
);
}
let labels = (!labels.is_empty()).then_some(labels);
let profile_name = build_labels
.and_then(|labels| labels.get("profile_name").or_else(|| labels.get("role")))
.or_else(|| {
request_labels
.and_then(|labels| labels.get("profile_name").or_else(|| labels.get("role")))
});
json!({
"scope_id": scope_id,
"session_id": labels.as_ref().and_then(|l| l.get("session_id")),
"profile_name": profile_name,
"model": &req.model,
"prompt": &req.prompt,
"labels": &labels,
"app_context": req.build.as_ref()
.and_then(|b| b.app_context.as_ref()),
})
}
#[async_trait]
impl SessionAgentBuilder for StdioCallbackAgentBuilder {
type Agent = FactoryAgent;
async fn abort_absent_session_compaction_stages(
&self,
session_id: &meerkat_core::SessionId,
) -> Result<(), SessionError> {
self.inner
.abort_absent_session_compaction_stages(session_id)
.await
}
async fn build_agent(
&self,
req: &CreateSessionRequest,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<Self::Agent, SessionError> {
if !self.has_session_builder {
let mut normalized_req = CreateSessionRequest {
model: req.model.clone(),
prompt: req.prompt.clone(),
system_prompt: req.system_prompt.clone(),
max_tokens: req.max_tokens,
event_tx: req.event_tx.clone(),
initial_turn: req.initial_turn.clone(),
build: req.build.clone(),
labels: req.labels.clone(),
deferred_prompt_policy: req.deferred_prompt_policy,
injected_context: req.injected_context.clone(),
};
ensure_shell_tooling_build_substrate(&mut normalized_req);
return self.inner.build_agent(&normalized_req, event_tx).await;
}
let scope_id = format!(
"build-{}",
self.bridge
.counter
.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
);
let options = callback_build_agent_options(req, &scope_id);
let params = json!({ "options": options });
let callback_result = self.bridge.call("callback/build_agent", params).await;
match callback_result {
Ok(result) => {
let mut modified_req = CreateSessionRequest {
model: req.model.clone(),
prompt: req.prompt.clone(),
system_prompt: req.system_prompt.clone(),
max_tokens: req.max_tokens,
event_tx: req.event_tx.clone(),
initial_turn: req.initial_turn.clone(),
build: req.build.clone(),
labels: req.labels.clone(),
deferred_prompt_policy: req.deferred_prompt_policy,
injected_context: req.injected_context.clone(),
};
if let Some(instructions) = result.get("additional_instructions") {
if let Some(arr) = instructions.as_array() {
let combined: Vec<&str> = arr.iter().filter_map(|v| v.as_str()).collect();
if !combined.is_empty() {
let extra = combined.join("\n");
use meerkat_core::config::SystemPromptOverride;
modified_req.system_prompt = match &modified_req.system_prompt {
SystemPromptOverride::Set(existing) => {
SystemPromptOverride::Set(format!("{existing}\n{extra}"))
}
SystemPromptOverride::Inherit => SystemPromptOverride::Set(extra),
SystemPromptOverride::Disable => {
tracing::warn!(
"callback/build_agent: additional_instructions ignored \
because system prompt is explicitly disabled"
);
SystemPromptOverride::Disable
}
};
}
}
}
if let Some(labels) = result.get("labels").and_then(|v| v.as_object()) {
let label_map = modified_req.labels.get_or_insert_with(Default::default);
for (k, v) in labels {
if let Some(s) = v.as_str() {
label_map.insert(k.clone(), s.to_string());
}
}
}
if let Some(resume_id) = result.get("resume_session_id").and_then(|v| v.as_str()) {
if let Some(ref store) = self.session_store {
let sid =
meerkat_core::types::SessionId::parse(resume_id).map_err(|_| {
SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: invalid resume_session_id: {resume_id}"
)))
})?;
if let Some(existing) = modified_req
.build
.as_ref()
.and_then(|b| b.resume_session.as_ref())
{
if existing.id() != &sid {
return Err(SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: resume_session_id conflict: \
spawn set {} but hook set {resume_id}",
existing.id()
))));
}
} else {
let session = store.load(&sid).await.map_err(|e| {
SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: failed to load resume session {resume_id}: {e}"
)))
})?;
let session = session.ok_or_else(|| {
SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: resume session not found: {resume_id}"
)))
})?;
let build = modified_req.build.get_or_insert_with(|| {
meerkat_core::service::SessionBuildOptions::default()
});
build.resume_session = Some(session);
}
} else {
return Err(SessionError::Agent(agent_tool_error(
"callback/build_agent: resume_session_id requires persistent mode \
(no session store available in ephemeral mode)"
.to_string(),
)));
}
}
if let Some(tools) = result.get("tools") {
match tools.as_array() {
Some(arr) => {
let mut tool_specs = Vec::with_capacity(arr.len());
for v in arr {
let spec = CallbackToolSpec::parse(v).map_err(|reason| {
SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: {reason}"
)))
})?;
tool_specs.push(spec);
}
if !tool_specs.is_empty() {
let dispatcher = CallbackToolDispatcher::new(
self.bridge.clone(),
scope_id.clone(),
tool_specs,
self.detached_jobs.clone(),
);
if dispatcher.reconcile_registered_catalog {
let reconciliation = dispatcher.clone();
tokio::spawn(async move {
if let Err(error) =
reconciliation.reconcile_detached_jobs().await
{
tracing::warn!(
%error,
"detached callback reconciliation failed"
);
}
});
}
let build = modified_req.build.get_or_insert_with(|| {
meerkat_core::service::SessionBuildOptions::default()
});
let pre_installed = build.external_tools.take();
build.external_tools = Some(
meerkat_mobkit::tool_compose::ComposedExternalTools::over(
Arc::new(dispatcher),
pre_installed,
),
);
}
}
None => {
return Err(SessionError::Agent(agent_tool_error(format!(
"callback/build_agent: tools must be a JSON array, got: {tools}"
))));
}
}
}
ensure_shell_tooling_build_substrate(&mut modified_req);
self.inner.build_agent(&modified_req, event_tx).await
}
Err(err) => {
Err(SessionError::Agent(agent_tool_error(format!(
"callback/build_agent failed: {err}"
))))
}
}
}
}
fn agent_tool_error(message: String) -> AgentError {
AgentError::Tool {
error: ToolError::execution_failed(message),
}
}
fn run_persistent() {
match tokio::runtime::Builder::new_multi_thread()
.enable_all()
.thread_stack_size(16 * 1024 * 1024)
.build()
{
Ok(runtime) => runtime.block_on(run_persistent_inner()),
Err(error) => {
eprintln!("failed to build tokio runtime: {error}");
std::process::exit(1);
}
}
}
async fn run_persistent_inner() {
use tokio::io::{AsyncBufReadExt, BufReader};
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new(DEFAULT_TRACING_FILTER)),
)
.with_writer(std::io::stderr)
.with_ansi(false)
.init();
let stdin = tokio::io::stdin();
let mut reader = BufReader::new(stdin);
let mut init_line = String::new();
if reader.read_line(&mut init_line).await.unwrap_or(0) == 0 {
eprintln!("rpc_gateway: stdin closed before init request");
std::process::exit(1);
}
let init_raw: Value = match serde_json::from_str(init_line.trim()) {
Ok(v) => v,
Err(e) => {
let error_response = json!({
"jsonrpc": "2.0",
"id": null,
"error": { "code": -32700, "message": format!("Parse error: {e}") }
});
println!(
"{}",
serde_json::to_string(&error_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string())
);
std::process::exit(1);
}
};
let request_id = init_raw.get("id").cloned().unwrap_or(Value::Null);
let method = init_raw
.get("method")
.and_then(|m| m.as_str())
.unwrap_or("");
if method != "mobkit/init" {
let error_response = json!({
"jsonrpc": "2.0",
"id": request_id,
"error": { "code": -32600, "message": format!("Expected mobkit/init, got {method}") }
});
println!(
"{}",
serde_json::to_string(&error_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string())
);
std::process::exit(1);
}
let params = init_raw.get("params").cloned().unwrap_or_else(|| json!({}));
let mob_config_param = params.get("mob_config").and_then(|v| v.as_str());
let is_workspace_config = mob_config_param.is_some();
let mob_config_toml = mob_config_param.unwrap_or(
r#"
[mob]
id = "persistent-gateway"
[profiles.default]
model = "gpt-5.5"
external_addressable = true
"#,
);
let definition = MobDefinition::from_toml(mob_config_toml).unwrap_or_else(|e| {
let error_response = json!({
"jsonrpc": "2.0",
"id": request_id,
"error": { "code": -32602, "message": format!("Invalid mob_config TOML: {e}") }
});
println!(
"{}",
serde_json::to_string(&error_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string())
);
std::process::exit(1);
});
let image_generation = mob_definition_may_use_image_generation(&definition);
let shell = mob_definition_may_use_shell(&definition);
{
let catalog = meerkat_models::catalog::catalog();
let known_models: std::collections::HashSet<&str> =
catalog.iter().map(|entry| entry.id).collect();
for (profile_name, binding) in &definition.profiles {
let Some(profile) = binding.as_inline() else {
continue; };
if !known_models.contains(profile.model.as_str()) {
let model = &profile.model;
let prefix = model.split('-').take(3).collect::<Vec<_>>().join("-");
let mut suggestions: Vec<&str> = known_models
.iter()
.filter(|m| {
m.starts_with(&prefix)
|| model
.starts_with(&m.split('-').take(3).collect::<Vec<_>>().join("-"))
})
.copied()
.collect();
suggestions.sort_unstable();
suggestions.truncate(5);
let hint = if suggestions.is_empty() {
String::new()
} else {
format!(". Did you mean one of: {}?", suggestions.join(", "))
};
fail_init(
&request_id,
-32602,
format!("Profile '{profile_name}' uses unknown model '{model}'{hint}"),
);
}
}
}
let (modules, pre_spawn) = parse_gateway_modules(¶ms);
let discovery_modules: Vec<String> = modules.iter().map(|m| m.id.clone()).collect();
let module_config = MobKitConfig {
modules,
discovery: DiscoverySpec {
namespace: "persistent-gateway".to_string(),
modules: discovery_modules,
},
pre_spawn,
};
let has_session_builder = params
.get("has_session_builder")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let persistent_state = params
.get("persistent_state")
.and_then(|v| v.as_str())
.map(std::path::PathBuf::from);
let has_roster_provider = params
.get("has_roster_provider")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let has_topology_provider = params
.get("has_topology_provider")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let has_agent_customizer = params
.get("has_agent_customizer")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let has_continuity_store = params
.get("has_continuity_store")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let has_lease_provider = params
.get("has_lease_provider")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let scratch_dir = params
.get("scratch_dir")
.and_then(|v| v.as_str())
.map(std::path::PathBuf::from);
let storage_layout = match persistent_state {
Some(ref state_path) => {
meerkat_mobkit::MobKitStorageLayout::with_injected_roots(state_path.clone(), None)
}
None => meerkat_mobkit::MobKitStorageLayout::declared_ephemeral(
meerkat_mobkit::storage_layout::default_ephemeral_scratch_root(),
),
};
let callback_job_store: Option<Arc<dyn meerkat::DetachedJobStore>> =
persistent_state.as_ref().map(|state_path| {
let path = meerkat_store::realm_paths_in(
state_path,
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
)
.jobs_sqlite_path;
if let Some(parent) = path.parent()
&& let Err(error) = std::fs::create_dir_all(parent)
{
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!(
"failed to create detached job store directory {}: {error}",
parent.display()
),
);
}
match meerkat::SqliteDetachedJobStore::open(path.clone()) {
Ok(store) => Arc::new(store) as Arc<dyn meerkat::DetachedJobStore>,
Err(error) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
meerkat_mobkit::storage_health::JobStoreResolutionError {
path,
message: error.to_string(),
}
.to_string(),
),
}
});
let (stdout_tx, mut stdout_rx) = mpsc::channel::<String>(64);
let (stdout_shutdown_tx, mut stdout_shutdown_rx) = oneshot::channel::<()>();
let mut stdout_writer = tokio::spawn(async move {
loop {
let line = tokio::select! {
shutdown = &mut stdout_shutdown_rx => {
let _ = shutdown;
while let Ok(line) = stdout_rx.try_recv() {
let mut stdout = std::io::stdout().lock();
let _ = writeln!(stdout, "{line}");
let _ = stdout.flush();
}
break;
}
line = stdout_rx.recv() => match line {
Some(line) => line,
None => break,
},
};
let mut stdout = std::io::stdout().lock();
let _ = writeln!(stdout, "{line}");
let _ = stdout.flush();
drop(stdout); }
});
let bridge = StdioCallbackBridge::new(stdout_tx.clone());
let (rpc_tx, mut rpc_rx) = mpsc::channel::<String>(64);
let shutdown_requested = Arc::new(AtomicBool::new(false));
let stdin_reader = tokio::spawn({
let bridge = bridge.clone();
let rpc_tx = rpc_tx.clone();
let stdout_tx = stdout_tx.clone();
let shutdown_requested = shutdown_requested.clone();
async move {
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line).await {
Ok(0) => break, Ok(_) => {}
Err(_) => break,
}
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let msg: Value = match serde_json::from_str(trimmed) {
Ok(v) => v,
Err(_) => continue,
};
let is_callback_response = msg
.get("id")
.and_then(|v| v.as_str())
.is_some_and(|id| id.starts_with("cb-"))
&& msg.get("method").is_none();
if is_callback_response {
bridge.route_callback_response(msg).await;
} else if gateway_shutdown_request(&msg).is_some() {
if shutdown_requested.swap(true, Ordering::AcqRel) {
if let Some(id) = msg.get("id") {
let response = json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": -32098,
"message": "gateway shutdown already in progress"
}
});
if let Ok(line) = serde_json::to_string(&response) {
let _ = stdout_tx.send(line).await;
}
}
} else if rpc_tx.send(trimmed.to_string()).await.is_err() {
break;
}
} else if shutdown_requested.load(Ordering::Acquire) {
if let Some(id) = msg.get("id") {
let response = json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": -32098,
"message": "gateway shutdown in progress"
}
});
if let Ok(line) = serde_json::to_string(&response) {
let _ = stdout_tx.send(line).await;
}
}
} else {
if rpc_tx.send(trimmed.to_string()).await.is_err() {
break;
}
}
}
bridge.close().await;
}
});
drop(rpc_tx);
fn fail_init(request_id: &Value, code: i64, message: String) -> ! {
let error_response = json!({
"jsonrpc": "2.0",
"id": request_id,
"error": { "code": code, "message": message }
});
let mut stdout = std::io::stdout().lock();
let _ = writeln!(
stdout,
"{}",
serde_json::to_string(&error_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string())
);
let _ = stdout.flush();
std::process::exit(1);
}
if (has_continuity_store || has_lease_provider || scratch_dir.is_some())
&& !(has_continuity_store && has_lease_provider && scratch_dir.is_some())
{
let mut missing = Vec::new();
if !has_continuity_store {
missing.push("continuity_store");
}
if !has_lease_provider {
missing.push("lease_provider");
}
if scratch_dir.is_none() {
missing.push("scratch_dir");
}
fail_init(
&request_id,
-32602,
format!(
"external-authoritative path requires continuity_store + lease_provider + scratch_dir; missing: {}",
missing.join(", ")
),
);
}
let mut gateway_options = parse_gateway_runtime_options(¶ms, persistent_state.as_deref())
.unwrap_or_else(|e| {
fail_init(&request_id, -32602, e);
});
validate_gateway_identity_bootstrap_intent(
gateway_options.identity_bootstrap_mode.as_ref(),
has_roster_provider,
)
.unwrap_or_else(|error| fail_init(&request_id, -32602, error));
if gateway_options.agent_memory.is_some() && !has_roster_provider {
fail_init(
&request_id,
-32602,
"runtime_options.agent_memory requires an identity-first roster provider".to_string(),
);
}
let default_llm_client: Option<Arc<dyn meerkat_client::LlmClient>> =
match gateway_options.demo_llm {
true => {
let client: Arc<dyn meerkat_client::LlmClient> =
Arc::new(meerkat_client::TestClient::default());
Some(client)
}
false => None,
};
let mut local_default_lease_provider: Option<
Arc<dyn meerkat_mobkit::identity_first::contracts::LeaseProvider>,
> = None;
let identity_continuity_store: Option<
Arc<dyn meerkat_mobkit::identity_first::ContinuityStore>,
> = if has_roster_provider {
Some(if has_continuity_store {
Arc::new(meerkat_mobkit::identity_first::GatewayContinuityStore::new(
bridge.clone(),
))
} else {
let db_path = match storage_layout.continuity_db() {
Ok(resolved) => resolved.path,
Err(e) => fail_init(&request_id, STORAGE_RESOLUTION_CODE, e.to_string()),
};
if storage_layout.is_declared_ephemeral()
&& let Err(e) = std::fs::create_dir_all(storage_layout.state_dir())
{
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to create the declared-ephemeral scratch root: {e}"),
);
}
let substrate = meerkat_mobkit::gateway_wiring::open_identity_substrate(&db_path)
.await
.unwrap_or_else(|e| fail_init(&request_id, STORAGE_RESOLUTION_CODE, e));
local_default_lease_provider = Some(substrate.lease_provider);
substrate.continuity_store
})
} else {
None
};
let identity_lease_provider: Option<
Arc<dyn meerkat_mobkit::identity_first::contracts::LeaseProvider>,
> = if has_roster_provider {
Some(if has_lease_provider {
Arc::new(meerkat_mobkit::identity_first::GatewayLeaseProvider::new(
bridge.clone(),
)) as Arc<dyn meerkat_mobkit::identity_first::contracts::LeaseProvider>
} else {
local_default_lease_provider
.clone()
.expect("local substrate initialized with the continuity store")
})
} else {
None
};
let identity_session_store_adapter = identity_continuity_store.as_ref().map(|store| {
Arc::new(meerkat_mobkit::identity_first::ContinuitySessionStoreAdapter::new(store.clone()))
});
let schedule_owner_id = definition.id.to_string();
let (
mob_spec,
_temp_dir,
schedule_host_inputs,
transcript_edit_service,
workgraph_service,
live_inputs,
gateway_detached_jobs,
) = if let Some(ref state_path) = persistent_state {
if let Err(e) = std::fs::create_dir_all(state_path) {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to create persistent state directory: {e}"),
);
}
let sqlite_path = match storage_layout.session_db() {
Ok(resolved) => resolved.path,
Err(e) => fail_init(&request_id, STORAGE_RESOLUTION_CODE, e.to_string()),
};
let session_store_kind = if identity_session_store_adapter.is_some() {
"ContinuitySessionStoreAdapter"
} else {
"SqliteSessionStore"
};
let session_store: Arc<dyn meerkat::SessionStore> =
if let Some(adapter) = identity_session_store_adapter.clone() {
adapter
} else {
match meerkat_store::SqliteSessionStore::open(sqlite_path) {
Ok(s) => Arc::new(s),
Err(e) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to open SQLite session store: {e}"),
),
}
};
let session_store_incremental =
meerkat_mobkit::storage_health::probe_session_store_incremental(
&session_store,
session_store_kind,
);
let mob_storage = MobStorage::in_memory();
let binary_blob_store: Arc<dyn BinaryBlobStore> =
match ObjectStoreBlobStore::local(storage_layout.blob_root()) {
Ok(store) => Arc::new(store),
Err(e) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to open binary blob store: {e}"),
),
};
let blob_store: Arc<dyn meerkat_core::BlobStore> =
Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
let runtime_db_path = storage_layout.runtime_db();
let (runtime_store, runtime_store_slot): (
Arc<dyn meerkat_runtime::RuntimeStore>,
meerkat_mobkit::storage_health::StorageSlotSummary,
) = if gateway_options.runtime_store_ephemeral {
(
Arc::new(meerkat_runtime::InMemoryRuntimeStore::new()),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"runtime",
"InMemoryRuntimeStore",
"explicitly declared via runtime_options.runtime_store: sessions do not \
survive gateway restart",
),
)
} else {
match meerkat_runtime::store::SqliteRuntimeStore::new(&runtime_db_path) {
Ok(store) => (
Arc::new(store) as Arc<dyn meerkat_runtime::RuntimeStore>,
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"runtime",
"SqliteRuntimeStore",
),
),
Err(err) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
meerkat_mobkit::storage_health::RuntimeStoreResolutionError {
path: runtime_db_path.clone(),
message: err.to_string(),
}
.to_string(),
),
}
};
let (runtime_store, session_write_epochs) =
meerkat_mobkit::mob_handle_runtime::epoch_tracking_runtime_store_with_durable_projection(
runtime_store,
session_store.clone(),
);
let adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
Arc::clone(&runtime_store),
Arc::clone(&blob_store),
));
let mut factory = AgentFactory::new(state_path)
.builtins(false)
.shell(shell)
.comms(true);
if image_generation {
factory = factory.with_image_generation_machine(adapter.clone());
}
let live_agent_factory = factory.clone();
let live_machine = Arc::clone(&adapter);
let mut inner_builder =
FactoryAgentBuilder::new(factory, gateway_agent_config(&gateway_options));
inner_builder.default_session_store = Some(Arc::new(meerkat_store::StoreAdapter::new(
session_store.clone(),
)));
inner_builder.default_blob_store = Some(blob_store.clone());
inner_builder.default_detached_job_store = callback_job_store.clone();
if let Some(job_store) = callback_job_store.as_ref() {
inner_builder.default_shell_job_delivery_projector =
Some(Arc::new(meerkat::JobOutboxProjector::new_for_realm(
Arc::clone(job_store),
meerkat_runtime::RuntimeDeliveryInbox::new(Arc::clone(&runtime_store)),
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
)));
}
let (schedule_tools, schedule_slot) =
match meerkat_mobkit::schedule_wiring::attach_schedule_tools_with_identity_targets_reporting(
&inner_builder,
storage_layout.state_dir(),
) {
Ok(tools) => (
Some(tools),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"schedule",
"SqliteScheduleStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"schedule",
format!("schedule store failed to open; schedule tools disabled: {error}"),
),
),
};
let (workgraph, workgraph_slot) = match &gateway_options.workgraph {
GatewayWorkgraphOption::Disabled => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"workgraph",
"disabled",
"explicitly disabled via runtime_options.workgraph = false",
),
),
GatewayWorkgraphOption::Enabled => {
match meerkat_mobkit::workgraph_wiring::attach_workgraph_tools_reporting(
&inner_builder,
storage_layout.state_dir(),
&schedule_owner_id,
) {
Ok((service, slot)) => (
Some((service, slot, storage_layout.state_dir().to_path_buf())),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"workgraph",
"SqliteWorkGraphStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"workgraph",
format!("workgraph store failed to open; workgraph disabled: {error}"),
),
),
}
}
GatewayWorkgraphOption::DurableDir(dir) => {
let _ = std::fs::create_dir_all(dir);
match meerkat_mobkit::workgraph_wiring::attach_workgraph_tools_reporting(
&inner_builder,
dir,
&schedule_owner_id,
) {
Ok((service, slot)) => (
Some((service, slot, dir.clone())),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"workgraph",
"SqliteWorkGraphStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"workgraph",
format!("workgraph store failed to open; workgraph disabled: {error}"),
),
),
}
}
};
let workgraph_service = workgraph.as_ref().map(|(service, _, _)| service.clone());
let agent_mob_tools_slot = Arc::clone(&inner_builder.default_mob_tools);
let detached_jobs = callback_job_store.as_ref().map(|store| {
DetachedCallbackJobRuntime::new(
meerkat_mobkit::storage_provider::MEERKAT_LEVEL_REALM_ID,
Arc::clone(store),
blob_store.clone(),
)
.with_runtime_delivery_store(Arc::clone(&runtime_store))
.with_monitor_shell(
state_path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.unwrap_or(state_path)
.to_path_buf(),
shell,
)
});
let callback_builder = StdioCallbackAgentBuilder {
inner: inner_builder,
bridge: bridge.clone(),
has_session_builder,
session_store: Some(session_store.clone()),
detached_jobs: detached_jobs.clone(),
};
let concrete_service = Arc::new(meerkat_session::PersistentSessionService::new(
callback_builder,
gateway_options.max_sessions,
session_store,
Arc::clone(&runtime_store),
blob_store,
));
if let Some(detached_jobs) = detached_jobs.as_ref() {
let delivery_service: Arc<dyn meerkat_mob::MobSessionService> =
concrete_service.clone();
detached_jobs.attach_delivery_service(delivery_service);
}
let schedule_host_inputs = schedule_tools.map(|tools| {
(
tools.service,
tools.mob_target_registry,
Arc::clone(&concrete_service),
adapter.clone(),
storage_layout.schedule_db(),
)
});
let transcript_edit_service: Option<
Arc<dyn meerkat_mobkit::memory::hygienist::TranscriptEditSessionService>,
> = Some(Arc::clone(&concrete_service) as _);
let live_inputs = if matches!(gateway_options.live, GatewayLiveOption::Enabled { .. }) {
Some((
Arc::clone(&concrete_service),
live_machine,
live_agent_factory,
))
} else {
None
};
let committed_boundary_recoverer: Arc<
dyn meerkat_mobkit::identity_first::CommittedBoundaryRecoverer,
> = Arc::clone(&concrete_service) as _;
let session_service: Arc<dyn meerkat_mob::MobSessionService> = concrete_service;
let mut spec = MobBootstrapSpec::new(definition, mob_storage, session_service)
.with_session_write_epochs(&session_write_epochs)
.with_runtime_archived_terminal_authority(Arc::clone(&runtime_store))
.with_session_runtime_adapter(adapter.clone())
.with_workgraph_service(workgraph_service.clone())
.with_agent_mob_tools(agent_mob_tools_slot)
.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: default_llm_client.clone(),
});
spec.committed_boundary_recoverer = Some(committed_boundary_recoverer);
if let Some((_, admission_slot, workgraph_state_dir)) = &workgraph {
spec = spec
.with_workgraph_admission_slot(admission_slot.clone())
.with_workgraph_admission_sidecar(workgraph_state_dir);
}
spec.runtime_adapter = Some(adapter);
spec.binary_blob_store = Some(binary_blob_store);
let mut slots = vec![
if identity_session_store_adapter.is_some() {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"sessions",
"ContinuitySessionStoreAdapter",
)
.with_detail("sessions ride the identity continuity store")
} else {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"sessions",
"SqliteSessionStore",
)
},
runtime_store_slot,
meerkat_mobkit::storage_health::blob_slot_summary(
meerkat_mobkit::storage_health::BlobDurability::PersistentDisk,
),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"console",
"SqliteConsoleLogStore",
),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"metadata",
"SqliteMetadataStore",
),
schedule_slot,
workgraph_slot,
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"jobs",
"SqliteDetachedJobStore",
),
gateway_event_log_slot(&gateway_options),
];
if let Some(agent_memory) = gateway_options.agent_memory.as_ref() {
slots.push(agent_memory_census_slot(agent_memory));
}
if identity_continuity_store.is_some() {
slots.push(if has_continuity_store {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"continuity",
"GatewayContinuityStore (SDK-hosted)",
)
.with_detail("durability rides with the SDK-hosted continuity store")
} else {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"continuity",
"LocalContinuityStore",
)
});
}
slots.extend(meerkat_mobkit::storage_health::scratch_ring_buffer_slots());
spec.resolved_storage = Some(
meerkat_mobkit::storage_health::ResolvedStorageSummary::new(
meerkat_mobkit::storage_health::BlobDurability::PersistentDisk,
Some(session_store_incremental),
)
.with_state_dir(storage_layout.state_dir())
.with_slots(slots),
);
(
spec,
None,
schedule_host_inputs,
transcript_edit_service,
workgraph_service,
live_inputs,
detached_jobs,
)
} else {
let temp_dir = if scratch_dir.is_none() {
Some(tempfile::tempdir().expect("create temp dir for agent working space"))
} else {
None
};
let agent_workspace = scratch_dir
.as_deref()
.or_else(|| temp_dir.as_ref().map(|dir| dir.path()))
.expect("scratch dir or temp dir");
if let Err(err) = std::fs::create_dir_all(agent_workspace) {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to create scratch directory: {err}"),
);
}
let binary_blob_store: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
let blob_store: Arc<dyn meerkat_core::BlobStore> =
Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
Arc::new(meerkat_runtime::InMemoryRuntimeStore::new());
let (runtime_store, session_write_epochs) =
meerkat_mobkit::mob_handle_runtime::epoch_tracking_runtime_store(runtime_store);
let adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
Arc::clone(&runtime_store),
Arc::clone(&blob_store),
));
let mut factory = AgentFactory::new(agent_workspace)
.builtins(false)
.shell(shell)
.comms(true)
.session_store(Arc::new(meerkat::MemoryStore::new()));
if image_generation {
factory = factory.with_image_generation_machine(adapter.clone());
}
let mut inner_builder =
FactoryAgentBuilder::new(factory, gateway_agent_config(&gateway_options));
inner_builder.default_blob_store = Some(blob_store.clone());
let mut workgraph_sidecar_dir: Option<PathBuf> = None;
let (workgraph_service, ephemeral_workgraph_slot) = match &gateway_options.workgraph {
GatewayWorkgraphOption::Disabled => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"workgraph",
"disabled",
"explicitly disabled via runtime_options.workgraph = false",
),
),
GatewayWorkgraphOption::DurableDir(dir) => {
let _ = std::fs::create_dir_all(dir);
workgraph_sidecar_dir = Some(dir.clone());
match meerkat_mobkit::workgraph_wiring::open_workgraph_service_reporting(
dir,
&schedule_owner_id,
) {
Ok(service) => (
Some(service),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"workgraph",
"SqliteWorkGraphStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"workgraph",
format!("workgraph store failed to open; workgraph disabled: {error}"),
),
),
}
}
GatewayWorkgraphOption::Enabled => {
if has_continuity_store {
tracing::warn!(
"workgraph is MEMORY-BACKED for this launch: the continuity store \
is SDK-hosted (no local path) and no persistent_state dir is set, \
so goals, work items and attention bindings will NOT survive a \
gateway restart. Set runtime_options.workgraph to a directory \
path (or provide persistent_state) to place a durable \
workgraph.sqlite3.",
);
}
(
Some(
meerkat_mobkit::workgraph_wiring::ephemeral_workgraph_service(
&schedule_owner_id,
),
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"workgraph",
"MemoryWorkGraphStore",
if has_continuity_store {
"memory-backed on an otherwise durable identity-first launch: set \
runtime_options.workgraph to a directory (or persistent_state) \
for a durable store"
} else {
"declared by the ephemeral launch mode"
},
),
)
}
};
let mut workgraph_admission_slots = Vec::new();
if let Some(service) = workgraph_service.as_ref() {
workgraph_admission_slots.push(
meerkat_mobkit::workgraph_wiring::install_workgraph_tools(&inner_builder, service),
);
}
let callback_builder = StdioCallbackAgentBuilder {
inner: inner_builder,
bridge: bridge.clone(),
has_session_builder,
session_store: None,
detached_jobs: None,
};
let mut transcript_edit_service: Option<
Arc<dyn meerkat_mobkit::memory::hygienist::TranscriptEditSessionService>,
> = None;
let mut committed_boundary_recoverer: Option<
Arc<dyn meerkat_mobkit::identity_first::CommittedBoundaryRecoverer>,
> = None;
let mut session_store_incremental: Option<bool> = None;
let session_service: Arc<dyn meerkat_mob::MobSessionService> =
if let Some(session_adapter) = identity_session_store_adapter.clone() {
let session_store: Arc<dyn meerkat::SessionStore> = session_adapter.clone();
session_store_incremental = Some(
meerkat_mobkit::storage_health::probe_session_store_incremental(
&session_store,
"ContinuitySessionStoreAdapter",
),
);
let mut factory = AgentFactory::new(agent_workspace)
.builtins(false)
.shell(shell)
.comms(true)
.session_store(Arc::new(meerkat::MemoryStore::new()));
if image_generation {
factory = factory.with_image_generation_machine(adapter.clone());
}
let mut inner_builder =
FactoryAgentBuilder::new(factory, gateway_agent_config(&gateway_options));
inner_builder.default_session_store = Some(Arc::new(
meerkat_store::StoreAdapter::new(session_store.clone()),
));
inner_builder.default_blob_store = Some(blob_store.clone());
if let Some(service) = workgraph_service.as_ref() {
workgraph_admission_slots.push(
meerkat_mobkit::workgraph_wiring::install_workgraph_tools(
&inner_builder,
service,
),
);
}
let callback_builder = StdioCallbackAgentBuilder {
inner: inner_builder,
bridge: bridge.clone(),
has_session_builder,
session_store: Some(session_store.clone()),
detached_jobs: None,
};
let concrete = Arc::new(meerkat_session::PersistentSessionService::new(
callback_builder,
gateway_options.max_sessions,
session_store,
Arc::clone(&runtime_store),
blob_store.clone(),
));
transcript_edit_service = Some(Arc::clone(&concrete) as _);
committed_boundary_recoverer = Some(Arc::clone(&concrete) as _);
concrete
} else {
Arc::new(EphemeralSessionService::new(
callback_builder,
gateway_options.max_sessions,
))
};
let mut spec = MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
.with_session_write_epochs(&session_write_epochs)
.with_session_runtime_adapter(adapter.clone())
.with_workgraph_service(workgraph_service.clone())
.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: default_llm_client.clone(),
});
spec.committed_boundary_recoverer = committed_boundary_recoverer;
for slot in workgraph_admission_slots {
spec = spec.with_workgraph_admission_slot(slot);
}
if let Some(dir) = workgraph_sidecar_dir.as_deref() {
spec = spec.with_workgraph_admission_sidecar(dir);
}
spec.runtime_adapter = Some(adapter);
spec.binary_blob_store = Some(binary_blob_store);
let mut slots = vec![
if identity_session_store_adapter.is_some() {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"sessions",
"ContinuitySessionStoreAdapter",
)
.with_detail("sessions ride the identity continuity store")
} else {
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"sessions",
"EphemeralSessionService",
"declared by the ephemeral launch mode",
)
},
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"runtime",
"InMemoryRuntimeStore",
"declared by the ephemeral launch mode",
),
meerkat_mobkit::storage_health::blob_slot_summary(
meerkat_mobkit::storage_health::BlobDurability::DeclaredEphemeral,
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"console",
"InMemoryConsoleLogStore",
"declared default of the ephemeral launch mode",
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"metadata",
"InMemoryMetadataStore",
"declared default of the ephemeral launch mode",
),
ephemeral_workgraph_slot,
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"jobs",
"disabled",
"semantic detached admission is unavailable in ephemeral gateway mode",
),
gateway_event_log_slot(&gateway_options),
];
if identity_continuity_store.is_some() {
slots.push(if has_continuity_store {
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"continuity",
"GatewayContinuityStore (SDK-hosted)",
)
.with_detail("durability rides with the SDK-hosted continuity store")
} else {
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"continuity",
"LocalContinuityStore",
"backed by the declared-ephemeral scratch root",
)
});
}
slots.extend(meerkat_mobkit::storage_health::scratch_ring_buffer_slots());
spec.resolved_storage = Some(
meerkat_mobkit::storage_health::ResolvedStorageSummary::new(
meerkat_mobkit::storage_health::BlobDurability::DeclaredEphemeral,
session_store_incremental,
)
.with_state_dir(storage_layout.state_dir())
.with_slots(slots),
);
(
spec,
temp_dir,
None,
transcript_edit_service,
workgraph_service,
None,
None,
)
};
let mob_spec = if has_session_builder {
let after_bridge = bridge.clone();
mob_spec.with_after_create_hook(Arc::new(
move |session_id: meerkat_core::types::SessionId, ctx| {
let b = after_bridge.clone();
Box::pin(async move {
b.notify_reliable(
"callback/after_create",
json!({
"session_id": session_id.to_string(),
"model": ctx.model,
"labels": ctx.labels,
"system_prompt": ctx.system_prompt,
}),
)
.await;
})
},
))
} else {
mob_spec
};
let timeout = GATEWAY_RUNTIME_EVENT_DRAIN_TIMEOUT;
let persistent_metadata: Arc<dyn PersistentMetadataStore> = if persistent_state.is_some() {
let metadata_path = match storage_layout.metadata_db() {
Ok(resolved) => resolved.path,
Err(e) => fail_init(&request_id, STORAGE_RESOLUTION_CODE, e.to_string()),
};
Arc::new(
SqliteMetadataStore::open(&metadata_path).unwrap_or_else(|e| {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!(
"failed to open the mobkit metadata store at {}: {e}",
metadata_path.display()
),
);
}),
)
} else {
Arc::new(InMemoryMetadataStore::new())
};
let mut runtime = Box::pin(UnifiedRuntime::bootstrap_with_options(
mob_spec,
module_config,
Vec::new(),
timeout,
gateway_options.runtime_options.clone(),
persistent_metadata,
))
.await
.unwrap_or_else(|e| {
let error_response = json!({
"jsonrpc": "2.0",
"id": request_id,
"error": { "code": -32603, "message": format!("Runtime bootstrap failed: {e}") }
});
let mut stdout = std::io::stdout().lock();
let _ = writeln!(
stdout,
"{}",
serde_json::to_string(&error_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string())
);
let _ = stdout.flush();
std::process::exit(1);
});
if persistent_state.is_some() {
let console_log_path = match storage_layout.console_db() {
Ok(resolved) => resolved.path,
Err(e) => fail_init(&request_id, STORAGE_RESOLUTION_CODE, e.to_string()),
};
let console_log_store = Arc::new(
SqliteConsoleLogStore::open(&console_log_path).unwrap_or_else(|e| {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!(
"failed to open the mobkit console store at {}: {e}",
console_log_path.display()
),
);
}),
);
runtime.set_console_log_store(console_log_store);
}
for route in gateway_options.routing_routes.iter().cloned() {
if let Err(err) = runtime.add_runtime_route(route).await {
fail_init(
&request_id,
-32602,
format!("runtime_options.routing_config_path route failed validation: {err}"),
);
}
}
if let Some(event_log_config) = gateway_options.event_log.take() {
runtime.start_event_log(event_log_config);
}
if let Some(access) = gateway_options.access.take() {
runtime.set_access_controller(access);
}
let gateway_error_hook: meerkat_mobkit::ErrorHook = {
let error_bridge = bridge.clone();
Arc::new(move |event| {
let b = error_bridge.clone();
Box::pin(async move {
if let Ok(params) = serde_json::to_value(&event) {
b.notify("mobkit/on_error", params);
}
})
})
};
runtime.set_error_hook(gateway_error_hook.clone());
let steward_late_runtime = StewardLateRuntime::default();
let steward_roster_slot: Arc<
std::sync::Mutex<Vec<meerkat_mobkit::identity_first::DurableAgentSpec>>,
> = Arc::new(std::sync::Mutex::new(Vec::new()));
let mut agent_memory_steward: Option<Arc<meerkat_mobkit::memory::steward::StewardEngine>> =
None;
let identity_ctx: Option<meerkat_mobkit::rpc::IdentityFirstContext> = if has_roster_provider {
use meerkat_mobkit::identity_first::{
AgentRuntimeServices, DurabilityPolicy, IdentityFirstRuntimeContext, IdentityRuntime,
IdentityRuntimeConfig, RosterContext,
gateway_bridges::{
GatewayAgentCustomizer, GatewayRosterProvider, GatewayTopologyProvider,
},
};
let continuity_store = identity_continuity_store
.clone()
.expect("identity continuity store initialized with roster provider");
let lease_provider = identity_lease_provider
.clone()
.expect("identity lease provider initialized with roster provider");
let mob_handle = runtime.mob_handle();
let mut identity_bridge = if let Some(adapter) = identity_session_store_adapter.clone() {
meerkat_mobkit::identity_first::MobSessionBridge::with_continuity_session_store(
mob_handle.clone(),
adapter,
runtime.mob_runtime().session_service().cloned(),
)
} else if let Some(session_service) = runtime.mob_runtime().session_service().cloned() {
meerkat_mobkit::identity_first::MobSessionBridge::with_session_service(
mob_handle.clone(),
session_service,
)
} else {
meerkat_mobkit::identity_first::MobSessionBridge::new(mob_handle.clone())
};
if let Some(recoverer) = runtime.mob_runtime().committed_boundary_recoverer() {
identity_bridge = identity_bridge.with_committed_boundary_recoverer(recoverer);
}
let bridge_arc: Arc<dyn meerkat_mobkit::identity_first::SessionBridge> =
Arc::new(identity_bridge);
let irt = Arc::new(
IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store,
lease_provider,
runtime_instance_id: format!("gateway-{}", std::process::id()),
has_runtime_store: identity_session_store_adapter.is_some()
|| persistent_state.is_some(),
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(bridge_arc),
default_timeout: None,
})
.with_runtime_services(AgentRuntimeServices::new(mob_handle)),
);
irt.set_error_hook(Some(gateway_error_hook.clone()));
let roster: Arc<dyn meerkat_mobkit::identity_first::contracts::RosterProvider> =
Arc::new(GatewayRosterProvider::new(bridge.clone()));
let mob_definition = runtime.mob_handle().definition().clone();
irt.set_reset_roster_provider_context(Some(roster.clone()), Some(mob_definition.clone()));
let topology: Option<Arc<dyn meerkat_mobkit::identity_first::contracts::TopologyProvider>> =
if has_topology_provider {
Some(Arc::new(GatewayTopologyProvider::new(bridge.clone())))
} else {
None
};
let base_customizer: Option<
Arc<dyn meerkat_mobkit::identity_first::contracts::AgentCustomizer>,
> = if has_agent_customizer {
Some(Arc::new(GatewayAgentCustomizer::new(bridge.clone())))
} else {
None
};
let mut agent_memory_taint: Option<meerkat_mobkit::SessionTaintTracker> = None;
let agent_memory_compaction_reset: Arc<
std::sync::OnceLock<meerkat_mobkit::AgentMemoryRuntimeInjector>,
> = Arc::new(std::sync::OnceLock::new());
let mut agent_memory_distiller: Option<
Arc<meerkat_mobkit::memory::distiller::DistillerEngine>,
> = None;
let mut agent_memory_hygienist: Option<
Arc<meerkat_mobkit::memory::hygienist::HygienistEngine>,
> = None;
let agent_memory_provider: Option<Arc<dyn meerkat_mobkit::AgentMemoryProvider>> =
if let Some(agent_memory) = gateway_options.agent_memory.as_ref() {
let provider: Arc<dyn meerkat_mobkit::AgentMemoryProvider> =
match agent_memory.store {
GatewayAgentMemoryStoreKind::Markdown => Arc::new(
meerkat_mobkit::MarkdownAgentMemoryStore::open(&agent_memory.path)
.unwrap_or_else(|e| {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to open agent memory store: {e}"),
);
}),
),
GatewayAgentMemoryStoreKind::Sqlite => Arc::new(
meerkat_mobkit::SqliteAgentMemoryStore::open(&agent_memory.path)
.unwrap_or_else(|e| {
fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("failed to open agent memory store: {e}"),
);
}),
),
};
if provider.as_taintable().is_some() {
let memory_transcript_store: Option<Arc<dyn meerkat::SessionStore>> =
if agent_memory.distiller.enabled || agent_memory.steward.enabled {
if persistent_state.is_none() {
fail_init(
&request_id,
-32602,
"agent memory distiller/steward require persistent_state"
.to_string(),
);
}
Some(
if let Some(adapter) = identity_session_store_adapter.clone() {
adapter
} else {
let session_db = match storage_layout.session_db() {
Ok(resolved) => resolved.path,
Err(e) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
e.to_string(),
),
};
match meerkat_store::SqliteSessionStore::open(session_db) {
Ok(store) => Arc::new(store),
Err(e) => fail_init(
&request_id,
STORAGE_RESOLUTION_CODE,
format!("agent memory session store: {e}"),
),
}
},
)
} else {
None
};
let memory_events = runtime.memory_event_sink();
let engines = meerkat_mobkit::memory_wiring::MemoryEnginesConfig {
distiller: agent_memory.distiller.clone(),
steward: agent_memory.steward.clone(),
};
let stack = match meerkat_mobkit::memory_wiring::attach_memory_engines(
provider.clone(),
&agent_memory.config,
&engines,
meerkat_mobkit::memory_wiring::MemoryStackSeams {
persistent_state: persistent_state.clone(),
transcript_store: memory_transcript_store,
event_sink: Some(memory_events.clone()),
mob_purpose: Some(Arc::new(GatewayMobPurposeSource {
mob: mob_definition.id.to_string(),
roster: steward_roster_slot.clone(),
})),
steward_gating: Some(Arc::new(GatewayMemoryGatingBridge {
runtime: steward_late_runtime.clone(),
})),
steward_conflicts: Some(Arc::new(GatewayMemoryConflictBridge {
runtime: steward_late_runtime.clone(),
handle: tokio::runtime::Handle::current(),
})),
},
) {
Ok(stack) => stack,
Err(e) => fail_init(&request_id, -32602, e),
};
let panel = stack.panel.clone();
let steward_store = stack.steward_store.clone();
let tracker = stack.taint;
let mut sinks = stack.sinks;
agent_memory_distiller = stack.distiller;
if let Some(engine) = agent_memory_distiller.as_ref() {
let _ = engine;
tracing::info!(
model = ?agent_memory.distiller.model,
"agent memory distiller installed"
);
}
if let Some(engine) = stack.steward.clone() {
if schedule_host_inputs.is_none() {
std::mem::forget(engine.spawn_dream_loop());
}
agent_memory_steward = Some(engine);
tracing::info!(
model = ?agent_memory.steward.model,
cadence = %agent_memory.steward.cadence,
per_mob = agent_memory.steward.per_mob,
"agent memory steward installed"
);
}
let taint_mob_handle = runtime.mob_handle();
tracker.set_outbound_taint_declarer(std::sync::Arc::new(
move |identity: &str,
taint: Option<meerkat_core::comms::SenderContentTaint>| {
let handle = taint_mob_handle.clone();
let member =
meerkat_mobkit::member_comms_id::mob_member_id(identity);
let identity_owned = identity.to_string();
tokio::spawn(async move {
if let Err(err) = handle
.declare_member_outbound_taint(member, taint)
.await
{
tracing::debug!(
identity = %identity_owned,
?taint,
error = %err,
"agent memory taint: outbound taint declaration \
failed (member may not be materialized)"
);
}
});
},
));
if let Some(panel) = panel {
runtime.set_memory_panel_store(panel);
}
if agent_memory.hygienist.enabled {
use meerkat_mobkit::memory::hygienist as memory_hygienist;
let Some(steward_spans) = steward_store.clone() else {
fail_init(
&request_id,
-32602,
"agent memory hygienist requires a provider with the \
steward surface (StewardStore)"
.to_string(),
);
};
let Some(edit_service) = transcript_edit_service.clone() else {
fail_init(
&request_id,
-32602,
"agent memory hygienist requires persistent sessions \
(runtime_options.persistent_state or an identity session \
store): the transcript-revision seam only exists on the \
persistent session service"
.to_string(),
);
};
let mut profile = memory_hygienist::HygienistProfile::embedded_default();
if let Some(model) = agent_memory.hygienist.model.as_deref() {
profile = profile.with_model_override(model).unwrap_or_else(|e| {
fail_init(
&request_id,
-32602,
format!("agent memory hygienist: {e}"),
);
});
}
let Some(state) = persistent_state.clone() else {
fail_init(
&request_id,
-32602,
"agent memory hygienist requires persistent_state".to_string(),
);
};
let handle = memory_hygienist::FactoryHygienistHandle::new(
state,
meerkat::Config::default(),
agent_memory.config.realm.clone(),
&profile,
);
let model = profile.model.clone();
let gate: Option<Arc<dyn memory_hygienist::DistillationGate>> =
agent_memory_distiller.clone().map(|engine| {
engine as Arc<dyn memory_hygienist::DistillationGate>
});
let engine = Arc::new(memory_hygienist::HygienistEngine::new(
profile,
agent_memory.hygienist.clone(),
Arc::new(handle),
Arc::new(memory_hygienist::SessionServiceRevisionSeam::new(
edit_service,
)),
Arc::new(memory_hygienist::StoreSpanReferenceSource::new(
steward_spans,
agent_memory.config.realm.clone(),
)),
gate,
agent_memory.config.realm.clone(),
));
engine.set_event_sink(memory_events.clone());
match agent_memory_distiller.as_ref() {
Some(distiller) => distiller.set_compaction_follow_up(
memory_hygienist::distiller_follow_up(engine.clone()),
),
None => sinks.push(Arc::new(memory_hygienist::HygienistTriggers::new(
engine.clone(),
))),
}
agent_memory_hygienist = Some(engine);
tracing::info!(
model = %model,
runs_per_day = agent_memory.hygienist.runs_per_day,
"agent memory hygienist installed"
);
}
let compaction_reset_slot = agent_memory_compaction_reset.clone();
sinks.push(Arc::new(meerkat_mobkit::CompactionResetSink::new(
Arc::new(move |session: &str| {
if let Some(injector) = compaction_reset_slot.get() {
injector.on_session_compacted(session);
}
}),
)));
std::mem::forget(meerkat_mobkit::spawn_member_event_observer(
runtime.mob_handle(),
sinks,
));
agent_memory_taint = Some(tracker);
Some(stack.provider)
} else {
if agent_memory.distiller.enabled
|| agent_memory.steward.enabled
|| agent_memory.hygienist.enabled
{
fail_init(
&request_id,
-32602,
"agent memory distiller/steward/hygienist require a provider \
with judgment-plane capabilities (the bundled sqlite store)"
.to_string(),
);
}
Some(provider)
}
} else {
None
};
{
use meerkat_mobkit::memory::selector as memory_selector;
let configured = gateway_options
.agent_memory
.as_ref()
.and_then(|agent_memory| agent_memory.selector.as_ref());
let spec = resolve_selector_spec(configured).unwrap_or_else(|e| {
fail_init(&request_id, -32602, format!("agent memory selector: {e}"));
});
let profile = memory_selector::profile_for_spec(&spec).unwrap_or_else(|e| {
fail_init(&request_id, -32602, format!("agent memory selector: {e}"));
});
if let Some(profile) = profile {
let Some(agent_memory) = gateway_options.agent_memory.as_ref() else {
fail_init(
&request_id,
-32602,
"agent memory selector requires runtime_options.agent_memory".to_string(),
);
};
let fetch = agent_memory_provider
.as_ref()
.and_then(|provider| provider.as_selected_record_fetch())
.unwrap_or_else(|| {
fail_init(
&request_id,
-32602,
"agent memory selector requires a provider with selected-record \
fetch support (SelectedRecordFetch)"
.to_string(),
);
});
let factory_state = storage_layout.state_dir().to_path_buf();
let handle = memory_selector::FactorySelectorHandle::new(
factory_state,
meerkat::Config::default(),
agent_memory.config.realm.clone(),
&profile,
);
let model = profile.model.clone();
memory_selector::install(Arc::new(memory_selector::SelectorRuntime {
stage: Arc::new(memory_selector::SelectorStage::new(
profile,
Arc::new(handle),
)),
fetch,
}));
tracing::info!(model = %model, "agent memory selector installed");
}
}
let agent_memory_operator_resolver: Option<
Arc<dyn meerkat_mobkit::memory::coordinator::OperatorResolver>,
> = if gateway_options.agent_memory.as_ref().is_some_and(|memory| {
memory.config.operator_scope == meerkat_mobkit::AgentMemoryOperatorScope::Provisional
}) {
let resolver = Arc::new(meerkat_mobkit::ConsolePrincipalOperatorResolver::new());
runtime.set_console_operator_resolver(resolver.clone());
tracing::info!(
"agent memory operator scope active (provisional keying: console auth principal)"
);
Some(resolver)
} else {
None
};
if operator_scope_recall_inert(
gateway_options.agent_memory.as_ref(),
agent_memory_operator_resolver.is_some(),
) {
tracing::warn!(
"agent_memory.operator_scope=\"provisional\" is configured but this gateway \
installs no operator resolver: operator-scope recall composition is INERT \
(records routed to the operator scope will not be recalled or injected); \
steward routing of operator-scope proposals remains active"
);
}
let agent_memory_mob_resolver: Option<
Arc<dyn meerkat_mobkit::memory::coordinator::MobScopeResolver>,
> = gateway_options.agent_memory.as_ref().map(|memory| {
Arc::new(meerkat_mobkit::memory::coordinator::StaticMobBinding {
realm: memory.config.realm.clone(),
mob: runtime.mob_handle().mob_id().to_string(),
}) as Arc<dyn meerkat_mobkit::memory::coordinator::MobScopeResolver>
});
let customizer: Option<
Arc<dyn meerkat_mobkit::identity_first::contracts::AgentCustomizer>,
> = if let (Some(agent_memory), Some(provider)) = (
gateway_options.agent_memory.as_ref(),
agent_memory_provider.clone(),
) {
Some(Arc::new(
meerkat_mobkit::AgentMemoryCustomizer::wrap(
base_customizer,
provider,
agent_memory.config.clone(),
)
.with_operator_resolver(agent_memory_operator_resolver.clone())
.with_mob_resolver(agent_memory_mob_resolver.clone()),
))
} else {
base_customizer
};
let agent_memory_injector = if let (Some(agent_memory), Some(provider)) = (
gateway_options.agent_memory.as_ref(),
agent_memory_provider.clone(),
) {
let mut injector = meerkat_mobkit::AgentMemoryRuntimeInjector::new(
provider,
agent_memory.config.clone(),
);
if let Some(tracker) = agent_memory_taint.clone() {
injector = injector.with_taint_tracker(tracker);
}
if let Some(distiller) = agent_memory_distiller.clone() {
injector = injector.with_distiller(distiller);
}
if let Some(steward) = agent_memory_steward.clone() {
injector = injector.with_steward(steward);
}
if let Some(hygienist) = agent_memory_hygienist.clone() {
injector = injector.with_hygienist(hygienist);
}
injector = injector.with_operator_resolver(agent_memory_operator_resolver.clone());
injector = injector.with_mob_resolver(agent_memory_mob_resolver.clone());
let _ = agent_memory_compaction_reset.set(injector.clone());
Some(injector)
} else {
None
};
irt.set_agent_memory(agent_memory_injector).await;
let roster_specs = roster
.roster(&RosterContext {
mob_definition: Some(mob_definition.clone()),
previous_identities: Vec::new(),
})
.await
.unwrap_or_else(|e| {
fail_init(&request_id, -32603, format!("roster provider failed: {e}"));
});
*steward_roster_slot
.lock()
.unwrap_or_else(|err| err.into_inner()) = roster_specs.clone();
let identity_context = Arc::new(IdentityFirstRuntimeContext::new_with_bootstrap_mode(
irt.clone(),
roster.clone(),
topology.clone(),
customizer.clone(),
Some(runtime.mob_handle().definition().clone()),
gateway_options
.identity_bootstrap_mode
.clone()
.unwrap_or_default(),
));
if let Err(e) = runtime
.install_and_bootstrap_identity_first_context(
Arc::clone(&identity_context),
&roster_specs,
)
.await
{
fail_init(
&request_id,
-32603,
format!("identity-first bootstrap failed: {e}"),
);
}
Some(meerkat_mobkit::rpc::IdentityFirstContext {
runtime: irt,
roster_provider: roster,
topology_provider: topology,
customizer,
agent_memory_provider,
mob_definition: Some(mob_definition),
})
} else {
None
};
let (_schedule_host, _schedule_watchdog) = if let Some((
schedule_service,
mob_target_registry,
service,
adapter,
schedule_store_path,
)) = schedule_host_inputs
{
let mob_state = runtime.mob_runtime().agent_mob_mcp_state();
mob_target_registry.set_mob_state(mob_state.clone());
match meerkat_mobkit::schedule_wiring::repair_resumable_session_targets_to_mob_members(
&schedule_service,
&mob_target_registry,
)
.await
{
Ok(repaired) if repaired > 0 => {
tracing::info!(
repaired,
"repaired persisted resumable-session schedules to identity mob targets"
);
}
Ok(_) => {}
Err(error) => {
tracing::warn!(
error = %error,
"failed to repair persisted resumable-session schedules to identity mob targets",
);
}
}
let callback_runnables = (!gateway_options.host_runnables.is_empty()).then(|| {
(
Arc::new(bridge.clone())
as Arc<dyn meerkat_mobkit::identity_first::gateway_bridges::CallbackBridge>,
gateway_options.host_runnables.clone(),
)
});
let runnable_host = match meerkat_mobkit::schedule_wiring::gateway_runnable_host(
agent_memory_steward.clone(),
callback_runnables,
) {
Ok(host) => host,
Err(error) => {
fail_init(
&request_id,
-32602,
format!("failed to compose schedule host runnables: {error}"),
);
}
};
if let Some(steward) = agent_memory_steward.as_ref()
&& let Err(error) = meerkat_mobkit::schedule_wiring::ensure_steward_dream_schedule(
&schedule_service,
steward.dream_cadence(),
chrono::Utc::now(),
)
.await
{
tracing::warn!(
error = %error,
"failed to ensure steward dream schedule; the steward will not dream on this gateway",
);
}
let watchdog = meerkat_mobkit::schedule_wiring::spawn_schedule_claim_watchdog(
schedule_service.clone(),
schedule_store_path,
Default::default(),
);
(
meerkat_mobkit::schedule_wiring::spawn_schedule_host_with_identity_runtime(
service,
adapter,
schedule_service,
mob_state,
runtime.mob_handle(),
runtime.identity_runtime().cloned(),
runnable_host,
workgraph_service.clone(),
schedule_owner_id.clone(),
),
Some(watchdog),
)
} else {
if !gateway_options.host_runnables.is_empty() {
tracing::warn!(
"runtime_options.host_runnables is configured but this gateway runs no \
schedule host (ephemeral mode, or schedule store unavailable): \
host-runnable schedule targets will never fire"
);
}
(None, None)
};
let runtime = Arc::new(runtime);
if let Some(detached_jobs) = gateway_detached_jobs.as_ref() {
match detached_jobs.health_projection().await {
Ok(projection) => runtime.set_job_health_projection(Some(projection)),
Err(error) => {
tracing::warn!(%error, "initial durable callback health projection failed");
runtime.set_job_health_projection(Some(json!({
"status": "degraded",
"detached_jobs": {
"status": "degraded",
"reason": "job_health_projection_failed"
}
})));
}
}
detached_jobs.arm_delivery_driver(Arc::clone(&runtime));
}
steward_late_runtime.bind(runtime.clone());
if let Some(steward) = agent_memory_steward.clone() {
runtime
.register_gating_resolution_observer(Arc::new(
meerkat_mobkit::PromotionGateResolver::new(
steward,
tokio::runtime::Handle::current(),
),
))
.await;
}
let event_drain_task = runtime.clone().spawn_event_drain_task();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind ephemeral port");
let port = listener.local_addr().expect("local addr").port();
let http_base_url = format!("http://127.0.0.1:{port}");
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let mut decision_state = gateway_options
.decisions
.clone()
.unwrap_or_else(minimal_decision_state);
decision_state.console.ui = gateway_options.console_ui.clone();
let app = runtime.build_reference_app_router(decision_state);
let (app, live_rpc) =
if let Some((live_service, live_machine, live_agent_factory)) = live_inputs {
let ws_base_url = match &gateway_options.live {
GatewayLiveOption::Enabled {
public_base_url: Some(public),
..
} => public.trim_end_matches('/').to_string(),
_ => format!("ws://127.0.0.1:{port}"),
};
let live_seed_max_chars = match &gateway_options.live {
GatewayLiveOption::Enabled { seed_max_chars, .. } => *seed_max_chars,
GatewayLiveOption::Disabled => None,
};
let live_ctx = Arc::new(meerkat_mobkit::live_wiring::attach_live(
Arc::clone(&live_service),
Arc::clone(&live_machine),
&live_agent_factory,
meerkat::Config::default(),
ws_base_url,
live_seed_max_chars,
));
let app = app.merge(meerkat_live::live_ws_router(Arc::clone(&live_ctx.ws_state)));
let handler =
meerkat_mobkit::live_wiring::live_rpc_handler(live_ctx, live_service, live_machine);
(app, Some(handler))
} else {
(app, None)
};
let mut serve_task = tokio::spawn({
let mut shutdown_rx = shutdown_rx.clone();
async move {
axum::serve(listener, app)
.with_graceful_shutdown(async move {
shutdown_rx.changed().await.ok();
})
.await
}
});
let loaded_modules = runtime.loaded_modules().await;
let runtime_origin = if is_workspace_config {
"workspace_config"
} else {
"fallback_minimal"
};
let runtime_fingerprint = {
let mut hasher = Sha256::new();
hasher.update(mob_config_toml.as_bytes());
hasher.update(
serde_json::to_string(&loaded_modules)
.unwrap_or_default()
.as_bytes(),
);
format!("{:x}", hasher.finalize())
};
let identity_bootstrap = identity_ctx
.as_ref()
.map(|ctx| ctx.runtime.identity_bootstrap_status());
let init_response = json!({
"jsonrpc": "2.0",
"id": request_id,
"result": {
"http_base_url": http_base_url,
"loaded_modules": loaded_modules,
"contract_version": MOBKIT_CONTRACT_VERSION,
"runtime_origin": runtime_origin,
"runtime_fingerprint": runtime_fingerprint,
"identity_bootstrap": identity_bootstrap,
"stdio_shutdown_handshake": true,
"stdio_shutdown_horizon_ms": GATEWAY_SHUTDOWN_HORIZON_MS,
}
});
let _ = stdout_tx
.send(
serde_json::to_string(&init_response)
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string()),
)
.await;
let identity_ctx = identity_ctx.map(Arc::new);
let http_base_url_shared: Arc<str> = http_base_url.clone().into();
let mut interrupted_with_open_stdin = false;
let mut gateway_shutdown = None;
{
let mut inflight = tokio::task::JoinSet::new();
loop {
let request_line = tokio::select! {
line = rpc_rx.recv() => line,
_ = tokio::signal::ctrl_c() => {
interrupted_with_open_stdin = true;
None
},
};
let Some(request_line) = request_line else {
break; };
if let Ok(message) = serde_json::from_str::<Value>(&request_line)
&& let Some(request) = gateway_shutdown_request(&message)
{
gateway_shutdown = Some(request);
break;
}
let request_line = apply_gateway_runtime_config_to_request(
&request_line,
&gateway_options.schedules,
&gateway_options.gating,
);
let runtime = runtime.clone();
let stdout_tx = stdout_tx.clone();
let identity_ctx = identity_ctx.clone();
let http_base_url = http_base_url_shared.clone();
let live_rpc = live_rpc.clone();
let gateway_detached_jobs = gateway_detached_jobs.clone();
let callback_bridge = bridge.clone();
inflight.spawn(async move {
let response = match handle_callback_job_rpc(
&request_line,
gateway_detached_jobs.as_ref(),
&callback_bridge,
)
.await
{
Some(response) => response,
None => {
meerkat_mobkit::rpc::handle_unified_rpc_json_with_live_arc(
&runtime,
&request_line,
timeout,
Some(http_base_url.as_ref()),
identity_ctx.as_deref(),
live_rpc.as_ref(),
)
.await
}
};
if !response.is_empty() {
let _ = stdout_tx.send(response).await;
}
});
while inflight.try_join_next().is_some() {}
}
if gateway_shutdown.is_none() {
bridge.close().await;
}
let drain = async { while inflight.join_next().await.is_some() {} };
if tokio::time::timeout(GATEWAY_RPC_DRAIN_TIMEOUT, drain)
.await
.is_err()
{
inflight.shutdown().await;
}
}
if gateway_shutdown.is_none() {
stdin_reader.abort();
}
let _ = shutdown_tx.send(true);
if tokio::time::timeout(GATEWAY_HTTP_DRAIN_TIMEOUT, &mut serve_task)
.await
.is_err()
{
serve_task.abort();
let _ = serve_task.await;
}
event_drain_task.abort();
let runtime_shutdown =
tokio::time::timeout(GATEWAY_RUNTIME_SHUTDOWN_TIMEOUT, runtime.shutdown()).await;
match runtime_shutdown.as_ref() {
Err(_) => {
tracing::warn!(
timeout_ms = GATEWAY_RUNTIME_SHUTDOWN_TIMEOUT.as_millis(),
"gateway runtime shutdown exceeded its bounded horizon"
);
}
Ok(report) if !report.cleanup_completed() => {
tracing::warn!(
drain_timed_out = report.drain.timed_out,
mob_stop = ?report.mob_stop,
identity_authority_release = ?report.identity_authority_release,
orphan_processes = report.module_shutdown.orphan_processes,
"gateway runtime shutdown completed without cleanup attestation"
);
}
Ok(_) => {}
}
bridge.close().await;
if let Some(request) = gateway_shutdown {
let response =
gateway_shutdown_response(request.response_id, runtime_shutdown.as_ref().ok());
if let Ok(line) = serde_json::to_string(&response) {
let _ = stdout_tx.send(line).await;
}
stdin_reader.abort();
}
drop(stdout_tx);
let _ = stdout_shutdown_tx.send(());
let _ = stdin_reader.await;
if tokio::time::timeout(GATEWAY_STDOUT_DRAIN_TIMEOUT, &mut stdout_writer)
.await
.is_err()
{
stdout_writer.abort();
let _ = stdout_writer.await;
}
drop(_temp_dir);
if interrupted_with_open_stdin {
std::process::exit(0);
}
}
fn main() {
let args: Vec<String> = std::env::args().collect();
if args.iter().skip(1).any(|a| a == "--version" || a == "-V") {
println!(
"rpc_gateway {} (meerkat-mobkit SDK stdin-RPC gateway)",
env!("CARGO_PKG_VERSION")
);
return;
}
if args.iter().any(|a| a == "--persistent") {
run_persistent();
} else {
run_single_shot();
}
}
#[derive(Clone, Default)]
struct StewardLateRuntime(Arc<tokio::sync::OnceCell<Arc<UnifiedRuntime>>>);
impl StewardLateRuntime {
fn bind(&self, runtime: Arc<UnifiedRuntime>) {
let _ = self.0.set(runtime);
}
fn get(&self) -> Option<Arc<UnifiedRuntime>> {
self.0.get().cloned()
}
}
struct GatewayMobPurposeSource {
mob: String,
roster: Arc<std::sync::Mutex<Vec<meerkat_mobkit::identity_first::DurableAgentSpec>>>,
}
impl meerkat_mobkit::MobPurposeSource for GatewayMobPurposeSource {
fn mob_contexts(&self) -> Vec<meerkat_mobkit::memory::steward::MobContext> {
let roster = self.roster.lock().unwrap_or_else(|err| err.into_inner());
let purpose = roster.iter().find_map(|spec| {
spec.labels
.get("mob_purpose")
.or_else(|| spec.labels.get("purpose"))
.cloned()
});
let member_labels = roster
.iter()
.map(|spec| (spec.identity.as_str().to_string(), spec.labels.clone()))
.collect();
vec![meerkat_mobkit::memory::steward::MobContext {
mob: self.mob.clone(),
purpose,
member_labels,
}]
}
}
struct GatewayMemoryGatingBridge {
runtime: StewardLateRuntime,
}
#[async_trait]
impl meerkat_mobkit::MemoryGatingBridge for GatewayMemoryGatingBridge {
async fn enqueue_promotion_gate(
&self,
realm: &str,
description: &str,
entity: &str,
topic: &str,
) -> Result<String, String> {
let Some(runtime) = self.runtime.get() else {
return Err("runtime not yet bound".to_string());
};
let result = runtime
.evaluate_gating_action(meerkat_mobkit::runtime::GatingEvaluateRequest {
action: description.to_string(),
actor_id: format!("memory-steward:{realm}"),
risk_tier: meerkat_mobkit::runtime::GatingRiskTier::R3,
rationale: Some(
"memory steward quarantine promotion (agent-memory §10.2)".to_string(),
),
requested_approver: None,
approval_recipient: None,
approval_channel: None,
approval_timeout_ms: None,
entity: Some(entity.to_string()),
topic: Some(topic.to_string()),
})
.await;
result.pending_id.ok_or_else(|| {
format!(
"gating evaluation returned outcome {:?} without a pending entry{}",
result.outcome,
result
.fallback_reason
.as_deref()
.map(|reason| format!(" ({reason})"))
.unwrap_or_default()
)
})
}
}
struct GatewayMemoryConflictBridge {
runtime: StewardLateRuntime,
handle: tokio::runtime::Handle,
}
impl meerkat_mobkit::MemoryConflictBridge for GatewayMemoryConflictBridge {
fn emit_conflict(&self, entity: &str, topic: &str, reason: &str) {
let Some(runtime) = self.runtime.get() else {
tracing::warn!(
entity,
topic,
"memory conflict bridge: runtime not yet bound"
);
return;
};
let entity = entity.to_string();
let topic = topic.to_string();
let reason = reason.to_string();
self.handle.spawn(async move {
let request = meerkat_mobkit::runtime::MemoryIndexRequest {
entity: entity.clone(),
topic: topic.clone(),
store: None,
fact: None,
metadata: None,
conflict: Some(true),
conflict_reason: Some(reason),
};
if let Err(err) = runtime.memory_index(request).await {
tracing::warn!(
entity,
topic,
error = ?err,
"memory conflict bridge: conflict signal write failed"
);
}
});
}
}