use std::collections::{BTreeMap, BTreeSet};
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::console_aggregator::is_implicit_delegate_member;
use crate::mob_handle_runtime::{topology_restore_failed_peer_ids, topology_restore_warning_json};
use crate::runtime::{
BigQuerySessionStoreAdapter, BigQuerySessionStoreError, ConsoleRestJsonRequest,
ConsoleRestJsonResponse, DeliveryHistoryRequest, DeliverySendError, DeliverySendRequest,
GatingDecideError, GatingDecideRequest, GatingDecision, GatingEvaluateRequest, GatingRiskTier,
LocalJsonMemoryStoreError, MemoryIndexError, MemoryIndexRequest, MemoryQueryRequest,
MobkitRuntimeHandle, ModuleRouteError, ModuleRouteRequest, ROUTING_RETRY_MAX_CAP,
RoutingResolveError, RoutingResolveRequest, RuntimeDecisionState, RuntimeRoute,
RuntimeRouteMutationError, ScheduleDefinition, ScheduleValidationError, SessionPersistenceRow,
SubscribeError, SubscribeRequest, SubscribeScope, handle_console_rest_json_route,
route_module_call, validate_schedules,
};
use crate::unified_runtime::{EventQuery, UnifiedRuntime};
mod console_ingress;
mod gating_methods;
pub(crate) mod memory_methods;
pub(crate) mod mob_methods;
pub(crate) mod params;
mod routing_delivery_methods;
mod scheduling_methods;
mod session_store_methods;
pub(crate) mod storage_methods;
mod subscribe_methods;
pub(crate) mod topology_methods;
pub(crate) mod workgraph_methods;
pub use console_ingress::handle_console_ingress_json;
use gating_methods::{
GatingParamsError, parse_gating_audit_params, parse_gating_decide_params,
parse_gating_evaluate_params, parse_gating_pending_params,
};
use memory_methods::{
MemoryParamsError, parse_agent_memory_forget_params, parse_agent_memory_manifest_params,
parse_agent_memory_recall_params, parse_agent_memory_remember_params,
parse_agent_memory_update_params, parse_memory_index_params, parse_memory_query_params,
parse_memory_stores_params,
};
use routing_delivery_methods::{
RoutingDeliveryParamsError, parse_delivery_history_params, parse_delivery_send_params,
parse_routing_resolve_params, parse_routing_route_add_params,
parse_routing_route_delete_params, parse_routing_routes_list_params,
};
use scheduling_methods::{format_schedule_validation_error, parse_scheduling_params};
use session_store_methods::{
BigQuerySessionStoreRpcError, format_bigquery_store_error, parse_bigquery_session_store_params,
run_bigquery_session_store_request,
};
use subscribe_methods::{SubscribeParamsError, parse_subscribe_request};
pub const JSONRPC_VERSION: &str = "2.0";
pub const MOBKIT_CONTRACT_VERSION: &str = "0.4.0";
pub const MAX_SCHEDULES_PER_REQUEST: usize = 256;
pub(crate) const MOBPACK_AUTHORING_METHODS: &[&str] = &[
"mobkit/mobpacks/schema",
"mobkit/mobpacks/catalogs",
"mobkit/tools/catalog",
"mobkit/skills/catalog",
"mobkit/agent_definitions/list",
"mobkit/mobpacks/templates",
"mobkit/mobpacks/validate",
"mobkit/mobpacks/source",
"mobkit/mobpacks/export",
"mobkit/mobpacks/import",
"mobkit/mobpacks/list",
"mobkit/mobpacks/get",
"mobkit/mobpacks/create",
"mobkit/mobpacks/save",
"mobkit/mobpacks/delete",
"mobkit/mobpacks/undo",
"mobkit/mobpacks/redo",
"mobkit/mobpacks/apply_operation",
"mobkit/mobpacks/graph_projection",
"mobkit/mobpacks/graph_to_flow",
"mobkit/mobpacks/deploy_command",
"mobkit/mobpacks/deploy",
];
pub(crate) fn mobpack_authoring_capabilities() -> Value {
serde_json::json!({
"domain": "mobpack_authoring",
"runtime_mutation": false,
"host_mutation_methods": {
"mobkit/mobpacks/deploy": "when execute=true, writes a mobpack archive and runs rkat mob run on the host",
"mobkit/mobpacks/validate": "when rkat_validate=true, writes a mobpack archive and runs rkat mob validate on the host"
},
"deploy_command": "rkat mob run",
"methods": MOBPACK_AUTHORING_METHODS,
"operations": crate::mobpack::mobpack_authoring_operations(),
})
}
async fn mobpack_runtime_catalog_state(
runtime: &UnifiedRuntime,
) -> crate::mobpack::MobpackRuntimeCatalogState {
let loaded_modules = runtime.loaded_modules().await;
let runtime_flow_rows = crate::mobpack::runtime_flow_registry_rows_from_definition(
runtime.mob_handle().definition(),
);
let runtime_agent_definition_sources =
crate::mobpack::runtime_agent_definition_sources_from_definition(
runtime.mob_handle().definition(),
);
let runtime_skill_realms =
crate::mobpack::runtime_skill_realms_from_definition(runtime.mob_handle().definition());
let mut runtime_methods = vec![
"mobkit/capabilities".to_string(),
"mobkit/models/catalog".to_string(),
"mobkit/spawn_member".to_string(),
"mobkit/list_members".to_string(),
"mobkit/get_member".to_string(),
"mobkit/run_flow".to_string(),
"mobkit/list_flows".to_string(),
"mobkit/list_runs".to_string(),
];
runtime_methods.extend(
MOBPACK_AUTHORING_METHODS
.iter()
.map(std::string::ToString::to_string),
);
if runtime.has_contact_directory() {
runtime_methods.push("mobkit/cross_mob/directory".to_string());
}
if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
runtime_methods.extend([
"mobkit/cross_mob/wire".to_string(),
"mobkit/cross_mob/unwire".to_string(),
"mobkit/cross_mob/send".to_string(),
]);
}
crate::mobpack::MobpackRuntimeCatalogState {
loaded_modules,
runtime_methods,
has_contact_directory: runtime.has_contact_directory(),
has_peer_mob_handles: runtime.has_peer_mob_handles().await,
has_inproc_contacts: runtime.has_inproc_contacts(),
runtime_flow_rows,
runtime_agent_definition_sources,
runtime_skill_realms,
}
}
async fn handle_unified_mobpack_authoring_rpc(
runtime: &UnifiedRuntime,
method: &str,
params: &Value,
response_id: Value,
) -> JsonRpcResponse {
let runtime_catalog_state = match method {
"mobkit/mobpacks/schema"
| "mobkit/mobpacks/catalogs"
| "mobkit/tools/catalog"
| "mobkit/skills/catalog"
| "mobkit/agent_definitions/list"
| "mobkit/mobpacks/templates"
| "mobkit/mobpacks/list"
| "mobkit/mobpacks/get"
| "mobkit/mobpacks/apply_operation" => Some(mobpack_runtime_catalog_state(runtime).await),
_ => None,
};
let result = match method {
"mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
runtime_catalog_state.as_ref(),
)),
"mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
runtime_catalog_state.as_ref(),
)),
"mobkit/skills/catalog" => Ok(
crate::mobpack::mobpack_skills_catalog_response_with_runtime(
runtime_catalog_state.as_ref(),
),
),
"mobkit/agent_definitions/list" => Ok(
crate::mobpack::mobpack_agent_definitions_response_with_runtime(
runtime_catalog_state.as_ref(),
),
),
"mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
runtime_catalog_state.as_ref(),
)),
_ => {
return handle_mobpack_authoring_rpc_with_runtime(
method,
params,
response_id.clone(),
runtime_catalog_state.as_ref(),
)
.unwrap_or_else(|| JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: "Method not found".to_string(),
data: None,
}),
});
}
};
match result {
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
},
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message,
data: None,
}),
},
}
}
pub(crate) fn handle_mobpack_authoring_rpc(
method: &str,
params: &Value,
response_id: Value,
) -> Option<JsonRpcResponse> {
handle_mobpack_authoring_rpc_with_runtime(method, params, response_id, None)
}
pub(crate) fn handle_mobpack_authoring_rpc_with_runtime(
method: &str,
params: &Value,
response_id: Value,
runtime: Option<&crate::mobpack::MobpackRuntimeCatalogState>,
) -> Option<JsonRpcResponse> {
let result = match method {
"mobkit/mobpacks/schema" => Ok(crate::mobpack::mobpack_schema_response_with_runtime(
runtime,
)),
"mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
runtime,
)),
"mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
runtime,
)),
"mobkit/skills/catalog" => {
Ok(crate::mobpack::mobpack_skills_catalog_response_with_runtime(runtime))
}
"mobkit/agent_definitions/list" => {
Ok(crate::mobpack::mobpack_agent_definitions_response_with_runtime(runtime))
}
"mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
runtime,
)),
"mobkit/mobpacks/validate" => crate::mobpack::validate_mobpack(params)
.and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
"mobkit/mobpacks/source" => crate::mobpack::source_mobpack(params)
.and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
"mobkit/mobpacks/export" => crate::mobpack::export_mobpack(params)
.and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
"mobkit/mobpacks/import" => crate::mobpack::import_mobpack(params),
"mobkit/mobpacks/list" => crate::mobpack::list_mobpack_drafts_with_runtime(params, runtime),
"mobkit/mobpacks/get" => crate::mobpack::get_mobpack_draft_with_runtime(params, runtime),
"mobkit/mobpacks/create" => crate::mobpack::create_mobpack_draft(params),
"mobkit/mobpacks/save" => crate::mobpack::save_mobpack_draft(params),
"mobkit/mobpacks/delete" => crate::mobpack::delete_mobpack_draft(params),
"mobkit/mobpacks/undo" => crate::mobpack::undo_mobpack_draft(params),
"mobkit/mobpacks/redo" => crate::mobpack::redo_mobpack_draft(params),
"mobkit/mobpacks/apply_operation" => {
crate::mobpack::apply_mobpack_authoring_operation_with_runtime(params, runtime)
}
"mobkit/mobpacks/graph_projection" => crate::mobpack::graph_projection_mobpack(params),
"mobkit/mobpacks/graph_to_flow" => crate::mobpack::graph_to_flow_mobpack(params),
"mobkit/mobpacks/deploy_command" => crate::mobpack::deploy_command_preview(params)
.and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
"mobkit/mobpacks/deploy" => crate::mobpack::deploy_mobpack(params)
.and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
_ => return None,
};
Some(match result {
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
},
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message,
data: None,
}),
},
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RpcCapabilitiesError {
InvalidJson,
InvalidSchema,
MissingContractVersion,
InvalidContractVersion,
}
impl std::fmt::Display for RpcCapabilitiesError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::InvalidJson => write!(f, "invalid JSON"),
Self::InvalidSchema => write!(f, "invalid schema"),
Self::MissingContractVersion => write!(f, "missing contract version"),
Self::InvalidContractVersion => write!(f, "invalid contract version"),
}
}
}
impl std::error::Error for RpcCapabilitiesError {}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RpcCapabilities {
pub contract_version: String,
#[serde(flatten)]
pub extra: BTreeMap<String, Value>,
}
pub fn parse_rpc_capabilities(line: &str) -> Result<RpcCapabilities, RpcCapabilitiesError> {
let raw: Value = serde_json::from_str(line).map_err(|_| RpcCapabilitiesError::InvalidJson)?;
let object = raw.as_object().ok_or(RpcCapabilitiesError::InvalidSchema)?;
let contract = object
.get("contract_version")
.ok_or(RpcCapabilitiesError::MissingContractVersion)?;
let contract_str = contract
.as_str()
.ok_or(RpcCapabilitiesError::InvalidContractVersion)?;
if contract_str.trim().is_empty() {
return Err(RpcCapabilitiesError::InvalidContractVersion);
}
serde_json::from_value(raw).map_err(|_| RpcCapabilitiesError::InvalidSchema)
}
pub const MOB_EVENTS_STALE_CURSOR_CODE: i64 = -32010;
pub const MEMORY_BACKEND_UNAVAILABLE_CODE: i64 = -32012;
pub const CONSOLE_TIMELINE_REPLAY_UNAVAILABLE_CODE: i64 = -32013;
pub const STORAGE_RESOLUTION_CODE: i64 = -32014;
pub const WORKGRAPH_UNAVAILABLE_CODE: i64 = -32041;
pub const WORKGRAPH_CONFLICT_CODE: i64 = -32042;
pub const WORKGRAPH_ERROR_CODE: i64 = -32000;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct JsonRpcRequest {
pub jsonrpc: String,
#[serde(default)]
pub id: Option<Value>,
pub method: String,
#[serde(default)]
pub params: Value,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JsonRpcError {
pub code: i64,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub data: Option<Value>,
}
impl JsonRpcError {
pub fn new(code: i64, message: impl Into<String>) -> Self {
Self {
code,
message: message.into(),
data: None,
}
}
#[must_use]
pub fn with_data(mut self, data: Value) -> Self {
self.data = Some(data);
self
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct JsonRpcResponse {
pub jsonrpc: String,
pub id: Value,
#[serde(skip_serializing_if = "Option::is_none")]
pub result: Option<Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<JsonRpcError>,
}
pub fn handle_mobkit_rpc_json(
runtime: &mut MobkitRuntimeHandle,
request_json: &str,
timeout: Duration,
) -> String {
let raw_request: Value = match serde_json::from_str(request_json) {
Ok(raw_request) => raw_request,
Err(_) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: Value::Null,
result: None,
error: Some(JsonRpcError {
code: -32700,
message: "Parse error".to_string(),
data: None,
}),
});
}
};
let response_id = raw_request
.as_object()
.and_then(|object| object.get("id"))
.cloned()
.unwrap_or(Value::Null);
let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
Ok(request) => request,
Err(_) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32600,
message: "Invalid Request".to_string(),
data: None,
}),
});
}
};
let is_notification = request.id.is_none();
let response_id = request.id.clone().unwrap_or(Value::Null);
if request.jsonrpc != "2.0" {
let response = JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32600,
message: "Invalid Request".to_string(),
data: None,
}),
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
if let Err(message) =
crate::member_comms_id::validate_public_rpc_member_aliases(&request.params)
{
let response = JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let response = match request.method.as_str() {
"mobkit/status" => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"contract_version": MOBKIT_CONTRACT_VERSION,
"running": runtime.is_running(),
"loaded_modules": runtime.loaded_modules(),
})),
error: None,
},
"mobkit/capabilities" => {
let mut methods = vec![
"mobkit/status",
"mobkit/capabilities",
"mobkit/reconcile",
"mobkit/spawn_member",
"mobkit/scheduling/evaluate",
"mobkit/scheduling/dispatch",
"mobkit/routing/resolve",
"mobkit/routing/routes/list",
"mobkit/routing/routes/add",
"mobkit/routing/routes/delete",
"mobkit/delivery/send",
"mobkit/delivery/history",
"mobkit/events/subscribe",
"mobkit/memory/stores",
"mobkit/memory/index",
"mobkit/memory/query",
"mobkit/session_store/bigquery",
"mobkit/gating/evaluate",
"mobkit/gating/pending",
"mobkit/gating/decide",
"mobkit/gating/audit",
"mobkit/call_tool",
"mobkit/models/catalog",
storage_methods::STORAGE_DOCTOR_METHOD,
];
methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"contract_version": MOBKIT_CONTRACT_VERSION,
"methods": methods,
"loaded_modules": runtime.loaded_modules(),
"runtime_capabilities": {
"can_spawn_members": false,
"can_send_messages": false,
"can_wire_members": false,
"can_retire_members": false,
"available_spawn_modes": ["module"],
},
"authoring_capabilities": mobpack_authoring_capabilities(),
})),
error: None,
}
}
"mobkit/models/catalog" => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(build_models_catalog_result()),
error: None,
},
storage_methods::STORAGE_DOCTOR_METHOD => {
match storage_methods::parse_storage_doctor_params(&request.params) {
Ok(Some(params)) => {
let diagnosis = crate::storage_doctor::diagnose_state_dir_blocking_with_options(
¶ms.scope(),
None,
params.doctor_options(),
);
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(storage_methods::storage_doctor_result_json(
¶ms, &diagnosis, None,
)),
error: None,
}
}
Ok(None) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: state_dir required (the module-only runtime \
has no state directory)"
.to_string(),
data: None,
}),
},
Err(reason) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {reason}"),
data: None,
}),
},
}
}
method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
handle_mobpack_authoring_rpc(method, &request.params, response_id.clone())
.unwrap_or_else(|| JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: "Method not found".to_string(),
data: None,
}),
})
}
"mobkit/reconcile" => {
let modules = match params::required_string_array(&request.params, "modules") {
Ok(m) => m,
Err(reason) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {reason}"),
data: None,
}),
});
}
};
match runtime.reconcile_modules(modules.clone(), timeout) {
Ok(added) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"accepted": true,
"reconciled_modules": modules,
"added": added
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err:?}"),
data: None,
}),
},
}
}
"mobkit/spawn_member" => {
let module_id = request
.params
.get("module_id")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
if module_id.is_empty() {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: module_id required".to_string(),
data: None,
}),
}
} else {
match runtime.spawn_member(&module_id, timeout) {
Ok(()) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"accepted": true,
"module_id": module_id
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err:?}"),
data: None,
}),
},
}
}
}
"mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
Ok((schedules, tick_ms)) => match runtime.evaluate_schedule_tick(&schedules, tick_ms) {
Ok(evaluation) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"Invalid params: {}",
format_schedule_validation_error(err)
),
data: None,
}),
},
},
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
},
"mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
Ok((schedules, tick_ms)) => match runtime.dispatch_schedule_tick(&schedules, tick_ms) {
Ok(dispatch) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"Invalid params: {}",
format_schedule_validation_error(err)
),
data: None,
}),
},
},
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
},
"mobkit/routing/resolve" => {
match parse_routing_resolve_params(&request.params).and_then(|resolve_request| {
runtime
.resolve_routing(resolve_request)
.map_err(RoutingDeliveryParamsError::Routing)
}) {
Ok(resolution) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
Ok(()) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"routes": runtime.list_runtime_routes()
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/routing/routes/add" => match parse_routing_route_add_params(&request.params)
.and_then(|route| {
runtime
.add_runtime_route(route)
.map_err(RoutingDeliveryParamsError::RouteMutation)
}) {
Ok(route) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "route": route })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/routing/routes/delete" => match parse_routing_route_delete_params(&request.params)
.and_then(|route_key| {
runtime
.delete_runtime_route(&route_key)
.map_err(RoutingDeliveryParamsError::RouteMutation)
}) {
Ok(route) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "deleted": route })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/delivery/send" => {
match parse_delivery_send_params(&request.params).and_then(|send_request| {
runtime
.send_delivery(send_request)
.map_err(RoutingDeliveryParamsError::Delivery)
}) {
Ok(record) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
Ok(history_request) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(
serde_json::to_value(runtime.delivery_history(history_request))
.unwrap_or(Value::Null),
),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/events/subscribe" => {
match parse_subscribe_request(&request.params).and_then(|subscribe_request| {
runtime
.subscribe_events(subscribe_request)
.map_err(SubscribeParamsError::Runtime)
}) {
Ok(subscribe_result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
Ok(()) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"stores": runtime.memory_stores(),
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/memory/index" => match parse_memory_index_params(&request.params) {
Ok(index_request) => match runtime.memory_index(index_request) {
Ok(indexed) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
error: None,
},
Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: MEMORY_BACKEND_UNAVAILABLE_CODE,
message: format!(
"Memory backend unavailable: {}",
MemoryParamsError::backend_message(&error)
),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"Invalid params: {}",
MemoryParamsError::Index(err).message()
),
data: None,
}),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/memory/query" => match parse_memory_query_params(&request.params) {
Ok(query_request) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(
serde_json::to_value(runtime.memory_query(query_request))
.unwrap_or(Value::Null),
),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/session_store/bigquery" => {
match parse_bigquery_session_store_params(&request.params)
.and_then(run_bigquery_session_store_request)
{
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
},
Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32011,
message: format!(
"BigQuery session store request failed: {}",
format_bigquery_store_error(&error)
),
data: None,
}),
},
}
}
"mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
Ok(gating_request) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(
serde_json::to_value(runtime.evaluate_gating_action(gating_request))
.unwrap_or(Value::Null),
),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
Ok(()) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"pending": runtime.list_gating_pending(),
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/gating/decide" => {
match parse_gating_decide_params(&request.params).and_then(|decide_request| {
runtime
.decide_gating_action(decide_request)
.map_err(GatingParamsError::Decision)
}) {
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
Ok(limit) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"entries": runtime.gating_audit_entries(limit),
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/call_tool" => {
let module_id = request.params.get("module_id").and_then(Value::as_str);
let tool = request.params.get("tool").and_then(Value::as_str);
let arguments = request
.params
.get("arguments")
.cloned()
.unwrap_or(serde_json::json!({}));
match (module_id, tool) {
(Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
let route = route_module_call(
runtime,
&ModuleRouteRequest {
module_id: module_id.to_string(),
method: tool.to_string(),
params: arguments,
},
timeout,
);
match route {
Ok(response) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"module_id": response.module_id,
"tool": response.method,
"result": response.payload
})),
error: None,
},
Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: format!("Module '{mid}' not loaded"),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("Tool call failed: {err:?}"),
data: None,
}),
},
}
}
_ => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: module_id and tool required".to_string(),
data: None,
}),
},
}
}
method if method.contains('/') && !method.starts_with("mobkit/") => {
let module_id = method
.split('/')
.next()
.map(ToString::to_string)
.unwrap_or_default();
let route = route_module_call(
runtime,
&ModuleRouteRequest {
module_id,
method: method.to_string(),
params: request.params,
},
timeout,
);
match route {
Ok(response) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"module_id": response.module_id,
"method": response.method,
"payload": response.payload
})),
error: None,
},
Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: format!("Module '{module_id}' not loaded"),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("Module route failed: {err:?}"),
data: None,
}),
},
}
}
_ => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: "Method not found".to_string(),
data: None,
}),
},
};
if is_notification {
String::new()
} else {
serialize_response(&response)
}
}
pub struct IdentityFirstContext {
pub runtime: std::sync::Arc<crate::identity_first::IdentityRuntime>,
pub roster_provider: std::sync::Arc<dyn crate::identity_first::contracts::RosterProvider>,
pub topology_provider:
Option<std::sync::Arc<dyn crate::identity_first::contracts::TopologyProvider>>,
pub customizer: Option<std::sync::Arc<dyn crate::identity_first::contracts::AgentCustomizer>>,
pub agent_memory_provider:
Option<std::sync::Arc<dyn crate::identity_first::AgentMemoryProvider>>,
pub mob_definition: Option<meerkat_mob::MobDefinition>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum IdentityMemberReadiness {
Ready,
TimedOut,
Failed(String),
}
struct IdentityStartupReadyWait {
status: crate::identity_first::IdentityBootstrapStatus,
timed_out: bool,
startup_ready: bool,
}
async fn wait_identity_startup_ready<F, Fut>(
identity_rt: &crate::identity_first::IdentityRuntime,
wait_timeout: Duration,
mut wait_members: F,
) -> Result<IdentityStartupReadyWait, String>
where
F: FnMut(Vec<meerkat_mob::ids::AgentIdentity>, Duration) -> Fut + Send,
Fut: Future<Output = IdentityMemberReadiness> + Send,
{
let started = std::time::Instant::now();
let (mut status, mut timed_out, mut generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(wait_timeout)
.await;
loop {
if timed_out || !status.ready {
return Ok(IdentityStartupReadyWait {
status,
timed_out,
startup_ready: false,
});
}
let mut bootstrap_changes = identity_rt.subscribe_identity_bootstrap_status();
let (current_generation, current_status) =
identity_rt.identity_bootstrap_status_with_generation();
if current_generation != generation || current_status != status {
let remaining = wait_timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
return Ok(IdentityStartupReadyWait {
status: current_status,
timed_out: true,
startup_ready: false,
});
}
(status, timed_out, generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(remaining)
.await;
continue;
}
let member_ids = identity_rt
.identity_bootstrap_member_ids_for_status(&status)
.await;
let (mapped_generation, mapped_status) =
identity_rt.identity_bootstrap_status_with_generation();
if mapped_generation != generation || mapped_status != status {
let remaining = wait_timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
return Ok(IdentityStartupReadyWait {
status: mapped_status,
timed_out: true,
startup_ready: false,
});
}
(status, timed_out, generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(remaining)
.await;
continue;
}
if member_ids.len() != status.identities.len() {
return Err("identity bootstrap readiness mapping is incomplete".to_string());
}
let remaining = wait_timeout.saturating_sub(started.elapsed());
let readiness = wait_members(member_ids, remaining);
tokio::pin!(readiness);
let readiness_result = tokio::select! {
result = &mut readiness => Some(result),
changed = bootstrap_changes.changed() => {
if changed.is_err() {
return Err(
"identity bootstrap readiness status channel closed".to_string(),
);
}
None
}
};
let Some(readiness_result) = readiness_result else {
let remaining = wait_timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
let (_, latest) = identity_rt.identity_bootstrap_status_with_generation();
return Ok(IdentityStartupReadyWait {
status: latest,
timed_out: true,
startup_ready: false,
});
}
(status, timed_out, generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(remaining)
.await;
continue;
};
let (result_generation, result_status) =
identity_rt.identity_bootstrap_status_with_generation();
if result_generation != generation || result_status != status {
let remaining = wait_timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
return Ok(IdentityStartupReadyWait {
status: result_status,
timed_out: true,
startup_ready: false,
});
}
(status, timed_out, generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(remaining)
.await;
continue;
}
match readiness_result {
IdentityMemberReadiness::Ready => {
let (ready_generation, ready_status) =
identity_rt.identity_bootstrap_status_with_generation();
if ready_generation == generation && ready_status == status {
return Ok(IdentityStartupReadyWait {
status,
timed_out: false,
startup_ready: true,
});
}
let remaining = wait_timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
return Ok(IdentityStartupReadyWait {
status: ready_status,
timed_out: true,
startup_ready: false,
});
}
(status, timed_out, generation) = identity_rt
.wait_identity_bootstrap_terminal_with_generation(remaining)
.await;
}
IdentityMemberReadiness::TimedOut => {
return Ok(IdentityStartupReadyWait {
status,
timed_out: true,
startup_ready: false,
});
}
IdentityMemberReadiness::Failed(error) => {
return Err(format!("identity bootstrap readiness failed: {error}"));
}
}
}
}
pub fn handle_unified_rpc_json<'a>(
runtime: &'a UnifiedRuntime,
request_json: &'a str,
timeout: Duration,
http_base_url: Option<&'a str>,
identity_ctx: Option<&'a IdentityFirstContext>,
) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
handle_unified_rpc_json_with_live(
runtime,
request_json,
timeout,
http_base_url,
identity_ctx,
None,
)
}
pub fn handle_unified_rpc_json_arc<'a>(
runtime: &'a Arc<UnifiedRuntime>,
request_json: &'a str,
timeout: Duration,
http_base_url: Option<&'a str>,
identity_ctx: Option<&'a IdentityFirstContext>,
) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
handle_unified_rpc_json_with_live_arc(
runtime,
request_json,
timeout,
http_base_url,
identity_ctx,
None,
)
}
pub fn handle_unified_rpc_json_with_live<'a>(
runtime: &'a UnifiedRuntime,
request_json: &'a str,
timeout: Duration,
http_base_url: Option<&'a str>,
identity_ctx: Option<&'a IdentityFirstContext>,
live: Option<&'a crate::live_wiring::LiveRpcHandler>,
) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
Box::pin(handle_unified_rpc_json_inner(
runtime,
None,
request_json,
timeout,
http_base_url,
identity_ctx,
live,
))
}
pub fn handle_unified_rpc_json_with_live_arc<'a>(
runtime: &'a Arc<UnifiedRuntime>,
request_json: &'a str,
timeout: Duration,
http_base_url: Option<&'a str>,
identity_ctx: Option<&'a IdentityFirstContext>,
live: Option<&'a crate::live_wiring::LiveRpcHandler>,
) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
Box::pin(handle_unified_rpc_json_inner(
runtime.as_ref(),
Some(runtime),
request_json,
timeout,
http_base_url,
identity_ctx,
live,
))
}
async fn handle_unified_rpc_json_inner(
runtime: &UnifiedRuntime,
runtime_owner: Option<&Arc<UnifiedRuntime>>,
request_json: &str,
timeout: Duration,
http_base_url: Option<&str>,
identity_ctx: Option<&IdentityFirstContext>,
live: Option<&crate::live_wiring::LiveRpcHandler>,
) -> String {
let raw_request: Value = match serde_json::from_str(request_json) {
Ok(raw_request) => raw_request,
Err(_) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: Value::Null,
result: None,
error: Some(JsonRpcError {
code: -32700,
message: "Parse error".to_string(),
data: None,
}),
});
}
};
let response_id = raw_request
.as_object()
.and_then(|object| object.get("id"))
.cloned()
.unwrap_or(Value::Null);
let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
Ok(request) => request,
Err(_) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32600,
message: "Invalid Request".to_string(),
data: None,
}),
});
}
};
let is_notification = request.id.is_none();
let response_id = request.id.clone().unwrap_or(Value::Null);
if request.jsonrpc != "2.0" {
let response = JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32600,
message: "Invalid Request".to_string(),
data: None,
}),
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
if let Err(message) =
crate::member_comms_id::validate_public_rpc_member_aliases(&request.params)
{
let response = JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let response = match request.method.as_str() {
"mobkit/status" => {
let mob_state = Some(runtime.mob_handle().status_observation_snapshot());
let is_running = runtime.module_is_running().await;
let loaded = runtime.loaded_modules().await;
let mut result = serde_json::json!({
"contract_version": MOBKIT_CONTRACT_VERSION,
"running": is_running,
"loaded_modules": loaded,
"mob_state": format!("{mob_state:?}"),
});
if let Some(url) = http_base_url {
result["http_base_url"] = Value::String(url.to_string());
}
if let Some(ctx) = identity_ctx {
result["identity_bootstrap"] =
serde_json::to_value(ctx.runtime.identity_bootstrap_status())
.unwrap_or(Value::Null);
}
if let Some(storage) = runtime.resolved_storage() {
result["storage"] = storage.status_json();
}
if let Some(job_health) = runtime.job_health_projection() {
result["detached_jobs"] = job_health
.get("detached_jobs")
.cloned()
.unwrap_or(Value::Null);
}
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
"mobkit/capabilities" => {
let loaded = runtime.loaded_modules().await;
let mut methods = vec![
"mobkit/init",
"mobkit/status",
"mobkit/capabilities",
"mobkit/reconcile",
"mobkit/spawn_member",
"mobkit/scheduling/evaluate",
"mobkit/scheduling/dispatch",
"mobkit/routing/resolve",
"mobkit/routing/routes/list",
"mobkit/routing/routes/add",
"mobkit/routing/routes/delete",
"mobkit/delivery/send",
"mobkit/delivery/history",
"mobkit/events/subscribe",
"mobkit/query_events",
"mobkit/memory/stores",
"mobkit/memory/index",
"mobkit/memory/query",
"mobkit/session_store/bigquery",
"mobkit/gating/evaluate",
"mobkit/gating/pending",
"mobkit/gating/decide",
"mobkit/gating/audit",
"mobkit/call_tool",
"mobkit/models/catalog",
"mobkit/blob/get",
"mobkit/send_message",
"mobkit/find_members",
"mobkit/ensure_member",
"mobkit/list_members",
"mobkit/get_member",
"mobkit/retire_member",
"mobkit/respawn_member",
"mobkit/reconcile_edges",
"mobkit/rediscover",
"mobkit/mob_events/query",
"mobkit/mob_events/subscribe",
"mobkit/cross_mob/peer_info",
"mobkit/cross_mob/wire_local",
"mobkit/cross_mob/unwire_local",
"mobkit/peer_pubkey",
"mobkit/member_status",
"mobkit/identity/resolved_tools",
"mobkit/force_cancel_member",
"mobkit/spawn_helper",
"mobkit/fork_helper",
"mobkit/attach_existing_session",
"mobkit/cancel_flow",
"mobkit/flow_status",
"mobkit/list_flows",
"mobkit/list_runs",
"mobkit/run_flow",
"mobkit/collect_completed",
"mobkit/wait_ready",
"mobkit/mob_labels/set",
"mobkit/mob_labels/get",
"mobkit/mob_labels/delete",
"mobkit/run_labels/set",
"mobkit/run_labels/get",
"mobkit/run_labels/delete",
storage_methods::STORAGE_DOCTOR_METHOD,
];
methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
let workgraph_configured = runtime.workgraph_service().is_some();
if workgraph_configured {
methods.extend_from_slice(workgraph_methods::WORKGRAPH_READ_METHODS);
methods.extend_from_slice(workgraph_methods::WORKGRAPH_MUTATE_METHODS);
}
if live.is_some() {
methods.extend_from_slice(&[
"mobkit/live/open",
"mobkit/live/status",
"mobkit/live/close",
"mobkit/live/refresh",
"mobkit/live/send_input",
"mobkit/live/commit_input",
"mobkit/live/interrupt",
"mobkit/live/truncate",
]);
}
if identity_ctx.is_some() {
methods.extend_from_slice(&[
"mobkit/send",
"mobkit/interact",
"mobkit/dispatch",
"mobkit/subscribe",
"mobkit/status_identity",
"mobkit/respawn",
"mobkit/retire",
"mobkit/reset",
"mobkit/delete_identity",
"mobkit/inspect_identity",
"mobkit/reconcile_identity",
"mobkit/status_identity_bootstrap",
"mobkit/wait_identity_bootstrap",
]);
}
if identity_ctx
.and_then(|ctx| ctx.agent_memory_provider.as_ref())
.is_some()
{
methods.push("mobkit/agent_memory/recall");
if identity_ctx
.and_then(|ctx| ctx.agent_memory_provider.as_ref())
.is_some_and(|provider| provider.supports_remember())
{
methods.push("mobkit/agent_memory/remember");
}
if identity_ctx
.and_then(|ctx| ctx.agent_memory_provider.as_ref())
.is_some_and(|provider| provider.supports_forget())
{
methods.push("mobkit/agent_memory/forget");
}
if identity_ctx
.and_then(|ctx| ctx.agent_memory_provider.as_ref())
.is_some_and(|provider| provider.supports_supersede())
{
methods.push("mobkit/agent_memory/update");
}
if identity_ctx
.and_then(|ctx| ctx.agent_memory_provider.as_ref())
.is_some_and(|provider| provider.supports_manifest())
{
methods.push("mobkit/agent_memory/manifest");
}
}
if runtime.has_contact_directory() {
methods.push("mobkit/cross_mob/directory");
}
if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
methods.extend_from_slice(&[
"mobkit/cross_mob/wire",
"mobkit/cross_mob/unwire",
"mobkit/cross_mob/send",
]);
}
let job_health = runtime.job_health_projection();
if job_health.is_some() {
methods.extend_from_slice(&[
"jobs/get",
"jobs/list",
"jobs/cancel",
"jobs/progress",
"jobs/result",
"jobs/artifacts",
"jobs/retry",
"jobs/health",
"jobs/subscribe",
"jobs/unsubscribe",
]);
if job_health
.as_ref()
.and_then(|projection| projection.get("monitors_available"))
.and_then(Value::as_bool)
.unwrap_or(false)
{
methods.push("monitors/start");
}
}
let topology = runtime.topology_runtime_handle();
let (topology_methods, topology_capabilities) =
topology_methods::capability_projection(&topology, None, false);
methods.extend(topology_methods);
let storage = runtime
.resolved_storage()
.map(|summary| summary.status_json())
.unwrap_or(Value::Null);
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"contract_version": MOBKIT_CONTRACT_VERSION,
"runtime_type": "unified",
"methods": methods,
"storage": storage,
"identity_first": identity_ctx.is_some(),
"workgraph": workgraph_configured,
"detached_jobs": job_health
.as_ref()
.and_then(|projection| projection.get("detached_jobs"))
.cloned()
.unwrap_or(Value::Null),
"loaded_modules": loaded,
"runtime_capabilities": {
"can_spawn_members": true,
"can_send_messages": true,
"can_wire_members": true,
"can_retire_members": true,
"available_spawn_modes": ["module", "profile"],
},
"authoring_capabilities": mobpack_authoring_capabilities(),
"topology_control": topology_capabilities,
})),
error: None,
}
}
topology_methods::TOPOLOGY_QUERY_METHOD => {
let topology = runtime.topology_runtime_handle();
topology_methods::handle_query(&topology, response_id, None, false).await
}
topology_methods::TOPOLOGY_PLAN_METHOD => {
let topology = runtime.topology_runtime_handle();
topology_methods::handle_plan(&topology, response_id, &request.params, None).await
}
topology_methods::TOPOLOGY_APPLY_METHOD => {
let topology = runtime.topology_runtime_handle();
topology_methods::handle_apply(
&topology,
response_id,
&request.params,
None,
Some("local-host"),
)
.await
}
topology_methods::TOPOLOGY_OPERATION_METHOD => {
let topology = runtime.topology_runtime_handle();
topology_methods::handle_operation(&topology, response_id, &request.params, None).await
}
topology_methods::TOPOLOGY_AUDIT_METHOD => {
let topology = runtime.topology_runtime_handle();
topology_methods::handle_audit(&topology, response_id, &request.params, None).await
}
"mobkit/reconcile" => {
let modules = match params::required_string_array(&request.params, "modules") {
Ok(m) => m,
Err(reason) => {
return serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {reason}"),
data: None,
}),
});
}
};
match runtime.reconcile_modules(modules.clone(), timeout).await {
Ok(added) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"accepted": true,
"reconciled_modules": modules,
"added": added
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err:?}"),
data: None,
}),
},
}
}
"mobkit/spawn_member" => {
let module_id = request.params.get("module_id").and_then(Value::as_str);
let profile = request.params.get("profile").and_then(Value::as_str);
let meerkat_id = request.params.get("meerkat_id").and_then(Value::as_str);
if let Some(module_id) = module_id {
if module_id.is_empty() {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: module_id required".to_string(),
data: None,
}),
}
} else {
match runtime.spawn_member(module_id, timeout).await {
Ok(()) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"accepted": true,
"module_id": module_id
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err:?}"),
data: None,
}),
},
}
}
} else if let (Some(profile), Some(meerkat_id)) = (profile, meerkat_id) {
let target_identity_runtime = identity_ctx
.map(|ctx| &ctx.runtime)
.or_else(|| runtime.identity_runtime());
let raw_target_validation = crate::member_comms_id::validate_raw_member_target(
target_identity_runtime,
meerkat_id,
)
.await;
if meerkat_id == crate::console_contracts::SYSTEM_EVENT_IDENTITY {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: '{meerkat_id}' is reserved"),
data: None,
}),
}
} else if let Err(message) = raw_target_validation.as_ref() {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
}
} else {
let compatibility_reservation = if runtime.identity_runtime().is_none() {
crate::member_comms_id::reserve_raw_member_target(
target_identity_runtime,
meerkat_id,
)
.await
.map(Some)
} else {
Ok(None)
};
match compatibility_reservation {
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
Ok(compatibility_reservation) => {
let member_id = compatibility_reservation
.as_ref()
.map(|reservation| reservation.alias().to_string())
.or_else(|| raw_target_validation.as_ref().ok().cloned())
.unwrap_or_else(|| meerkat_id.trim().to_string());
let spec = meerkat_mob::SpawnMemberSpec::from_wire(
profile.to_string(),
member_id.clone(),
request
.params
.get("initial_message")
.and_then(Value::as_str)
.map(|s| meerkat_core::ContentInput::from(s.to_string())),
None,
None,
);
let spawn_result = Box::pin(runtime.spawn(spec)).await;
drop(compatibility_reservation);
match spawn_result {
Ok(_member_ref) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"accepted": true,
"meerkat_id": member_id
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err}"),
data: None,
}),
},
}
}
}
}
} else {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: module_id or (profile + meerkat_id) required"
.to_string(),
data: None,
}),
}
}
}
"mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
Ok((schedules, tick_ms)) => {
match runtime.evaluate_schedule_tick(&schedules, tick_ms).await {
Ok(evaluation) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"Invalid params: {}",
format_schedule_validation_error(err)
),
data: None,
}),
},
}
}
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
},
"mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
Ok((schedules, tick_ms)) => {
match runtime.dispatch_schedule_tick(&schedules, tick_ms).await {
Ok(dispatch) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err}"),
data: None,
}),
},
}
}
Err(message) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
},
"mobkit/routing/resolve" => {
let resolve_result = match parse_routing_resolve_params(&request.params) {
Ok(resolve_request) => runtime
.resolve_routing(resolve_request)
.await
.map_err(RoutingDeliveryParamsError::Routing),
Err(e) => Err(e),
};
match resolve_result {
Ok(resolution) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
Ok(()) => {
let routes = runtime.list_runtime_routes().await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"routes": routes
})),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/routing/routes/add" => {
let add_result = match parse_routing_route_add_params(&request.params) {
Ok(route) => runtime
.add_runtime_route(route)
.await
.map_err(RoutingDeliveryParamsError::RouteMutation),
Err(e) => Err(e),
};
match add_result {
Ok(route) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "route": route })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/routing/routes/delete" => {
let delete_result = match parse_routing_route_delete_params(&request.params) {
Ok(route_key) => runtime
.delete_runtime_route(&route_key)
.await
.map_err(RoutingDeliveryParamsError::RouteMutation),
Err(e) => Err(e),
};
match delete_result {
Ok(route) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "deleted": route })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/delivery/send" => {
let send_result = match parse_delivery_send_params(&request.params) {
Ok(send_request) => runtime
.send_delivery(send_request)
.await
.map_err(RoutingDeliveryParamsError::Delivery),
Err(e) => Err(e),
};
match send_result {
Ok(record) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
Ok(history_request) => {
let history = runtime.delivery_history(history_request).await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(history).unwrap_or(Value::Null)),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/events/subscribe" => match parse_subscribe_request(&request.params) {
Ok(subscribe_request) => match runtime.subscribe_events(subscribe_request).await {
Ok(subscribe_result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {err}"),
data: None,
}),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/query_events" => {
let query: EventQuery = if request.params.is_null() {
EventQuery::default()
} else {
match serde_json::from_value(request.params.clone()) {
Ok(query) => query,
Err(err) => {
return serde_json::to_string(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: invalid query params: {err}"),
data: None,
}),
})
.unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string());
}
}
};
match runtime.event_log_store() {
Some(store) => match store.query(query).await {
Ok(events) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(events).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32603,
message: format!("query_events failed: {err}"),
data: None,
}),
},
},
None => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"status": "no_event_log_configured",
"events": [],
})),
error: None,
},
}
}
"mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
Ok(()) => {
let stores = runtime.memory_stores().await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"stores": stores,
})),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/memory/index" => match parse_memory_index_params(&request.params) {
Ok(index_request) => match runtime.memory_index(index_request).await {
Ok(indexed) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
error: None,
},
Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: MEMORY_BACKEND_UNAVAILABLE_CODE,
message: format!(
"Memory backend unavailable: {}",
MemoryParamsError::backend_message(&error)
),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"Invalid params: {}",
MemoryParamsError::Index(err).message()
),
data: None,
}),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/memory/query" => match parse_memory_query_params(&request.params) {
Ok(query_request) => {
let query_result = runtime.memory_query(query_request).await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(query_result).unwrap_or(Value::Null)),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/agent_memory/remember" => {
let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
Some(runtime) => runtime,
None => {
return maybe_error_response(
is_notification,
response_id,
-32601,
"agent memory is not configured".to_string(),
);
}
};
match parse_agent_memory_remember_params(&request.params) {
Ok(remember_request) => match runtime
.remember_agent_memory(
&remember_request.realm,
&remember_request.identity,
remember_request.memory,
)
.await
{
Ok(record) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(agent_memory_rpc_error("write", err)),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/agent_memory/forget" => {
let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
Some(runtime) => runtime,
None => {
return maybe_error_response(
is_notification,
response_id,
-32601,
"agent memory is not configured".to_string(),
);
}
};
match parse_agent_memory_forget_params(&request.params) {
Ok(forget_request) => match runtime
.forget_agent_memory(
&forget_request.realm,
&forget_request.identity,
&forget_request.memory_id,
)
.await
{
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(agent_memory_rpc_error("forget", err)),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/agent_memory/recall" => {
let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
Some(runtime) => runtime,
None => {
return maybe_error_response(
is_notification,
response_id,
-32601,
"agent memory is not configured".to_string(),
);
}
};
match parse_agent_memory_recall_params(&request.params) {
Ok(recall_request) => {
match runtime.recall_agent_memory(recall_request.request).await {
Ok(records) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "records": records })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(agent_memory_rpc_error("recall", err)),
},
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/agent_memory/update" => {
let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
Some(runtime) => runtime,
None => {
return maybe_error_response(
is_notification,
response_id,
-32601,
"agent memory is not configured".to_string(),
);
}
};
match parse_agent_memory_update_params(&request.params) {
Ok(update_request) => match runtime
.update_agent_memory(
&update_request.realm,
&update_request.identity,
&update_request.memory_id,
update_request.memory,
)
.await
{
Ok(new_id) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"memory_id": new_id,
"supersedes": update_request.memory_id,
})),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(agent_memory_rpc_error("update", err)),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/agent_memory/manifest" => {
let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
Some(runtime) => runtime,
None => {
return maybe_error_response(
is_notification,
response_id,
-32601,
"agent memory is not configured".to_string(),
);
}
};
match parse_agent_memory_manifest_params(&request.params) {
Ok(manifest_request) => match runtime
.manifest_agent_memory(
&manifest_request.realm,
&manifest_request.identity,
manifest_request.tier,
)
.await
{
Ok(records) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({ "records": records })),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(agent_memory_rpc_error("manifest", err)),
},
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/session_store/bigquery" => {
match parse_bigquery_session_store_params(&request.params)
.and_then(run_bigquery_session_store_request)
{
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
},
Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {message}"),
data: None,
}),
},
Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32011,
message: format!(
"BigQuery session store request failed: {}",
format_bigquery_store_error(&error)
),
data: None,
}),
},
}
}
"mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
Ok(gating_request) => {
let gating_result = runtime.evaluate_gating_action(gating_request).await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(gating_result).unwrap_or(Value::Null)),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
Ok(()) => {
let pending = runtime.list_gating_pending().await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"pending": pending,
})),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/gating/decide" => {
let decide_result = match parse_gating_decide_params(&request.params) {
Ok(decide_request) => runtime
.decide_gating_action(decide_request)
.await
.map_err(GatingParamsError::Decision),
Err(e) => Err(e),
};
match decide_result {
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
error: None,
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
}
}
"mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
Ok(limit) => {
let entries = runtime.gating_audit_entries(limit).await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"entries": entries,
})),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {}", err.message()),
data: None,
}),
},
},
"mobkit/call_tool" => {
let module_id = request.params.get("module_id").and_then(Value::as_str);
let tool = request.params.get("tool").and_then(Value::as_str);
let arguments = request
.params
.get("arguments")
.cloned()
.unwrap_or(serde_json::json!({}));
match (module_id, tool) {
(Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
let route = runtime
.route_module_call(
&ModuleRouteRequest {
module_id: module_id.to_string(),
method: tool.to_string(),
params: arguments,
},
timeout,
)
.await;
match route {
Ok(response) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"module_id": response.module_id,
"tool": response.method,
"result": response.payload
})),
error: None,
},
Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: format!("Module '{mid}' not loaded"),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("Tool call failed: {err:?}"),
data: None,
}),
},
}
}
_ => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "Invalid params: module_id and tool required".to_string(),
data: None,
}),
},
}
}
"mobkit/models/catalog" => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(build_models_catalog_result()),
error: None,
},
storage_methods::STORAGE_DOCTOR_METHOD => {
match storage_methods::parse_storage_doctor_params(&request.params) {
Ok(Some(params)) => {
let result =
storage_methods::run_storage_doctor(¶ms, runtime.resolved_storage())
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
Ok(None) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(storage_methods::storage_doctor_state_dir_unavailable_error()),
},
Err(reason) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!("Invalid params: {reason}"),
data: None,
}),
},
}
}
method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
handle_unified_mobpack_authoring_rpc(runtime, method, &request.params, response_id)
.await
}
"mobkit/blob/get" => {
mob_methods::handle_blob_get(runtime, response_id, &request.params).await
}
"mobkit/send_message" => {
Box::pin(mob_methods::handle_send_message(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/find_members" => {
mob_methods::handle_find_members(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/ensure_member" => {
Box::pin(mob_methods::handle_ensure_member(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/list_members" => {
mob_methods::handle_list_members(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
)
.await
}
"mobkit/get_member" => {
mob_methods::handle_get_member(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/retire_member" => {
mob_methods::handle_retire_member(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/respawn_member" => {
Box::pin(mob_methods::handle_respawn_member(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/reconcile_edges" => mob_methods::handle_reconcile_edges(runtime, response_id).await,
"mobkit/rediscover" => mob_methods::handle_rediscover(runtime, response_id).await,
"mobkit/mob_events/query" => {
mob_methods::handle_mob_events_query(runtime, response_id, request.params).await
}
"mobkit/mob_events/subscribe" => {
mob_methods::handle_mob_events_subscribe(runtime, response_id, request.params).await
}
"mobkit/cross_mob/wire" => {
Box::pin(mob_methods::handle_cross_mob_wire(
runtime,
runtime_owner.cloned(),
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/cross_mob/unwire" => {
Box::pin(mob_methods::handle_cross_mob_unwire(
runtime,
runtime_owner.cloned(),
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/cross_mob/send" => {
mob_methods::handle_cross_mob_send(
runtime,
runtime_owner.cloned(),
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/cross_mob/directory" => {
mob_methods::handle_cross_mob_directory(runtime, response_id).await
}
"mobkit/cross_mob/peer_info" => {
mob_methods::handle_cross_mob_peer_info(runtime, response_id, &request.params).await
}
"mobkit/cross_mob/wire_local" => {
mob_methods::handle_cross_mob_wire_local(
runtime,
runtime_owner.cloned(),
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/cross_mob/unwire_local" => {
mob_methods::handle_cross_mob_unwire_local(
runtime,
runtime_owner.cloned(),
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/peer_pubkey" => mob_methods::handle_peer_pubkey(runtime, response_id).await,
"mobkit/member_status" => {
mob_methods::handle_member_status(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/identity/resolved_tools" => {
mob_methods::handle_identity_resolved_tools(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/force_cancel_member" => {
mob_methods::handle_force_cancel_member(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
)
.await
}
"mobkit/spawn_helper" => {
Box::pin(mob_methods::handle_spawn_helper(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/fork_helper" => {
Box::pin(mob_methods::handle_fork_helper(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/attach_existing_session" => {
Box::pin(mob_methods::handle_attach_existing_session(
runtime,
identity_ctx.map(|ctx| &ctx.runtime),
response_id,
&request.params,
))
.await
}
"mobkit/cancel_flow" => {
mob_methods::handle_cancel_flow(runtime, response_id, &request.params).await
}
"mobkit/flow_status" => {
mob_methods::handle_flow_status(runtime, response_id, &request.params).await
}
"mobkit/list_flows" => mob_methods::handle_list_flows(runtime, response_id).await,
"mobkit/list_runs" => {
mob_methods::handle_list_runs(runtime, response_id, &request.params).await
}
"mobkit/run_flow" => {
Box::pin(mob_methods::handle_run_flow(
runtime,
response_id,
&request.params,
))
.await
}
"mobkit/collect_completed" => {
mob_methods::handle_collect_completed(runtime, response_id).await
}
"mobkit/wait_ready" => {
mob_methods::handle_wait_ready(runtime, response_id, &request.params).await
}
"mobkit/mob_labels/set" => {
mob_methods::handle_mob_labels_set(runtime, response_id, &request.params).await
}
"mobkit/mob_labels/get" => mob_methods::handle_mob_labels_get(runtime, response_id).await,
"mobkit/mob_labels/delete" => {
mob_methods::handle_mob_labels_delete(runtime, response_id).await
}
"mobkit/run_labels/set" => {
mob_methods::handle_run_labels_set(runtime, response_id, &request.params).await
}
"mobkit/run_labels/get" => {
mob_methods::handle_run_labels_get(runtime, response_id, &request.params).await
}
"mobkit/run_labels/delete" => {
mob_methods::handle_run_labels_delete(runtime, response_id, &request.params).await
}
"mobkit/send" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let content_val = request
.params
.get("content")
.cloned()
.unwrap_or(Value::Null);
let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
Ok(content) => content,
Err(err) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid content: {err}"),
);
}
};
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let send_result = identity_rt
.send_admission_tracked(
&identity,
expected_alias,
&content,
meerkat_core::types::HandlingMode::Queue,
None,
)
.await;
match send_result {
Ok(admission) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"fencing_token": admission.fencing_token.get(),
"completion_baseline": completion_cursor_json(admission.completion_baseline),
})),
error: None,
},
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/interact" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let content_val = request
.params
.get("content")
.cloned()
.unwrap_or(Value::Null);
let content =
match serde_json::from_value::<meerkat_core::ContentInput>(content_val.clone()) {
Ok(content) => content,
Err(err) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid content: {err}"),
);
}
};
let origin = request
.params
.get("origin")
.and_then(|v| v.as_str())
.unwrap_or("console");
let interaction_id = request
.params
.get("interaction_id")
.and_then(|v| v.as_str())
.map(ToString::to_string)
.unwrap_or_else(|| meerkat_core::types::SessionId::new().to_string());
let runtime_member_id = identity_rt
.status(&identity)
.await
.ok()
.and_then(|status| status.agent_runtime_id.map(|id| id.as_str().to_string()));
if let Err(err) = runtime
.reserve_identity_interaction(
identity.as_str(),
runtime_member_id.as_deref(),
&interaction_id,
origin,
content_val,
)
.await
{
return maybe_error_response(
is_notification,
response_id,
-32003,
format!("failed to reserve interaction: {err}"),
);
}
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let send_result = identity_rt
.send_admission_tracked(
&identity,
expected_alias,
&content,
meerkat_core::types::HandlingMode::Queue,
None,
)
.await;
match send_result {
Ok(admission) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"interaction_id": interaction_id,
"fencing_token": admission.fencing_token.get(),
"completion_baseline": completion_cursor_json(admission.completion_baseline),
"stream": {
"route": format!("/console/identity/{}/stream", identity.as_str()),
"identity": identity.as_str(),
}
})),
error: None,
},
Err(e) => {
runtime
.record_console_lifecycle(
identity.as_str(),
"interaction_failed",
serde_json::json!({
"interaction_id": interaction_id,
"origin": origin,
"error": e.to_string(),
}),
)
.await;
identity_error_response(response_id, &e)
}
}
}
"mobkit/dispatch" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let di_val = request
.params
.get("dispatch_input")
.cloned()
.unwrap_or(Value::Null);
let content_val = di_val
.get("content")
.cloned()
.unwrap_or_else(|| Value::String(String::new()));
let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
Ok(content) => content,
Err(err) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid dispatch_input.content: {err}"),
);
}
};
let origin_str = di_val
.get("origin")
.and_then(|v| v.as_str())
.unwrap_or("system");
let origin = match origin_str {
"connector" => crate::identity_first::DispatchOrigin::Connector,
"scheduler" => crate::identity_first::DispatchOrigin::Scheduler,
"policy" => crate::identity_first::DispatchOrigin::Policy,
"flow" => crate::identity_first::DispatchOrigin::Flow,
_ => crate::identity_first::DispatchOrigin::System,
};
let correlation_id = di_val
.get("correlation_id")
.and_then(|v| v.as_str())
.map(crate::identity_first::CorrelationId::new);
let idempotency_key = di_val
.get("idempotency_key")
.and_then(|v| v.as_str())
.map(crate::identity_first::DispatchIdempotencyKey::new);
let dispatch_input = crate::identity_first::DispatchInput {
content,
origin,
correlation_id,
idempotency_key,
};
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let dispatch_result = identity_rt
.dispatch_admission_tracked(&identity, expected_alias, &dispatch_input)
.await;
match dispatch_result {
Ok(admission) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"fencing_token": admission.fencing_token.get(),
"durable": admission.durable,
"completion_baseline": completion_cursor_json(admission.completion_baseline),
})),
error: None,
},
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/subscribe" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
match identity_rt.subscribe(&identity).await {
Ok(_receiver) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"identity": identity.as_str(),
"stream_id": identity.as_str(),
"subscribed": true,
})),
error: None,
},
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/status_identity_bootstrap" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(
serde_json::to_value(identity_rt.identity_bootstrap_status())
.unwrap_or(Value::Null),
),
error: None,
}
}
"mobkit/wait_identity_bootstrap" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let params = match request.params.as_object() {
Some(params) => params,
None => {
return maybe_error_response(
is_notification,
response_id,
-32602,
"params must be an object".to_string(),
);
}
};
if let Some(field) = params
.keys()
.find(|field| !matches!(field.as_str(), "target" | "timeout_ms"))
{
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("unsupported parameter: {field}"),
);
}
let target = match params.get("target") {
None => "materialized".to_string(),
Some(Value::String(target)) => target.clone(),
Some(_) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
"target must be a string".to_string(),
);
}
};
if !matches!(target.as_str(), "materialized" | "startup_ready") {
return maybe_error_response(
is_notification,
response_id,
-32602,
"target must be 'materialized' or 'startup_ready'".to_string(),
);
}
let wait_timeout = match params.get("timeout_ms") {
None => timeout,
Some(value) => match value.as_u64() {
Some(value) => Duration::from_millis(value),
None => {
return maybe_error_response(
is_notification,
response_id,
-32602,
"timeout_ms must be a non-negative integer".to_string(),
);
}
},
};
let (status, timed_out, startup_ready) = if target == "startup_ready" {
let mob_handle = runtime.mob_handle();
let wait = wait_identity_startup_ready(
identity_rt,
wait_timeout,
move |member_ids, remaining| {
let mob_handle = mob_handle.clone();
async move {
match mob_handle
.wait_for_members_ready(&member_ids, Some(remaining))
.await
{
Ok(_) => IdentityMemberReadiness::Ready,
Err(error)
if crate::unified_runtime::mob_ops::is_ready_wait_timeout(
&error,
) =>
{
IdentityMemberReadiness::TimedOut
}
Err(error) => IdentityMemberReadiness::Failed(error.to_string()),
}
}
},
)
.await;
match wait {
Ok(wait) => (wait.status, wait.timed_out, Some(wait.startup_ready)),
Err(error) => {
return maybe_error_response(is_notification, response_id, -32000, error);
}
}
} else {
let (status, timed_out) = identity_rt
.wait_identity_bootstrap_terminal(wait_timeout)
.await;
(status, timed_out, None)
};
let mut result = serde_json::to_value(status).unwrap_or_else(|_| serde_json::json!({}));
if let Some(object) = result.as_object_mut() {
object.insert("timed_out".to_string(), Value::Bool(timed_out));
object.insert("target".to_string(), Value::String(target));
if let Some(startup_ready) = startup_ready {
object.insert("startup_ready".to_string(), Value::Bool(startup_ready));
}
}
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
"mobkit/status_identity" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
match identity_rt.status(&identity).await {
Ok(status) => {
let continuity_health =
serde_json::to_value(&status.continuity_health).unwrap_or(Value::Null);
let result = serde_json::json!({
"state": identity_lifecycle_state_json(status.state),
"identity": status.identity.as_str(),
"agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
"session_id": status.session_id.as_ref().map(ToString::to_string),
"profile": status.profile.as_ref().map(meerkat_mob::ProfileName::as_str),
"addressability": addressability_json(status.addressability),
"display_name": status.display_name.as_ref().map(super::identity_first::DisplayName::as_str),
"labels": status.labels,
"generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
"checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
"continuity_health": continuity_health,
"lease_healthy": status.lease.as_ref().map(|lease| lease.healthy),
"lease": status.lease.as_ref().map(|lease| serde_json::json!({
"fencing_token": lease.fencing_token.get(),
"ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
"healthy": lease.healthy,
})),
});
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if rpc_live_only_fallback_allowed(&target, identity_str)
&& let Some(live) = target.live.as_ref()
{
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(rpc_live_identity_status_json(live)),
error: None,
}
} else {
identity_error_response(response_id, &e)
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/respawn" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let respawn_result = identity_rt
.respawn_identity_in_place_tracked(&identity, expected_alias)
.await;
match respawn_result {
Ok(record) => {
let live_respawn_warning: Option<Value> = None;
let cleanup_warning: Option<Value> = None;
runtime
.record_console_lifecycle(
identity.as_str(),
"identity_respawned",
serde_json::json!({
"generation": record.generation.get(),
"checkpoint_version": record.checkpoint_version.get(),
"live_respawn_warning": live_respawn_warning.clone(),
"cleanup_warning": cleanup_warning.clone(),
}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"identity": record.identity.as_str(),
"agent_runtime_id": record.agent_runtime_id.as_str(),
"session_id": record.session_id.to_string(),
"generation": record.generation.get(),
"checkpoint_version": record.checkpoint_version.get(),
"live_respawn_warning": live_respawn_warning,
"cleanup_warning": cleanup_warning,
})),
error: None,
}
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if rpc_live_only_fallback_allowed(&target, identity_str)
&& let Some(live) = target.live.as_ref()
{
match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
Ok(result) => {
runtime
.record_console_lifecycle(
live.identity.as_str(),
"identity_respawned",
serde_json::json!({}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("respawn failed: {err}"),
data: None,
}),
},
}
} else {
identity_error_response(response_id, &e)
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/retire" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let cleanup_handle = runtime.mob_handle();
let cleanup_identity = identity.clone();
let include_current = !identity_rt.has_session_bridge();
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let retire_result = identity_rt
.retire_and_cleanup_live_members_tracked(
&identity,
expected_alias,
move |retired_alias| async move {
let stale_member_ids = stale_rpc_member_ids_for_identity_with_handle(
&cleanup_handle,
cleanup_identity.as_str(),
retired_alias
.as_ref()
.map(crate::identity_first::AgentRuntimeId::as_str),
include_current,
)
.await;
retire_rpc_member_ids_with_handle(&cleanup_handle, stale_member_ids)
.await
.err()
.map(|error| {
serde_json::json!({
"kind": "stale_member_cleanup_failed_after_identity_retire",
"message": error,
"identity": cleanup_identity.as_str(),
})
})
},
)
.await;
match retire_result {
Ok((token, cleanup_warning)) => {
runtime
.record_console_lifecycle(
identity.as_str(),
"identity_retired",
serde_json::json!({
"fencing_token": token.get(),
"cleanup_warning": cleanup_warning.clone(),
}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"fencing_token": token.get(),
"cleanup_warning": cleanup_warning,
})),
error: None,
}
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if rpc_live_only_fallback_allowed(&target, identity_str)
&& let Some(live) = target.live.as_ref()
{
match retire_rpc_live_identity(runtime, live).await {
Ok(()) => {
runtime
.record_console_lifecycle(
live.identity.as_str(),
"identity_retired",
serde_json::json!({}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(
serde_json::json!({ "identity": live.identity.as_str() }),
),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("retire failed: {err}"),
data: None,
}),
},
}
} else {
identity_error_response(response_id, &e)
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/reset" => {
let identity_reset_ctx = match identity_ctx {
Some(ctx) => ctx,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_rt = &identity_reset_ctx.runtime;
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let _registered_status = match identity_rt.status(&identity).await {
Ok(status) => {
if !identity_rt.has_session_bridge() {
let response = rpc_reset_requires_session_bridge_response(response_id);
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
status
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if rpc_live_only_fallback_allowed(&target, identity_str)
&& let Some(live) = target.live.as_ref()
{
let response =
match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
Ok(result) => {
runtime
.record_console_lifecycle(
live.identity.as_str(),
"identity_reset",
serde_json::json!({}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("reset failed: {err}"),
data: None,
}),
},
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let response = identity_error_response(response_id, &e);
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
Err(e) => {
let response = identity_error_response(response_id, &e);
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
};
identity_rt.set_reset_roster_provider_context(
Some(identity_reset_ctx.roster_provider.clone()),
identity_reset_ctx.mob_definition.clone(),
);
let reset_result = if crate::member_comms_id::is_reserved_generated_alias(identity_str)
{
identity_rt
.reset_member_alias_tracked(&identity, identity_str)
.await
} else {
identity_rt.reset_tracked(&identity).await
};
match reset_result {
Ok(record) => {
let cleanup_warning = Some(serde_json::json!({
"kind": "stale_member_cleanup_skipped_after_identity_reset",
"message": "reset published the new generation without retiring stale live mob members; identity control calls reject stale runtime ids",
"identity": identity.as_str(),
"agent_runtime_id": record.agent_runtime_id.as_str(),
}));
runtime
.record_console_lifecycle(
identity.as_str(),
"identity_reset",
serde_json::json!({
"generation": record.generation.get(),
"checkpoint_version": record.checkpoint_version.get(),
"cleanup_warning": cleanup_warning.clone(),
}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"identity": record.identity.as_str(),
"agent_runtime_id": record.agent_runtime_id.as_str(),
"session_id": record.session_id.to_string(),
"generation": record.generation.get(),
"checkpoint_version": record.checkpoint_version.get(),
"cleanup_warning": cleanup_warning,
})),
error: None,
}
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if rpc_live_only_fallback_allowed(&target, identity_str)
&& let Some(live) = target.live.as_ref()
{
match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
Ok(result) => {
runtime
.record_console_lifecycle(
live.identity.as_str(),
"identity_reset",
serde_json::json!({}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("reset failed: {err}"),
data: None,
}),
},
}
} else {
identity_error_response(response_id, &e)
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/delete_identity" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let _registered_status = match identity_rt.status(&identity).await {
Ok(status) => status,
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if target.live.is_some() {
let response = JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: format!(
"delete_identity requires durable identity: {} is live-only",
identity.as_str()
),
data: Some(serde_json::json!({
"kind": "live_only_identity_delete_unsupported",
"identity": identity.as_str(),
})),
}),
};
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
let response = identity_error_response(response_id, &e);
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
Err(e) => {
let response = identity_error_response(response_id, &e);
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
};
let cleanup_handle = runtime.mob_handle();
let cleanup_identity = identity.clone();
let include_current = !identity_rt.has_session_bridge();
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then_some(identity_str);
let delete_result = identity_rt
.delete_identity_and_cleanup_live_members_tracked(
&identity,
expected_alias,
move |deleted_alias| async move {
let stale_member_ids = stale_rpc_member_ids_for_identity_with_handle(
&cleanup_handle,
cleanup_identity.as_str(),
deleted_alias
.as_ref()
.map(crate::identity_first::AgentRuntimeId::as_str),
include_current,
)
.await;
retire_rpc_member_ids_with_handle(&cleanup_handle, stale_member_ids)
.await
.err()
.map(|error| {
serde_json::json!({
"kind": "stale_member_cleanup_failed_after_identity_delete",
"identity": cleanup_identity.as_str(),
"message": error,
})
})
},
)
.await;
match delete_result {
Ok(cleanup_warning) => {
runtime
.record_console_lifecycle(
identity.as_str(),
"identity_deleted",
serde_json::json!({
"cleanup_warning": cleanup_warning,
}),
)
.await;
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"identity": identity.as_str(),
"cleanup_warning": cleanup_warning,
})),
error: None,
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/inspect_identity" => {
let identity_rt = match identity_ctx {
Some(ctx) => &ctx.runtime,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let identity_str = request
.params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target =
match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32602,
format!("invalid identity: {e}"),
);
}
};
let identity = target.identity.clone();
let status = identity_rt.status(&identity).await;
let completion_cursor = identity_rt.completion_cursor(&identity).await;
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return if is_notification {
String::new()
} else {
serialize_response(&response)
};
}
match identity_rt.inspect(&identity).await {
Ok(inspection) => {
let status = status.ok();
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"identity": identity.as_str(),
"state": status.as_ref().map(|status| identity_lifecycle_state_json(status.state)),
"profile": status.as_ref().and_then(|status| status.profile.as_ref().map(meerkat_mob::ProfileName::as_str)),
"addressability": status.as_ref().map(|status| addressability_json(status.addressability)),
"display_name": status.as_ref().and_then(|status| status.display_name.as_ref().map(super::identity_first::DisplayName::as_str)),
"labels": status.as_ref().map(|status| status.labels.clone()).unwrap_or_default(),
"generation": status.as_ref().and_then(|status| status.generation.map(super::identity_first::ContinuityGeneration::get)),
"checkpoint_version": status.as_ref().and_then(|status| status.checkpoint_version.map(super::identity_first::CheckpointVersion::get)),
"continuity_health": status.as_ref().and_then(|status| serde_json::to_value(&status.continuity_health).ok()).unwrap_or(Value::Null),
"lease_healthy": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| lease.healthy)),
"continuity": status.as_ref().map(|status| serde_json::json!({
"generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
"checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
"session_id": status.session_id.as_ref().map(ToString::to_string),
"agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
})).unwrap_or_else(|| serde_json::json!({})),
"lease": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| serde_json::json!({
"fencing_token": lease.fencing_token.get(),
"ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
"healthy": lease.healthy,
}))),
"output_preview": inspection.output_preview,
"is_final": inspection.is_final,
"peer_reachable_count": inspection.peer_reachable_count,
"completion_cursor": completion_cursor_json(completion_cursor),
})),
error: None,
}
}
Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
if let Some(live) = target.live.as_ref() {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(rpc_live_identity_inspect_json(runtime, live).await),
error: None,
}
} else {
identity_error_response(response_id, &e)
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
"mobkit/reconcile_identity" => {
let ctx = match identity_ctx {
Some(ctx) => ctx,
None => return maybe_identity_not_configured(is_notification, response_id),
};
let reconciled = match runtime.refresh_desired_topology().await {
Ok(Some(result)) => Ok(result),
Ok(None) => {
let roster_specs = match ctx
.roster_provider
.roster(&crate::identity_first::RosterContext {
mob_definition: ctx.mob_definition.clone(),
previous_identities: Vec::new(),
})
.await
{
Ok(specs) => specs,
Err(e) => {
return maybe_error_response(
is_notification,
response_id,
-32603,
format!("roster provider failed: {e}"),
);
}
};
ctx.runtime
.restore_flow_tracked(
roster_specs,
ctx.topology_provider.clone(),
ctx.customizer.clone(),
)
.await
}
Err(error) => Err(error),
};
match reconciled {
Ok(result) => {
let outcomes: serde_json::Map<String, Value> = result
.outcomes
.iter()
.map(|(id, outcome)| {
let val = match outcome {
crate::identity_first::RestoreOutcome::Created {
record, ..
} => {
serde_json::json!({
"outcome": "created",
"identity": record.identity.as_str(),
"agent_runtime_id": record.agent_runtime_id.as_str(),
"session_id": record.session_id.to_string(),
"generation": record.generation.get(),
})
}
crate::identity_first::RestoreOutcome::Dormant {
record, ..
} => {
serde_json::json!({
"outcome": "dormant",
"identity": id.as_str(),
"agent_runtime_id": record.as_ref().map(|record| record.agent_runtime_id.as_str()),
"session_id": record.as_ref().map(|record| record.session_id.to_string()),
"generation": record.as_ref().map(|record| record.generation.get()),
})
}
crate::identity_first::RestoreOutcome::Resumed {
record, ..
} => {
serde_json::json!({
"outcome": "resumed",
"identity": record.identity.as_str(),
"agent_runtime_id": record.agent_runtime_id.as_str(),
"session_id": record.session_id.to_string(),
"generation": record.generation.get(),
})
}
crate::identity_first::RestoreOutcome::Broken(failure) => {
serde_json::json!({
"outcome": "broken",
"identity": failure.identity.as_str(),
"detail": failure.detail,
})
}
};
(id.to_string(), val)
})
.collect();
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"outcomes": outcomes,
"managed_edges": result.managed_edges.len(),
})),
error: None,
}
}
Err(e) => identity_error_response(response_id, &e),
}
}
method if method.contains('/') && !method.starts_with("mobkit/") => {
let module_id = method
.split('/')
.next()
.map(ToString::to_string)
.unwrap_or_default();
let route = runtime
.route_module_call(
&ModuleRouteRequest {
module_id: module_id.clone(),
method: method.to_string(),
params: request.params,
},
timeout,
)
.await;
match route {
Ok(response) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(serde_json::json!({
"module_id": response.module_id,
"method": response.method,
"payload": response.payload
})),
error: None,
},
Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: format!("Module '{module_id}' not loaded"),
data: None,
}),
},
Err(err) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("Module route failed: {err:?}"),
data: None,
}),
},
}
}
method if method.starts_with("mobkit/live/") => match live {
None => crate::live_wiring::live_unavailable_response(response_id),
Some(live) => {
let params = request.params.clone();
let member_alias = live_member_alias(¶ms);
let identity_runtime = identity_ctx
.map(|context| &context.runtime)
.or_else(|| runtime.identity_runtime());
let authority_target = if let Some(alias) = member_alias.as_deref()
&& let Some(identity_runtime) = identity_runtime
{
identity_runtime.member_alias_lifecycle_target(alias).await
} else {
Ok(None)
};
match authority_target {
Err(error) => identity_error_response(response_id, &error),
Ok(None)
if member_alias
.as_deref()
.is_some_and(crate::member_comms_id::is_reserved_generated_alias) =>
{
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!(
"generated live target requires current identity authority: {}",
member_alias.as_deref().unwrap_or_default()
),
data: None,
}),
}
}
Ok(Some(target)) => {
let handle = runtime.mob_handle();
let identity_runtime = identity_runtime.cloned();
let live = Arc::clone(live);
let method = method.to_string();
let operation_response_id = response_id.clone();
match crate::identity_first::IdentityRuntime::run_member_alias_targets_operation_tracked(
vec![target],
move || async move {
let session = resolve_live_target(
&handle,
identity_runtime.as_ref(),
true,
¶ms,
)
.await?;
Ok(live(session, method, params, operation_response_id).await)
},
)
.await
{
Ok(response) => response,
Err(error) => identity_error_response(response_id, &error),
}
}
Ok(None) => {
match resolve_live_target(&runtime.mob_handle(), None, false, ¶ms).await
{
Ok(session) => {
live(session, method.to_string(), params, response_id).await
}
Err(error) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!("live target resolution failed: {error}"),
data: None,
}),
},
}
}
}
}
},
method if workgraph_methods::is_workgraph_method(method) => {
let service = runtime.workgraph_service();
let admission = runtime.workgraph_admission();
match workgraph_methods::handle_workgraph_method(
service.as_ref(),
&admission,
workgraph_methods::WorkgraphSurface::HostStdin,
method,
&request.params,
)
.await
{
Ok(result) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
},
Err(error) => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(error),
},
}
}
_ => JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32601,
message: "Method not found".to_string(),
data: None,
}),
},
};
if is_notification {
String::new()
} else {
serialize_response(&response)
}
}
fn build_models_catalog_result() -> Value {
let entries: Vec<Value> = meerkat_models::catalog()
.iter()
.filter_map(|e| {
let mut val = serde_json::to_value(e).ok()?;
if let Some(provider) = meerkat_core::Provider::parse_strict(e.provider)
&& let Some(profile) = meerkat_models::profile_for(provider, e.id)
&& let Ok(p) = serde_json::to_value(&profile)
{
val["profile"] = p;
}
Some(val)
})
.collect();
let defaults: Vec<Value> = meerkat_models::provider_defaults()
.iter()
.filter_map(|d| serde_json::to_value(d).ok())
.collect();
serde_json::json!({
"models": entries,
"provider_defaults": defaults,
})
}
#[derive(Debug, Clone)]
struct RpcLiveIdentityAlias {
identity: crate::identity_first::AgentIdentity,
runtime_member_id: String,
member: meerkat_mob::runtime::MobMemberListEntry,
session_id: Option<String>,
}
#[derive(Debug, Clone)]
struct RpcIdentityControlTarget {
identity: crate::identity_first::AgentIdentity,
live: Option<RpcLiveIdentityAlias>,
was_registered: bool,
}
fn rpc_live_only_fallback_allowed(
target: &RpcIdentityControlTarget,
requested_identity: &str,
) -> bool {
!target.was_registered
&& !crate::member_comms_id::is_reserved_generated_alias(requested_identity)
&& target.live.as_ref().is_some_and(|live| {
!crate::member_comms_id::is_reserved_generated_alias(&live.runtime_member_id)
})
}
fn rpc_member_durable_identity(member: &meerkat_mob::runtime::MobMemberListEntry) -> String {
crate::member_comms_id::durable_identity_label(&member.labels)
.map(str::to_owned)
.unwrap_or_else(|| {
crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
})
}
async fn resolve_rpc_live_identity_alias(
runtime: &UnifiedRuntime,
requested_identity: &str,
) -> Result<Option<RpcLiveIdentityAlias>, String> {
let matches = resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
if matches.len() > 1 {
return Err(format!(
"ambiguous live identity alias {requested_identity}: candidates [{}]",
matches
.iter()
.map(|entry| entry.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
}
Ok(matches.into_iter().next())
}
async fn resolve_rpc_live_runtime_member_alias(
runtime: &UnifiedRuntime,
runtime_member_id: &str,
) -> Result<Option<RpcLiveIdentityAlias>, String> {
let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
let handle = runtime.mob_handle();
let Some(member) = handle
.list_members_including_retiring()
.await
.into_iter()
.find(|entry| entry.agent_identity == requested_member_id)
else {
return Ok(None);
};
if !rpc_live_identity_alias_member_visible(&member) {
return Ok(None);
}
let durable_identity = rpc_member_durable_identity(&member);
let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
.map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
let session_id = handle
.resolve_bridge_session_id_observation(&member.agent_identity)
.await
.map(|session_id| session_id.to_string());
Ok(Some(RpcLiveIdentityAlias {
identity,
runtime_member_id: crate::member_comms_id::runtime_alias_str(
member.agent_identity.as_str(),
)
.into_owned(),
member,
session_id,
}))
}
async fn rpc_runtime_member_alias_exists_hidden(
runtime: &UnifiedRuntime,
runtime_member_id: &str,
) -> bool {
let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
runtime
.mob_handle()
.list_members_including_retiring()
.await
.into_iter()
.find(|entry| entry.agent_identity == requested_member_id)
.is_some_and(|member| !rpc_live_identity_alias_member_visible(&member))
}
async fn rpc_live_identity_alias_exists_hidden(
runtime: &UnifiedRuntime,
requested_identity: &str,
) -> bool {
let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
runtime
.mob_handle()
.list_members_including_retiring()
.await
.into_iter()
.any(|member| {
(member.agent_identity == requested_member_id
|| crate::member_comms_id::durable_identity_label(&member.labels)
.is_some_and(|identity| identity == requested_identity))
&& !rpc_live_identity_alias_member_visible(&member)
})
}
async fn resolve_rpc_live_identity_alias_candidates(
runtime: &UnifiedRuntime,
requested_identity: &str,
) -> Result<Vec<RpcLiveIdentityAlias>, String> {
let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
let handle = runtime.mob_handle();
let members = handle.list_members_including_retiring().await;
let exact_matches = members
.iter()
.filter(|entry| entry.agent_identity == requested_member_id)
.cloned()
.collect::<Vec<_>>();
let label_matches = members
.iter()
.filter(|entry| {
crate::member_comms_id::durable_identity_label(&entry.labels)
.is_some_and(|identity| identity == requested_identity)
})
.cloned()
.collect::<Vec<_>>();
let mut matches = exact_matches;
matches.extend(label_matches);
let mut seen_member_ids = BTreeSet::new();
matches.retain(|entry| seen_member_ids.insert(entry.agent_identity.to_string()));
let mut aliases = Vec::with_capacity(matches.len());
for member in matches {
if !rpc_live_identity_alias_member_visible(&member) {
continue;
}
let durable_identity = rpc_member_durable_identity(&member);
let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
.map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
let session_id = handle
.resolve_bridge_session_id_observation(&member.agent_identity)
.await
.map(|session_id| session_id.to_string());
aliases.push(RpcLiveIdentityAlias {
identity,
runtime_member_id: crate::member_comms_id::runtime_alias_str(
member.agent_identity.as_str(),
)
.into_owned(),
member,
session_id,
});
}
Ok(aliases)
}
fn rpc_live_identity_alias_member_visible(
member: &meerkat_mob::runtime::MobMemberListEntry,
) -> bool {
rpc_live_identity_alias_visible(member.role.as_str(), &member.labels)
}
fn rpc_live_identity_alias_visible(
member_role: &str,
labels: &std::collections::BTreeMap<String, String>,
) -> bool {
let projected_role = labels
.get("role")
.map(String::as_str)
.unwrap_or(member_role);
!is_implicit_delegate_member(member_role, labels)
&& !is_implicit_delegate_member(projected_role, labels)
}
async fn resolve_rpc_identity_control_target(
runtime: &UnifiedRuntime,
identity_rt: &crate::identity_first::IdentityRuntime,
requested_identity: &str,
) -> Result<RpcIdentityControlTarget, String> {
let requested_identity = crate::member_comms_id::runtime_alias_str(requested_identity);
let requested_identity = requested_identity.as_ref();
if crate::member_comms_id::is_reserved_generated_alias(requested_identity) {
for status in identity_rt.statuses().await {
if status
.agent_runtime_id
.as_ref()
.is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
{
let identity = status.identity;
let registered_live =
resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
if let Some(registered) = registered_live {
return Ok(RpcIdentityControlTarget {
identity,
live: Some(registered),
was_registered: true,
});
}
if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
return Err(format!("identity hidden by policy: {requested_identity}"));
}
let durable_live_candidates =
resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
let durable_live = if durable_live_candidates.len() > 1 {
return Err(format!(
"ambiguous live identity alias {}: candidates [{}]",
identity.as_str(),
durable_live_candidates
.iter()
.map(|alias| alias.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
} else {
durable_live_candidates.into_iter().next()
};
return Ok(RpcIdentityControlTarget {
identity,
live: durable_live,
was_registered: true,
});
}
}
let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
if let Some(live_alias) = live {
if let Ok(status) = identity_rt.status(&live_alias.identity).await
&& !rpc_live_alias_matches_status_runtime(Some(&live_alias), &status)
{
return Ok(RpcIdentityControlTarget {
identity: live_alias.identity.clone(),
live: Some(live_alias),
was_registered: true,
});
}
let live_identity_candidates =
resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
.await?;
if live_identity_candidates.len() > 1 {
return Err(format!(
"ambiguous live identity alias {}: candidates [{}]",
live_alias.identity.as_str(),
live_identity_candidates
.iter()
.map(|alias| alias.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
}
return Ok(RpcIdentityControlTarget {
identity: live_alias.identity.clone(),
live: Some(live_alias),
was_registered: false,
});
}
if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
return Err(format!("identity hidden by policy: {requested_identity}"));
}
return Err(format!("runtime identity not found: {requested_identity}"));
}
if let Ok(identity) = crate::identity_first::AgentIdentity::parse(requested_identity) {
match identity_rt.status(&identity).await {
Ok(status) => {
let registered_live = match status.agent_runtime_id.as_ref() {
Some(runtime_id) => {
resolve_rpc_live_runtime_member_alias(runtime, runtime_id.as_str()).await?
}
None => None,
};
if let Some(registered) = registered_live {
return Ok(RpcIdentityControlTarget {
identity,
live: Some(registered),
was_registered: true,
});
}
if let Some(runtime_id) = status.agent_runtime_id.as_ref()
&& rpc_runtime_member_alias_exists_hidden(runtime, runtime_id.as_str()).await
{
return Err(format!("identity hidden by policy: {requested_identity}"));
}
let requested_live_candidates =
resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
let requested_live = if requested_live_candidates.len() > 1 {
return Err(format!(
"ambiguous live identity alias {requested_identity}: candidates [{}]",
requested_live_candidates
.iter()
.map(|alias| alias.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
} else {
requested_live_candidates.into_iter().next()
};
return Ok(RpcIdentityControlTarget {
identity,
live: requested_live,
was_registered: true,
});
}
Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {}
Err(err) => return Err(err.to_string()),
}
}
for status in identity_rt.statuses().await {
if status
.agent_runtime_id
.as_ref()
.is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
{
let identity = status.identity;
let registered_live =
resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
let durable_live_candidates =
resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
let durable_live = if durable_live_candidates.len() > 1 {
return Err(format!(
"ambiguous live identity alias {}: candidates [{}]",
identity.as_str(),
durable_live_candidates
.iter()
.map(|alias| alias.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
} else {
durable_live_candidates.into_iter().next()
};
let live = match (registered_live, durable_live) {
(Some(registered), Some(durable))
if registered.runtime_member_id == durable.runtime_member_id =>
{
Some(registered)
}
(Some(registered), None) => Some(registered),
(Some(_registered), Some(durable)) => Some(durable),
(None, durable) => durable,
};
return Ok(RpcIdentityControlTarget {
identity,
live,
was_registered: true,
});
}
}
let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
if let Some(live_alias) = live {
if let Some(bound_status) = identity_rt.statuses().await.into_iter().find(|status| {
status
.agent_runtime_id
.as_ref()
.is_some_and(|runtime_id| runtime_id.as_str() == live_alias.runtime_member_id)
}) && bound_status.identity != live_alias.identity
{
return Err(format!(
"stale live identity alias: live console alias {} resolves to {}, but identity runtime binding belongs to {}",
live_alias.identity.as_str(),
live_alias.runtime_member_id,
bound_status.identity.as_str(),
));
}
let live_identity_candidates =
resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
.await?;
if live_identity_candidates.len() > 1 {
return Err(format!(
"ambiguous live identity alias {}: candidates [{}]",
live_alias.identity.as_str(),
live_identity_candidates
.iter()
.map(|alias| alias.runtime_member_id.clone())
.collect::<Vec<_>>()
.join(", ")
));
}
return Ok(RpcIdentityControlTarget {
identity: live_alias.identity.clone(),
live: Some(live_alias),
was_registered: false,
});
}
if rpc_live_identity_alias_exists_hidden(runtime, requested_identity).await {
return Err(format!("identity hidden by policy: {requested_identity}"));
}
let identity = crate::identity_first::AgentIdentity::parse(requested_identity)
.map_err(|err| err.to_string())?;
Ok(RpcIdentityControlTarget {
identity,
live: None,
was_registered: false,
})
}
fn rpc_reset_requires_session_bridge_response(response_id: Value) -> JsonRpcResponse {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32602,
message: "reset requires an identity runtime with a session bridge".to_string(),
data: Some(serde_json::json!({
"kind": "identity_reset_requires_session_bridge",
})),
}),
}
}
fn rpc_live_alias_matches_status_runtime(
alias: Option<&RpcLiveIdentityAlias>,
status: &crate::identity_first::IdentityStatus,
) -> bool {
let Some(alias) = alias else {
return true;
};
let session_matches = match (
status.session_id.as_ref().map(ToString::to_string),
alias.session_id.as_deref(),
) {
(Some(status_session), Some(live_session)) => status_session == live_session,
_ => true,
};
status
.agent_runtime_id
.as_ref()
.is_some_and(|runtime_id| runtime_id.as_str() == alias.runtime_member_id)
&& alias.identity == status.identity
&& session_matches
}
async fn rpc_stale_live_alias_error_response(
identity_rt: &crate::identity_first::IdentityRuntime,
target: &RpcIdentityControlTarget,
response_id: Value,
) -> Option<JsonRpcResponse> {
let live = target.live.as_ref()?;
let Ok(status) = identity_rt.status(&target.identity).await else {
return None;
};
if rpc_live_alias_matches_status_runtime(Some(live), &status) {
return None;
}
Some(JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code: -32000,
message: format!(
"identity runtime binding for {} points at {}, but requested live member is {}",
target.identity.as_str(),
status
.agent_runtime_id
.as_ref()
.map(crate::identity_first::AgentRuntimeId::as_str)
.unwrap_or("<none>"),
live.runtime_member_id
),
data: Some(serde_json::json!({
"kind": "stale_identity_runtime_binding",
"identity": target.identity.as_str(),
"registered_runtime_member_id": status.agent_runtime_id.as_ref().map(crate::identity_first::AgentRuntimeId::as_str),
"live_runtime_member_id": live.runtime_member_id,
"registered_session_id": status.session_id.as_ref().map(ToString::to_string),
"live_session_id": live.session_id,
})),
}),
})
}
fn rpc_member_is_addressable(member: &meerkat_mob::runtime::MobMemberListEntry) -> bool {
member
.labels
.get("addressable")
.map(|value| !value.eq_ignore_ascii_case("false"))
.unwrap_or(true)
}
fn rpc_live_identity_status_json(alias: &RpcLiveIdentityAlias) -> Value {
serde_json::json!({
"state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
"identity": alias.identity.as_str(),
"agent_runtime_id": alias.runtime_member_id,
"session_id": alias.session_id,
"profile": alias.member.role.to_string(),
"addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
"display_name": alias.member.labels.get("display_name"),
"labels": alias.member.labels,
"generation": Value::Null,
"checkpoint_version": Value::Null,
"continuity_health": Value::Null,
"lease_healthy": Value::Null,
"lease": Value::Null,
})
}
async fn rpc_live_identity_inspect_json(
runtime: &UnifiedRuntime,
alias: &RpcLiveIdentityAlias,
) -> Value {
let snapshot = runtime
.mob_handle()
.member_status(&crate::member_comms_id::mob_member_id(
alias.runtime_member_id.as_str(),
))
.await
.ok();
serde_json::json!({
"identity": alias.identity.as_str(),
"state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
"profile": alias.member.role.to_string(),
"addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
"display_name": alias.member.labels.get("display_name"),
"labels": alias.member.labels,
"generation": Value::Null,
"checkpoint_version": Value::Null,
"continuity_health": Value::Null,
"lease_healthy": Value::Null,
"continuity": {
"generation": Value::Null,
"checkpoint_version": Value::Null,
"session_id": alias.session_id,
"agent_runtime_id": alias.runtime_member_id,
},
"lease": Value::Null,
"output_preview": snapshot.as_ref().and_then(|snapshot| snapshot.output_preview.clone()),
"is_final": snapshot.as_ref().map(|snapshot| snapshot.is_final).unwrap_or(false),
"peer_reachable_count": alias.member.wired_to.len(),
"completion_cursor": Value::Null,
"progress": snapshot.as_ref().and_then(|snapshot| snapshot.progress.clone()),
})
}
async fn retire_rpc_live_identity(
runtime: &UnifiedRuntime,
alias: &RpcLiveIdentityAlias,
) -> Result<(), String> {
retire_rpc_runtime_member_id(runtime, alias.runtime_member_id.as_str()).await
}
async fn retire_rpc_runtime_member_id(
runtime: &UnifiedRuntime,
runtime_member_id: &str,
) -> Result<(), String> {
retire_rpc_runtime_member_id_with_handle(&runtime.mob_handle(), runtime_member_id).await
}
async fn retire_rpc_runtime_member_id_with_handle(
handle: &meerkat_mob::MobHandle,
runtime_member_id: &str,
) -> Result<(), String> {
match handle
.retire(crate::member_comms_id::mob_member_id(runtime_member_id))
.await
{
Ok(()) => Ok(()),
Err(err) if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) => Ok(()),
Err(err) => Err(err.to_string()),
}
}
fn rpc_member_id_matches_durable_identity(member_id: &str, durable_identity: &str) -> bool {
crate::member_comms_id::runtime_alias_str(member_id) == durable_identity
}
fn rpc_runtime_alias_generation(alias: &str, durable_identity: &str) -> Option<u64> {
let alias = crate::member_comms_id::runtime_alias_str(alias);
let rest = alias.strip_prefix("rt:")?;
let (identity, generation) = rest.rsplit_once(':')?;
if identity != durable_identity {
return None;
}
generation.parse().ok()
}
async fn stale_rpc_member_ids_for_identity_with_handle(
handle: &meerkat_mob::MobHandle,
durable_identity: &str,
current_runtime_member_id: Option<&str>,
include_current: bool,
) -> Vec<String> {
let Some(current_generation) = current_runtime_member_id
.and_then(|alias| rpc_runtime_alias_generation(alias, durable_identity))
else {
return Vec::new();
};
handle
.list_members_including_retiring()
.await
.into_iter()
.filter(|member| {
if !rpc_live_identity_alias_member_visible(member) {
return false;
}
let matches_identity =
rpc_member_id_matches_durable_identity(
member.agent_identity.as_str(),
durable_identity,
) || crate::member_comms_id::durable_identity_label(&member.labels)
.is_some_and(|identity| identity == durable_identity);
let public_alias =
crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str());
matches_identity
&& rpc_runtime_alias_generation(public_alias.as_ref(), durable_identity)
.is_some_and(|generation| {
generation < current_generation
|| (include_current && generation == current_generation)
})
})
.map(|member| {
crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
})
.collect()
}
async fn retire_rpc_member_ids_with_handle(
handle: &meerkat_mob::MobHandle,
member_ids: Vec<String>,
) -> Result<(), String> {
for member_id in member_ids {
retire_rpc_runtime_member_id_with_handle(handle, &member_id).await?;
}
Ok(())
}
async fn respawn_rpc_live_identity(
runtime: &UnifiedRuntime,
alias: &RpcLiveIdentityAlias,
) -> Result<Value, String> {
let mut result = Box::pin(respawn_rpc_runtime_member_id(
runtime,
alias.runtime_member_id.as_str(),
))
.await?;
result["identity"] = serde_json::json!(alias.identity.as_str());
Ok(result)
}
async fn respawn_rpc_runtime_member_id(
runtime: &UnifiedRuntime,
runtime_member_id: &str,
) -> Result<Value, String> {
respawn_rpc_runtime_member_id_with_handle(&runtime.mob_handle(), runtime_member_id).await
}
async fn respawn_rpc_runtime_member_id_with_handle(
handle: &meerkat_mob::MobHandle,
runtime_member_id: &str,
) -> Result<Value, String> {
let member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
let entry_before_respawn = handle.get_member(&member_id).await.ok().flatten();
let mut topology_restore_warning = None;
match handle.respawn(member_id.clone(), None).await {
Ok(_receipt) => {}
Err(err) => {
if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(&err) {
tracing::warn!(
member_id = %member_id,
failed_peer_count = failed_peer_ids.len(),
failed_peer_ids = ?failed_peer_ids,
"rpc member respawn restored member with isolated peer edges; continuing degraded respawn"
);
topology_restore_warning = Some(topology_restore_warning_json(&failed_peer_ids));
} else if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) {
if handle
.get_member(&member_id)
.await
.map_err(|lookup_err| lookup_err.to_string())?
.is_none()
&& let Some(entry) = entry_before_respawn
{
let mut spec =
meerkat_mob::SpawnMemberSpec::new(entry.role.clone(), member_id.clone());
if !entry.labels.is_empty() {
spec = spec.with_labels(entry.labels.clone());
}
handle
.ensure_member(spec)
.await
.map_err(|ensure_err| ensure_err.to_string())?;
}
} else {
return Err(err.to_string());
}
}
}
let session_id = handle
.resolve_bridge_session_id_observation(&member_id)
.await
.map(|session_id| session_id.to_string());
Ok(serde_json::json!({
"agent_runtime_id": runtime_member_id,
"session_id": session_id,
"generation": Value::Null,
"checkpoint_version": Value::Null,
"topology_restore_warning": topology_restore_warning,
}))
}
fn identity_not_configured(response_id: Value) -> String {
error_response(response_id, -32601, "identity-first runtime not configured")
}
fn maybe_identity_not_configured(is_notification: bool, response_id: Value) -> String {
if is_notification {
String::new()
} else {
identity_not_configured(response_id)
}
}
fn completion_cursor_json(cursor: crate::identity_first::CompletionCursor) -> Value {
serde_json::json!({
"epoch": cursor.epoch.get(),
"turns": cursor.turns,
})
}
fn addressability_json(addressability: crate::identity_first::AgentAddressability) -> &'static str {
match addressability {
crate::identity_first::AgentAddressability::Addressable => "addressable",
crate::identity_first::AgentAddressability::InternalOnly => "internal_only",
}
}
fn identity_lifecycle_state_json(
state: crate::identity_first::IdentityLifecycleState,
) -> &'static str {
state.wire_str()
}
fn identity_error_response(
response_id: Value,
err: &crate::identity_first::IdentityRuntimeError,
) -> JsonRpcResponse {
use crate::identity_first::IdentityRuntimeError;
let (code, message) = match err {
IdentityRuntimeError::UnknownIdentity(id) => (-32001, format!("unknown identity: {id}")),
IdentityRuntimeError::NotAddressable(na) => {
(-32002, format!("not addressable: {}", na.identity))
}
IdentityRuntimeError::NoActiveLease(id) => (-32003, format!("no active lease: {id}")),
IdentityRuntimeError::LeaseLost(id) => (-32005, format!("lease lost: {id}")),
_ => (-32603, format!("{err}")),
};
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code,
message,
data: None,
}),
}
}
fn error_response(response_id: Value, code: i64, message: impl Into<String>) -> String {
let message = message.into();
let ambiguous_alias_rest = message
.strip_prefix("ambiguous live identity alias ")
.or_else(|| message.strip_prefix("invalid identity: ambiguous live identity alias "));
let stale_live_alias_rest = message
.strip_prefix("stale live identity alias: live console alias ")
.or_else(|| {
message.strip_prefix("invalid identity: stale live identity alias: live console alias ")
});
let hidden_policy_identity = message
.strip_prefix("identity hidden by policy: ")
.or_else(|| message.strip_prefix("invalid identity: identity hidden by policy: "));
let data = if let Some(rest) = ambiguous_alias_rest {
let (identity, candidates) = rest
.split_once(": candidates [")
.map(|(identity, candidates)| {
(
identity.to_string(),
candidates
.trim_end_matches(']')
.split(',')
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string)
.collect::<Vec<_>>(),
)
})
.unwrap_or_else(|| (rest.to_string(), Vec::new()));
Some(serde_json::json!({
"kind": "ambiguous_live_identity_alias",
"identity": identity,
"candidates": candidates,
}))
} else if let Some(rest) = stale_live_alias_rest {
let (identity, rest) = rest.split_once(" resolves to ").unwrap_or((rest, ""));
let (runtime_member_id, bound_identity) = rest
.split_once(", but identity runtime binding belongs to ")
.unwrap_or((rest, ""));
Some(serde_json::json!({
"kind": "stale_live_identity_alias",
"identity": identity,
"live_runtime_member_id": runtime_member_id,
"bound_identity": bound_identity,
}))
} else {
hidden_policy_identity.map(|identity| {
serde_json::json!({
"kind": "identity_hidden_by_policy",
"identity": identity,
})
})
};
serialize_response(&JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code,
message,
data,
}),
})
}
fn live_member_alias(params: &Value) -> Option<String> {
if params.get("session_id").and_then(Value::as_str).is_some() {
return None;
}
params
.get("identity")
.and_then(Value::as_str)
.or_else(|| params.get("member_id").and_then(Value::as_str))
.map(|raw| crate::member_comms_id::runtime_alias_str(raw).into_owned())
}
async fn resolve_live_target(
handle: &meerkat_mob::MobHandle,
identity_runtime: Option<&Arc<crate::identity_first::IdentityRuntime>>,
identity_authoritative: bool,
params: &Value,
) -> Result<Option<meerkat_core::types::SessionId>, String> {
if let Some(raw) = params.get("session_id").and_then(Value::as_str) {
return Ok(meerkat_core::types::SessionId::parse(raw).ok());
}
let Some(raw) = live_member_alias(params) else {
return Ok(None);
};
let current_alias = if identity_authoritative {
let identity_runtime = identity_runtime
.ok_or_else(|| "durable live target lost its IdentityRuntime authority".to_string())?;
let identity =
crate::identity_first::IdentityRuntime::identity_for_generated_member_alias(&raw)
.or_else(|| crate::identity_first::AgentIdentity::parse(&raw).ok())
.ok_or_else(|| format!("invalid durable live target {raw:?}"))?;
identity_runtime
.status(&identity)
.await
.map_err(|error| error.to_string())?
.agent_runtime_id
.map(|runtime_id| runtime_id.to_string())
.ok_or_else(|| format!("identity {identity} has no current runtime member"))?
} else {
let direct = crate::member_comms_id::mob_member_id(&raw);
if handle
.get_member(&direct)
.await
.map_err(|error| error.to_string())?
.is_some()
{
raw.clone()
} else {
let candidates = handle
.list_members_including_retiring()
.await
.into_iter()
.filter(|entry| {
crate::member_comms_id::durable_identity_label(&entry.labels)
.is_some_and(|identity| identity == raw)
})
.map(|entry| {
crate::member_comms_id::runtime_alias_str(entry.agent_identity.as_str())
.into_owned()
})
.collect::<BTreeSet<_>>();
match candidates.len() {
0 => raw.clone(),
1 => candidates
.into_iter()
.next()
.ok_or_else(|| "live member alias candidate disappeared".to_string())?,
_ => {
return Err(format!(
"ambiguous live member alias {raw}: candidates [{}]",
candidates.into_iter().collect::<Vec<_>>().join(", ")
));
}
}
}
};
let member_id = crate::member_comms_id::mob_member_id(¤t_alias);
Ok(handle.resolve_bridge_session_id(&member_id).await)
}
fn maybe_error_response(
is_notification: bool,
response_id: Value,
code: i64,
message: impl Into<String>,
) -> String {
if is_notification {
String::new()
} else {
error_response(response_id, code, message)
}
}
pub(crate) fn agent_memory_rpc_error(
operation: &str,
err: crate::identity_first::AgentMemoryError,
) -> JsonRpcError {
let code = match &err {
crate::identity_first::AgentMemoryError::InvalidConfig(_)
| crate::identity_first::AgentMemoryError::InvalidRecord(_) => -32602,
crate::identity_first::AgentMemoryError::Unsupported(_) => -32601,
crate::identity_first::AgentMemoryError::Io(_)
| crate::identity_first::AgentMemoryError::Parse(_)
| crate::identity_first::AgentMemoryError::Timeout(_) => -32603,
};
JsonRpcError {
code,
message: format!("agent memory {operation} failed: {err}"),
data: None,
}
}
fn serialize_response(response: &JsonRpcResponse) -> String {
serde_json::to_string(response).unwrap_or_else(|_| {
r#"{"jsonrpc":"2.0","id":null,"error":{"code":-32603,"message":"Internal error"}}"#
.to_string()
})
}
#[cfg(test)]
#[allow(clippy::expect_used)]
mod tests {
use super::{
IdentityMemberReadiness, error_response, handle_unified_rpc_json, identity_error_response,
resolve_live_target, resolve_rpc_identity_control_target, rpc_live_identity_alias_visible,
rpc_member_id_matches_durable_identity, rpc_runtime_alias_generation,
wait_identity_startup_ready,
};
use crate::identity_first::contracts::{ContinuityStore, LeaseProvider, RosterProvider};
use crate::identity_first::{
AgentAddressability, AgentBuildDraft, AgentIdentity, AgentRuntimeId, BridgeError,
CheckpointVersion, ContinuityGeneration, ContinuityRecord, DurabilityPolicy,
DurableAgentSpec, FencingToken, IdentityLifecycleState, IdentityRuntime,
IdentityRuntimeConfig, LeaseAcquireResult, LeaseGrant, LocalContinuityStore,
LocalLeaseProvider, MarkdownAgentMemoryStore, ResumeSessionOutcome, RosterContext,
RosterError, SessionBridge, SessionSnapshot,
};
use crate::{
DiscoverySpec, IdentityFirstContext, MobBootstrapOptions, MobBootstrapSpec, MobKitConfig,
UnifiedRuntime,
};
use async_trait::async_trait;
use meerkat::{AgentFactory, Config, build_ephemeral_service};
use meerkat_client::TestClient;
use meerkat_mob::{MobDefinition, MobStorage, SpawnMemberSpec};
use serde_json::{Value, json};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
#[derive(Debug, Default)]
struct EmptyRosterProvider;
#[async_trait]
impl RosterProvider for EmptyRosterProvider {
async fn roster(
&self,
_context: &RosterContext,
) -> Result<Vec<DurableAgentSpec>, RosterError> {
Ok(Vec::new())
}
}
#[derive(Debug, Default)]
struct ContextRequiredEmptyRosterProvider {
missing_definition_calls: std::sync::atomic::AtomicUsize,
}
impl ContextRequiredEmptyRosterProvider {
fn missing_definition_calls(&self) -> usize {
self.missing_definition_calls
.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[async_trait]
impl RosterProvider for ContextRequiredEmptyRosterProvider {
async fn roster(
&self,
context: &RosterContext,
) -> Result<Vec<DurableAgentSpec>, RosterError> {
if context.mob_definition.is_none() {
self.missing_definition_calls
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return Err(RosterError::ProviderUnavailable(
"mob definition required".to_string(),
));
}
Ok(Vec::new())
}
}
#[derive(Debug)]
struct ContextRequiredStaticRosterProvider {
specs: Vec<DurableAgentSpec>,
missing_definition_calls: std::sync::atomic::AtomicUsize,
}
impl ContextRequiredStaticRosterProvider {
fn new(specs: Vec<DurableAgentSpec>) -> Self {
Self {
specs,
missing_definition_calls: std::sync::atomic::AtomicUsize::new(0),
}
}
fn missing_definition_calls(&self) -> usize {
self.missing_definition_calls
.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[async_trait]
impl RosterProvider for ContextRequiredStaticRosterProvider {
async fn roster(
&self,
context: &RosterContext,
) -> Result<Vec<DurableAgentSpec>, RosterError> {
if context.mob_definition.is_none() {
self.missing_definition_calls
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return Err(RosterError::ProviderUnavailable(
"mob definition required".to_string(),
));
}
Ok(self.specs.clone())
}
}
#[derive(Debug, Default)]
struct RpcResetTestBridge {
create_calls: std::sync::atomic::AtomicUsize,
last_create_spec: tokio::sync::Mutex<Option<DurableAgentSpec>>,
}
impl RpcResetTestBridge {
async fn last_create_spec(&self) -> Option<DurableAgentSpec> {
self.last_create_spec.lock().await.clone()
}
}
#[async_trait]
impl SessionBridge for RpcResetTestBridge {
async fn create_session(
&self,
_identity: &AgentIdentity,
_runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
_draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
self.create_calls
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
*self.last_create_spec.lock().await = Some(spec.clone());
Ok(session_id.clone())
}
async fn resume_session(
&self,
_identity: &AgentIdentity,
_runtime_id: &AgentRuntimeId,
_spec: &DurableAgentSpec,
_draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
_snapshot: &SessionSnapshot,
) -> Result<ResumeSessionOutcome, BridgeError> {
Ok(ResumeSessionOutcome::Resumed {
session_id: session_id.clone(),
})
}
async fn deliver(
&self,
_runtime_id: &AgentRuntimeId,
_content: &meerkat_core::ContentInput,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
Ok(meerkat_core::types::SessionId::new())
}
async fn checkpoint_session(
&self,
_runtime_id: &AgentRuntimeId,
_session_id: &meerkat_core::types::SessionId,
) -> Result<SessionSnapshot, BridgeError> {
Ok(SessionSnapshot { data: Vec::new() })
}
async fn retire_member(&self, _runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
Ok(())
}
}
struct ReadOnlyAgentMemoryProvider;
#[async_trait]
impl crate::identity_first::AgentMemoryProvider for ReadOnlyAgentMemoryProvider {
async fn recall(
&self,
_request: crate::identity_first::AgentMemoryRecallRequest,
) -> Result<
Vec<crate::identity_first::AgentMemoryRecord>,
crate::identity_first::AgentMemoryError,
> {
Ok(Vec::new())
}
}
fn rpc_test_mob_spec(
temp_dir: &tempfile::TempDir,
) -> Result<MobBootstrapSpec, Box<dyn std::error::Error + Send + Sync>> {
let session_path = temp_dir.path().join("sessions");
std::fs::create_dir_all(&session_path)?;
let factory = AgentFactory::new(&session_path).comms(true);
let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
let definition = MobDefinition::from_toml(
r#"
[mob]
id = "rpc-identity-alias-test"
[profiles.worker]
model = "gpt-5.5"
external_addressable = true
[profiles.worker.tools]
comms = true
"#,
)?;
Ok(
MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: Some(Arc::new(TestClient::default())),
}),
)
}
async fn spawn_identity_projection_fixture(
runtime: &UnifiedRuntime,
mut spec: SpawnMemberSpec,
) -> Result<(), meerkat_mob::MobError> {
let runtime_alias = spec.identity.to_string();
spec.identity = crate::member_comms_id::mob_member_id(&runtime_alias);
Box::pin(runtime.mob_handle().spawn_spec(spec)).await?;
Ok(())
}
fn rpc_reprofile_mob_spec(
temp_dir: &tempfile::TempDir,
) -> Result<MobBootstrapSpec, Box<dyn std::error::Error + Send + Sync>> {
let session_path = temp_dir.path().join("sessions");
std::fs::create_dir_all(&session_path)?;
let factory = AgentFactory::new(&session_path).comms(true);
let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
let definition = MobDefinition::from_toml(
r#"
[mob]
id = "rpc-reset-reprofile-test"
[profiles.domain]
model = "gpt-5.5"
external_addressable = true
[profiles.domain.tools]
comms = true
[profiles.security]
model = "gpt-5.5"
external_addressable = true
[profiles.security.tools]
comms = true
shell = true
"#,
)?;
Ok(
MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: Some(Arc::new(TestClient::default())),
}),
)
}
fn rpc_durable_spec(identity: &str, profile: &str) -> DurableAgentSpec {
DurableAgentSpec {
identity: AgentIdentity::parse(identity).expect("valid identity"),
profile: meerkat_mob::ProfileName::from(profile),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
}
}
#[tokio::test]
async fn unified_rpc_spawn_rejects_reserved_system_identity()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-reserved-identity-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/spawn_member",
"params": {
"profile": "worker",
"meerkat_id": crate::console_contracts::SYSTEM_EVENT_IDENTITY,
},
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#?}");
assert!(
response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("reserved")),
"rejection must name the reservation: {response:#?}"
);
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn unified_rpc_spawn_member_rejects_generated_runtime_alias_namespace()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-generated-alias-reservation-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/spawn_member",
"params": {
"profile": "worker",
"meerkat_id": "rt:user:forged:0",
},
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#?}");
assert!(
response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("reserved")),
"generated aliases must be rejected at admission: {response:#?}"
);
assert!(
runtime.mob_handle().list_members().await.is_empty(),
"reserved namespace rejection must happen before raw member spawn"
);
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn unified_rpc_spawn_holds_explicit_identity_context_alias_reservation()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-explicit-identity-context-reservation-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
assert!(runtime.identity_runtime().is_none());
let identity_runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-explicit-identity-context-reservation-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity_ctx = IdentityFirstContext {
runtime: identity_runtime.clone(),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: Some(runtime.mob_handle().definition().clone()),
};
let held_lock = identity_runtime
.raw_member_alias_lock("compat-worker")
.await
.lock_owned()
.await;
let request = json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/spawn_member",
"params": {
"profile": "worker",
"meerkat_id": " compat-worker ",
},
})
.to_string();
assert!(
tokio::time::timeout(
Duration::from_millis(25),
handle_unified_rpc_json(
&runtime,
&request,
Duration::from_secs(1),
None,
Some(&identity_ctx),
),
)
.await
.is_err(),
"explicit identity authority must block the raw spawn on its canonical alias lock"
);
drop(held_lock);
let response: Value = serde_json::from_str(
&tokio::time::timeout(
Duration::from_secs(2),
handle_unified_rpc_json(
&runtime,
&request,
Duration::from_secs(1),
None,
Some(&identity_ctx),
),
)
.await?,
)?;
assert!(response["error"].is_null(), "{response:#?}");
assert_eq!(response["result"]["meerkat_id"], json!("compat-worker"));
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn unified_rpc_rejects_encoded_roster_ingress_and_authoritative_identity_labels()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-public-alias-ingress-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let encoded = crate::member_comms_id::mob_member_id_str("rt:secret:0");
let encoded_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/get_member",
"params": { "member_id": encoded.as_ref() },
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert_eq!(
encoded_response["error"]["code"],
json!(-32602),
"encoded roster spelling must be rejected before resolution: {encoded_response:#?}"
);
let forged_label_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/ensure_member",
"params": {
"role": "worker",
"agent_identity": "raw-worker",
"labels": { "agent_identity": "secret" },
},
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert_eq!(
forged_label_response["error"]["code"],
json!(-32602),
"raw creation must not mint runtime-authoritative identity labels: {forged_label_response:#?}"
);
assert!(runtime.mob_handle().list_members().await.is_empty());
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn unified_capabilities_separate_mobpack_authoring_from_runtime_controls()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-authoring-capabilities-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/capabilities",
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert!(response["error"].is_null(), "{response:#?}");
let methods = response["result"]["methods"]
.as_array()
.expect("methods array")
.iter()
.filter_map(Value::as_str)
.collect::<Vec<_>>();
for method in super::MOBPACK_AUTHORING_METHODS {
assert!(
methods.contains(method),
"missing authoring method {method}"
);
}
assert_eq!(
response["result"]["authoring_capabilities"]["domain"],
json!("mobpack_authoring")
);
assert_eq!(
response["result"]["authoring_capabilities"]["runtime_mutation"],
json!(false)
);
assert_eq!(
response["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
);
assert_eq!(
response["result"]["authoring_capabilities"]["methods"]
.as_array()
.expect("authoring methods")
.iter()
.filter_map(Value::as_str)
.collect::<Vec<_>>(),
super::MOBPACK_AUTHORING_METHODS
);
assert_eq!(
response["result"]["authoring_capabilities"]["deploy_command"],
json!("rkat mob run")
);
Ok(())
}
#[tokio::test]
async fn unified_rpc_dispatches_mobpack_authoring_methods()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-authoring-dispatch-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/mobpacks/schema",
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert!(response["error"].is_null(), "{response:#?}");
assert_eq!(
response["result"]["media_type"],
json!("application/vnd.meerkat.mobpack")
);
assert_eq!(
response["result"]["commands"]["deploy_rpc"],
json!("mobkit/mobpacks/deploy")
);
assert_eq!(
response["result"]["deploy_settings"]["runtime_backed"],
json!(true)
);
assert_eq!(
response["result"]["deploy_settings"]["authoring_provider"]["runtime_binding"],
json!("bound")
);
assert_eq!(
response["result"]["deploy_settings"]["provenance"]["source"],
json!("UnifiedRuntime.authoring_provider.deploy_target")
);
assert!(response["result"]["sample_mobpacks"].is_null());
assert!(response["result"]["agent_definitions"].is_null());
let catalogs: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/mobpacks/catalogs",
})
.to_string(),
Duration::from_secs(1),
None,
None,
)
.await,
)?;
assert!(catalogs["error"].is_null(), "{catalogs:#?}");
assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
assert_eq!(
catalogs["result"]["authoring_provider"]["id"],
json!("unified_runtime")
);
assert_eq!(
catalogs["result"]["authoring_provider"]["runtime_binding"],
json!("bound")
);
assert_eq!(
catalogs["result"]["sources"]["runtime"],
json!("unified_runtime")
);
assert_eq!(
catalogs["result"]["sources"]["runtime_binding"],
json!("bound")
);
assert!(
catalogs["result"]["runtime_unavailable_reason"].is_null(),
"{catalogs:#?}"
);
assert_eq!(
catalogs["result"]["catalog_snapshot"]["runtime_backed"],
json!(true)
);
assert_eq!(
catalogs["result"]["authoring_provider"]["deploy_target"]["command"],
json!("rkat mob run")
);
assert!(
catalogs["result"]["authoring_provider"]["runtime_methods"]
.as_array()
.is_some_and(|methods| methods.contains(&json!("mobkit/mobpacks/deploy"))),
"{catalogs:#?}"
);
Ok(())
}
#[tokio::test]
async fn reconcile_identity_passes_mob_definition_to_roster_provider()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-reconcile-roster-context-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-reconcile-roster-context-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let roster_provider = Arc::new(ContextRequiredEmptyRosterProvider::default());
let identity_ctx = IdentityFirstContext {
runtime: identity_rt,
roster_provider: roster_provider.clone(),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: Some(runtime.mob_handle().definition().clone()),
};
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/reconcile_identity",
"params": {},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(response["error"].is_null(), "{response:#?}");
assert_eq!(
roster_provider.missing_definition_calls(),
0,
"reconcile_identity must preserve mob_definition in roster context"
);
Ok(())
}
#[tokio::test]
async fn startup_ready_wait_retries_stale_member_success_and_error_after_generation_change()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
for stale_result in [
IdentityMemberReadiness::Ready,
IdentityMemberReadiness::Failed("stale member failure".to_string()),
] {
let identity_runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-startup-ready-generation-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let callback_runtime = identity_runtime.clone();
let callback_calls = calls.clone();
let wait = wait_identity_startup_ready(
&identity_runtime,
Duration::from_secs(1),
move |_member_ids, _remaining| {
let call = callback_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let callback_runtime = callback_runtime.clone();
let stale_result = stale_result.clone();
async move {
if call == 0 {
callback_runtime.test_supersede_identity_bootstrap_ready();
stale_result
} else {
IdentityMemberReadiness::Ready
}
}
},
)
.await?;
assert!(wait.startup_ready);
assert!(!wait.timed_out);
assert!(wait.status.ready);
assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 2);
}
Ok(())
}
#[tokio::test]
async fn startup_ready_wait_reports_timeout_and_terminal_broken_status()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let timeout_runtime = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-startup-ready-timeout-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let timed_out = wait_identity_startup_ready(
&timeout_runtime,
Duration::from_secs(1),
|_member_ids, _remaining| async { IdentityMemberReadiness::TimedOut },
)
.await?;
assert!(timed_out.timed_out);
assert!(!timed_out.startup_ready);
let broken_runtime = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-startup-ready-broken-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
broken_runtime.test_fail_identity_bootstrap("injected terminal failure");
let readiness_calls = std::sync::atomic::AtomicUsize::new(0);
let broken = wait_identity_startup_ready(
&broken_runtime,
Duration::from_secs(1),
|_member_ids, _remaining| {
readiness_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
async { IdentityMemberReadiness::Ready }
},
)
.await?;
assert!(!broken.timed_out);
assert!(!broken.startup_ready);
assert!(!broken.status.ready);
assert_eq!(broken.status.counts.broken, 1);
assert!(
broken
.status
.error
.as_deref()
.is_some_and(|error| error.contains("injected terminal failure"))
);
assert_eq!(
readiness_calls.load(std::sync::atomic::Ordering::SeqCst),
0,
"terminal bootstrap failure must not enter member readiness"
);
Ok(())
}
#[tokio::test]
async fn identity_bootstrap_status_and_wait_rpc_are_typed_and_advertised()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-identity-bootstrap-status-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-identity-bootstrap-status-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
})),
roster_provider: Arc::new(ContextRequiredEmptyRosterProvider::default()),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: Some(runtime.mob_handle().definition().clone()),
};
let runtime_ref = &runtime;
let identity_ctx_ref = &identity_ctx;
let call = move |method: &'static str, params: Value| {
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": method,
"method": method,
"params": params,
})
.to_string();
async move {
handle_unified_rpc_json(
runtime_ref,
&request,
Duration::from_secs(1),
None,
Some(identity_ctx_ref),
)
.await
}
};
let status: Value = serde_json::from_str(
&call("mobkit/status_identity_bootstrap", serde_json::json!({})).await,
)?;
assert_eq!(status["result"]["mode"]["mode"], json!("eager_materialize"));
assert_eq!(status["result"]["complete"], json!(true));
assert_eq!(status["result"]["ready"], json!(true));
let waited: Value = serde_json::from_str(
&call(
"mobkit/wait_identity_bootstrap",
serde_json::json!({"target": "materialized", "timeout_ms": 0}),
)
.await,
)?;
assert_eq!(waited["result"]["timed_out"], json!(false));
assert_eq!(waited["result"]["target"], json!("materialized"));
let default_waited: Value = serde_json::from_str(
&call("mobkit/wait_identity_bootstrap", serde_json::json!({})).await,
)?;
assert_eq!(default_waited["result"]["target"], json!("materialized"));
let startup_ready: Value = serde_json::from_str(
&call(
"mobkit/wait_identity_bootstrap",
serde_json::json!({"target": "startup_ready", "timeout_ms": 0}),
)
.await,
)?;
assert_eq!(startup_ready["result"]["timed_out"], json!(false));
assert_eq!(startup_ready["result"]["startup_ready"], json!(true));
for (params, expected_message) in [
(serde_json::json!(null), "params must be an object"),
(serde_json::json!([]), "params must be an object"),
(
serde_json::json!({"target": null}),
"target must be a string",
),
(serde_json::json!({"target": 42}), "target must be a string"),
(
serde_json::json!({"target": "not_ready"}),
"target must be 'materialized' or 'startup_ready'",
),
(
serde_json::json!({"timeout_ms": null}),
"timeout_ms must be a non-negative integer",
),
(
serde_json::json!({"timeout_ms": -1}),
"timeout_ms must be a non-negative integer",
),
(
serde_json::json!({"timeout_ms": 1.5}),
"timeout_ms must be a non-negative integer",
),
(
serde_json::json!({"timeout_ms": true}),
"timeout_ms must be a non-negative integer",
),
(
serde_json::json!({"unexpected": true}),
"unsupported parameter: unexpected",
),
] {
let response: Value =
serde_json::from_str(&call("mobkit/wait_identity_bootstrap", params).await)?;
assert_eq!(response["error"]["code"], json!(-32602));
assert_eq!(response["error"]["message"], json!(expected_message));
}
let capabilities: Value =
serde_json::from_str(&call("mobkit/capabilities", serde_json::json!({})).await)?;
let methods = capabilities["result"]["methods"]
.as_array()
.expect("methods array");
assert!(methods.contains(&json!("mobkit/status_identity_bootstrap")));
assert!(methods.contains(&json!("mobkit/wait_identity_bootstrap")));
Ok(())
}
#[tokio::test]
async fn reset_identity_preserves_mob_definition_for_roster_reprofile()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_reprofile_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-reset-roster-context-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let continuity_store = Arc::new(LocalContinuityStore::in_memory()?);
let lease_provider = Arc::new(LocalLeaseProvider::new());
let bridge = Arc::new(RpcResetTestBridge::default());
let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: continuity_store.clone(),
lease_provider: lease_provider.clone(),
runtime_instance_id: "rpc-reset-roster-context-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(bridge.clone()),
default_timeout: None,
}));
let identity = AgentIdentity::parse("domain:security")?;
let initial_grants = lease_provider
.acquire_leases(
std::slice::from_ref(&identity),
"rpc-reset-roster-context-test",
)
.await?;
let initial_grant = match initial_grants.get(&identity) {
Some(LeaseAcquireResult::Acquired(grant)) => grant.clone(),
other => return Err(format!("expected acquired lease, got {other:?}").into()),
};
let initial_record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:domain:security:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
continuity_store
.upsert_continuity_record(&initial_record, initial_grant.fencing_token)
.await?;
identity_rt
.register(
rpc_durable_spec(identity.as_str(), "domain"),
IdentityLifecycleState::Active,
Some(initial_record),
Some(initial_grant),
)
.await;
let roster_provider = Arc::new(ContextRequiredStaticRosterProvider::new(vec![
rpc_durable_spec(identity.as_str(), "security"),
]));
let identity_ctx = IdentityFirstContext {
runtime: identity_rt.clone(),
roster_provider: roster_provider.clone(),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: Some(runtime.mob_handle().definition().clone()),
};
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/reset",
"params": { "identity": identity.as_str() },
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(response["error"].is_null(), "{response:#?}");
assert_eq!(
roster_provider.missing_definition_calls(),
0,
"mobkit/reset must preserve mob_definition when installing the reset roster provider"
);
assert_eq!(
bridge
.create_calls
.load(std::sync::atomic::Ordering::SeqCst),
1,
"reset should rebuild through the bridge"
);
let created = bridge
.last_create_spec()
.await
.expect("reset should record created spec");
assert_eq!(created.profile.as_str(), "security");
assert_eq!(
identity_rt
.status(&identity)
.await?
.profile
.expect("identity should keep profile")
.as_str(),
"security"
);
Ok(())
}
#[tokio::test]
async fn unified_rpc_agent_memory_remember_writes_identity_scoped_record()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-agent-memory-remember-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-agent-memory-remember-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let store = Arc::new(MarkdownAgentMemoryStore::open(
temp_dir.path().join("agent-memory"),
)?);
let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> = store.clone();
identity_rt
.set_agent_memory(Some(
crate::identity_first::AgentMemoryRuntimeInjector::new(
provider.clone(),
crate::identity_first::AgentMemoryConfig::default(),
),
))
.await;
let memory_identity = AgentIdentity::parse("identity:luka")?;
identity_rt
.register(
DurableAgentSpec {
identity: memory_identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: Default::default(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
None,
None,
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: Some(provider),
mob_definition: None,
};
let capabilities: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/capabilities",
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/remember"))),
"{capabilities:#?}"
);
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/recall"))),
"{capabilities:#?}"
);
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/forget"))),
"{capabilities:#?}"
);
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| !methods.contains(&json!("mobkit/agent_memory/update"))),
"{capabilities:#?}"
);
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| !methods.contains(&json!("mobkit/agent_memory/manifest"))),
"{capabilities:#?}"
);
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/agent_memory/remember",
"params": {
"identity": "identity:luka",
"realm": "family",
"title": "School pickup",
"body": "Pickup is before calendar planning.",
"tags": ["family", "calendar"]
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(response["error"].is_null(), "{response:#?}");
let memory_id = response["result"]["memory_id"]
.as_str()
.ok_or("memory_id should be present")?
.to_string();
assert_eq!(response["result"]["title"], json!("School pickup"));
assert_eq!(
response["result"]["body"],
json!("Pickup is before calendar planning.")
);
assert_eq!(response["result"]["tags"], json!(["calendar", "family"]));
let records = store.read_records("family", &memory_identity)?;
assert_eq!(records.len(), 1);
assert_eq!(records[0].title, "School pickup");
assert_eq!(records[0].tags, vec!["calendar", "family"]);
let recall_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 3,
"method": "mobkit/agent_memory/recall",
"params": {
"identity": "identity:luka",
"realm": "family",
"selection": "contextual",
"query_terms": ["pickup"],
"max_entries": 4
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(recall_response["error"].is_null(), "{recall_response:#?}");
assert_eq!(
recall_response["result"]["records"]
.as_array()
.map(Vec::len),
Some(1)
);
assert_eq!(
recall_response["result"]["records"][0]["body"],
json!("Pickup is before calendar planning.")
);
let forget_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 4,
"method": "mobkit/agent_memory/forget",
"params": {
"identity": "identity:luka",
"realm": "family",
"memory_id": memory_id.clone()
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(forget_response["error"].is_null(), "{forget_response:#?}");
assert_eq!(forget_response["result"]["memory_id"], json!(memory_id));
assert_eq!(forget_response["result"]["deleted"], json!(true));
assert!(store.read_records("family", &memory_identity)?.is_empty());
let recall_after_forget_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 5,
"method": "mobkit/agent_memory/recall",
"params": {
"identity": "identity:luka",
"realm": "family",
"selection": "always"
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(
recall_after_forget_response["error"].is_null(),
"{recall_after_forget_response:#?}"
);
assert_eq!(
recall_after_forget_response["result"]["records"]
.as_array()
.map(Vec::len),
Some(0)
);
let unknown_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 6,
"method": "mobkit/agent_memory/remember",
"params": {
"identity": "identity:unknown",
"title": "Orphan",
"body": "This should not be written."
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert_eq!(unknown_response["error"]["code"], json!(-32602));
assert!(
unknown_response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("unknown identity")),
"{unknown_response:#?}"
);
let unknown_forget_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 7,
"method": "mobkit/agent_memory/forget",
"params": {
"identity": "identity:unknown",
"memory_id": "mem-missing"
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert_eq!(unknown_forget_response["error"]["code"], json!(-32602));
assert!(
unknown_forget_response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("unknown identity")),
"{unknown_forget_response:#?}"
);
Ok(())
}
#[tokio::test]
async fn unified_rpc_agent_memory_update_and_manifest_over_sqlite_store()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-agent-memory-sqlite-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-agent-memory-sqlite-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let store = Arc::new(crate::memory::SqliteAgentMemoryStore::open(
temp_dir.path().join("agent-memory"),
)?);
let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> = store.clone();
identity_rt
.set_agent_memory(Some(
crate::identity_first::AgentMemoryRuntimeInjector::new(
provider.clone(),
crate::identity_first::AgentMemoryConfig::default(),
),
))
.await;
let memory_identity = AgentIdentity::parse("identity:luka")?;
identity_rt
.register(
DurableAgentSpec {
identity: memory_identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: Default::default(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
None,
None,
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: Some(provider),
mob_definition: None,
};
let capabilities: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/capabilities",
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
for method in [
"mobkit/agent_memory/recall",
"mobkit/agent_memory/remember",
"mobkit/agent_memory/forget",
"mobkit/agent_memory/update",
"mobkit/agent_memory/manifest",
] {
assert!(
capabilities["result"]["methods"]
.as_array()
.is_some_and(|methods| methods.contains(&json!(method))),
"missing {method}: {capabilities:#?}"
);
}
let remember_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/agent_memory/remember",
"params": {
"identity": "identity:luka",
"realm": "family",
"title": "School pickup",
"body": "Pickup is before calendar planning.",
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(
remember_response["error"].is_null(),
"{remember_response:#?}"
);
let memory_id = remember_response["result"]["memory_id"]
.as_str()
.ok_or("memory_id should be present")?
.to_string();
let update_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 3,
"method": "mobkit/agent_memory/update",
"params": {
"identity": "identity:luka",
"realm": "family",
"memory_id": memory_id.clone(),
"title": "School pickup",
"body": "Pickup moved to after calendar planning.",
"tags": ["family"]
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(update_response["error"].is_null(), "{update_response:#?}");
let new_id = update_response["result"]["memory_id"]
.as_str()
.ok_or("updated memory_id should be present")?
.to_string();
assert_ne!(new_id, memory_id);
assert_eq!(update_response["result"]["supersedes"], json!(memory_id));
let recall_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 4,
"method": "mobkit/agent_memory/recall",
"params": {
"identity": "identity:luka",
"realm": "family",
"selection": "always"
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(recall_response["error"].is_null(), "{recall_response:#?}");
assert_eq!(
recall_response["result"]["records"]
.as_array()
.map(Vec::len),
Some(1)
);
assert_eq!(
recall_response["result"]["records"][0]["memory_id"],
json!(new_id)
);
let manifest_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 5,
"method": "mobkit/agent_memory/manifest",
"params": {
"identity": "identity:luka",
"realm": "family",
"tier": "working_set",
"k": 4
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert!(
manifest_response["error"].is_null(),
"{manifest_response:#?}"
);
let records = manifest_response["result"]["records"]
.as_array()
.ok_or("manifest records array")?;
assert_eq!(records.len(), 1, "{manifest_response:#?}");
assert_eq!(records[0]["id"], json!(new_id));
assert_eq!(records[0]["kind"], json!("fact"));
assert_eq!(records[0]["age_days"], json!(0));
assert!(
records[0].get("body").is_none(),
"manifest is an index, never a dump: {manifest_response:#?}"
);
let bad_tier_response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 6,
"method": "mobkit/agent_memory/manifest",
"params": {
"identity": "identity:luka",
"tier": "everything"
},
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert_eq!(bad_tier_response["error"]["code"], json!(-32602));
Ok(())
}
#[tokio::test]
async fn unified_rpc_agent_memory_capabilities_do_not_advertise_read_only_writes()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-agent-memory-read-only-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-agent-memory-read-only-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> =
Arc::new(ReadOnlyAgentMemoryProvider);
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: Some(provider),
mob_definition: None,
};
let capabilities: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/capabilities",
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
let methods = capabilities["result"]["methods"]
.as_array()
.ok_or("methods should be an array")?;
assert!(methods.contains(&json!("mobkit/agent_memory/recall")));
assert!(!methods.contains(&json!("mobkit/agent_memory/remember")));
assert!(!methods.contains(&json!("mobkit/agent_memory/forget")));
Ok(())
}
#[test]
fn identity_lease_lost_maps_off_capability_unavailable_code() {
let identity = AgentIdentity::parse("review:singleton").expect("valid identity");
let err = crate::identity_first::IdentityRuntimeError::LeaseLost(identity);
let response = identity_error_response(json!("req-1"), &err);
let error = response.error.expect("lease-lost must surface an error");
assert_ne!(
error.code, -32004,
"LeaseLost must not use the capability code"
);
assert_eq!(
error.code, -32005,
"LeaseLost has its own identity-plane code"
);
assert!(error.message.contains("lease lost"));
}
#[test]
fn mobpack_authoring_rpc_helper_preserves_runtime_catalog_binding()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let runtime = crate::mobpack::MobpackRuntimeCatalogState {
loaded_modules: vec!["editor-host".to_string()],
runtime_methods: vec![
"mobkit/mobpacks/catalogs".to_string(),
"mobkit/mobpacks/apply_operation".to_string(),
"mobkit/mobpacks/deploy".to_string(),
],
has_contact_directory: true,
has_peer_mob_handles: false,
has_inproc_contacts: false,
runtime_flow_rows: vec![json!({
"id": "runtime_rpc_main",
"source": "mobkit/runtime/flow_projection",
"document": {
"mob_id": "runtime_rpc",
"flow": { "name": "main", "steps": [] },
"members": []
},
"validation": { "ok": true }
})],
runtime_agent_definition_sources: vec![json!({
"id": "runtime_profiles_rpc",
"name": "Runtime RPC profiles",
"source": "mobkit/runtime/agent-definitions",
"document": {
"mob_id": "runtime_rpc",
"members": [{
"id": "m_runtime_reviewer",
"name": "Runtime reviewer",
"role": "runtime_reviewer",
"profileBinding": "inline",
"model": "gpt-5.5",
"runtimeMode": "turn_driven",
"tools": ["builtins"],
"skills": ["mob.runtime.review"],
"schema": ""
}],
"schemas": []
}
})],
runtime_skill_realms: vec![json!({
"id": "runtime_rpc",
"label": "Runtime RPC",
"source": "mobkit/runtime/skills",
"skills": [{
"id": "mob.runtime.review",
"label": "Runtime review",
"source": "inline",
"content": "Review runtime work."
}]
})],
};
let catalogs = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/mobpacks/catalogs",
&json!({}),
json!(1),
Some(&runtime),
)
.expect("catalogs method");
let catalogs: Value = serde_json::to_value(catalogs)?;
assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
assert_eq!(
catalogs["result"]["authoring_provider"]["runtime_binding"],
json!("bound")
);
assert_eq!(
catalogs["result"]["runtime_flows"][0]["id"],
json!("runtime_rpc_main")
);
let listed = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/mobpacks/list",
&json!({}),
json!(2),
Some(&runtime),
)
.expect("list method");
let listed: Value = serde_json::to_value(listed)?;
assert_eq!(listed["result"]["runtime_backed"], json!(true));
assert_eq!(listed["result"]["rows"][0]["id"], json!("runtime_rpc_main"));
let fetched = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/mobpacks/get",
&json!({ "id": "runtime_rpc_main" }),
json!(3),
Some(&runtime),
)
.expect("get method");
let fetched: Value = serde_json::to_value(fetched)?;
assert_eq!(fetched["result"]["runtime_backed"], json!(true));
assert_eq!(fetched["result"]["row"]["id"], json!("runtime_rpc_main"));
let definitions = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/agent_definitions/list",
&json!({}),
json!(4),
Some(&runtime),
)
.expect("agent definitions method");
let definitions: Value = serde_json::to_value(definitions)?;
let runtime_definition = definitions["result"]["agent_definitions"]
.as_array()
.and_then(|rows| {
rows.iter()
.find(|row| row["sourceOrigin"] == "mobkit/runtime/agent-definitions")
})
.expect("runtime profile definition");
assert_eq!(runtime_definition["role"], json!("runtime_reviewer"));
assert_eq!(
runtime_definition["toolDefinitions"][0]["id"],
json!("builtins")
);
assert_eq!(
runtime_definition["skillDefinitions"][0]["id"],
json!("mob.runtime.review")
);
assert_eq!(
definitions["result"]["catalog_snapshot"]["runtime_backed"],
json!(true)
);
let tools = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/tools/catalog",
&json!({}),
json!(5),
Some(&runtime),
)
.expect("tools catalog method");
let tools: Value = serde_json::to_value(tools)?;
let mob_tool = tools["result"]["tool_catalog"]
.as_array()
.expect("tools")
.iter()
.find(|tool| tool["id"] == "mob")
.expect("mob tool");
assert_eq!(mob_tool["runtime_availability"]["available"], json!(false));
let agents = super::handle_mobpack_authoring_rpc_with_runtime(
"mobkit/agent_definitions/list",
&json!({}),
json!(3),
Some(&runtime),
)
.expect("agent definitions method");
let agents: Value = serde_json::to_value(agents)?;
assert_eq!(agents["result"]["runtime_backed"], json!(true));
let planner = agents["result"]["agent_definitions"]
.as_array()
.expect("agent definitions")
.iter()
.find(|definition| definition["role"] == "planner")
.expect("planner definition");
let planner_mob_tool = planner["toolDefinitions"]
.as_array()
.expect("planner tools")
.iter()
.find(|tool| tool["id"] == "mob")
.expect("planner mob tool");
assert_eq!(
planner_mob_tool["runtimeAvailability"]["state"],
json!("unavailable")
);
Ok(())
}
#[test]
fn module_rpc_dispatches_mobpack_authoring_methods()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let config = MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "module-rpc-authoring-dispatch-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
};
let mut runtime = crate::start_mobkit_runtime(config, Vec::new(), Duration::from_secs(1))?;
let capabilities: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
&mut runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/capabilities",
})
.to_string(),
Duration::from_secs(1),
))?;
let methods = capabilities["result"]["methods"]
.as_array()
.expect("methods array")
.iter()
.filter_map(Value::as_str)
.collect::<Vec<_>>();
for method in super::MOBPACK_AUTHORING_METHODS {
assert!(
methods.contains(method),
"missing authoring method {method}"
);
}
assert_eq!(
capabilities["result"]["authoring_capabilities"]["runtime_mutation"],
json!(false)
);
assert_eq!(
capabilities["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
);
let schema: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
&mut runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/mobpacks/schema",
})
.to_string(),
Duration::from_secs(1),
))?;
assert!(schema["error"].is_null(), "{schema:#?}");
assert_eq!(
schema["result"]["commands"]["deploy_rpc"],
json!("mobkit/mobpacks/deploy")
);
assert!(schema["result"]["agent_definitions"].is_null());
let _ = runtime.shutdown();
Ok(())
}
#[test]
fn generated_runtime_ids_match_their_durable_identity_prefix() {
assert!(!rpc_member_id_matches_durable_identity(
"rt:review:singleton:0",
"review:singleton",
));
assert!(!rpc_member_id_matches_durable_identity(
"review:singleton:gen1",
"review:singleton",
));
assert!(!rpc_member_id_matches_durable_identity(
"review:singleton:1",
"review:singleton",
));
assert!(!rpc_member_id_matches_durable_identity(
"rt:reviewer:singleton:0",
"review:singleton",
));
assert!(!rpc_member_id_matches_durable_identity(
"rt:review:singleton:qa:0",
"review:singleton",
));
assert!(!rpc_member_id_matches_durable_identity(
"review:singleton:qa",
"review:singleton",
));
assert_eq!(
rpc_runtime_alias_generation("rt:review:singleton:7", "review:singleton"),
Some(7)
);
assert_eq!(
rpc_runtime_alias_generation("rt:review:singleton:8", "review:other"),
None
);
}
#[test]
fn rpc_live_identity_visibility_matches_delegate_projection_labels() {
assert!(rpc_live_identity_alias_visible("worker", &BTreeMap::new()));
let mut labels = BTreeMap::new();
labels.insert("role".to_string(), "delegate".to_string());
labels.insert("source_mob_id".to_string(), "mob-a".to_string());
labels.insert("agent_identity".to_string(), "review:singleton".to_string());
assert!(!rpc_live_identity_alias_visible("worker", &labels));
assert!(!rpc_live_identity_alias_visible("delegate", &labels));
}
#[test]
fn ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error> {
let response: Value = serde_json::from_str(&error_response(
json!(1),
-32602,
"ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
))?;
assert_eq!(
response["error"]["data"]["kind"],
json!("ambiguous_live_identity_alias")
);
assert_eq!(
response["error"]["data"]["identity"],
json!("review:singleton")
);
assert_eq!(
response["error"]["data"]["candidates"],
json!(["rt:review:singleton:0", "rt:review:singleton:1"])
);
Ok(())
}
#[test]
fn wrapped_ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error>
{
let response: Value = serde_json::from_str(&error_response(
json!(1),
-32602,
"invalid identity: ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
))?;
assert_eq!(
response["error"]["data"]["kind"],
json!("ambiguous_live_identity_alias")
);
assert_eq!(
response["error"]["data"]["identity"],
json!("review:singleton")
);
assert_eq!(
response["error"]["data"]["candidates"],
json!(["rt:review:singleton:0", "rt:review:singleton:1"])
);
Ok(())
}
#[tokio::test]
async fn runtime_id_live_only_resolution_rejects_duplicate_projected_identity()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-identity-alias-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
let mut labels = BTreeMap::new();
labels.insert("agent_identity".to_string(), "review:singleton".to_string());
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_id.to_string(),
Some("You are a duplicate Review Agent.".into()),
None,
None,
)
.with_labels(labels),
)
.await?;
}
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-identity-alias-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let err =
resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:0")
.await
.expect_err("runtime-id live-only fallback should reject duplicate durable alias");
assert!(
err.contains("ambiguous live identity alias review:singleton"),
"unexpected error: {err}"
);
Ok(())
}
#[tokio::test]
async fn durable_resolution_prefers_registered_live_binding_over_stale_duplicates()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-identity-alias-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
let mut labels = BTreeMap::new();
labels.insert("agent_identity".to_string(), "review:singleton".to_string());
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_id.to_string(),
Some("You are a duplicate Review Agent.".into()),
None,
None,
)
.with_labels(labels),
)
.await?;
}
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-identity-alias-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse("review:singleton")?;
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:1")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
};
identity_rt
.register(
DurableAgentSpec {
identity,
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(record),
None,
)
.await;
let target =
resolve_rpc_identity_control_target(&runtime, &identity_rt, "review:singleton").await?;
assert_eq!(target.identity.as_str(), "review:singleton");
assert_eq!(
target
.live
.as_ref()
.map(|alias| alias.runtime_member_id.as_str()),
Some("rt:review:singleton:1")
);
let target =
resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:1")
.await?;
assert_eq!(
target
.live
.as_ref()
.map(|alias| alias.runtime_member_id.as_str()),
Some("rt:review:singleton:1")
);
let stale_target =
resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:0")
.await?;
assert_eq!(
stale_target
.live
.as_ref()
.map(|alias| alias.runtime_member_id.as_str()),
Some("rt:review:singleton:0")
);
let stale_response =
super::rpc_stale_live_alias_error_response(&identity_rt, &stale_target, json!(99))
.await
.expect("old reset generation should be rejected as stale");
assert_eq!(
stale_response
.error
.as_ref()
.and_then(|error| error.data.as_ref())
.and_then(|data| data.get("kind")),
Some(&json!("stale_identity_runtime_binding"))
);
assert_eq!(
stale_response
.error
.as_ref()
.and_then(|error| error.data.as_ref())
.and_then(|data| data.get("registered_runtime_member_id")),
Some(&json!("rt:review:singleton:1"))
);
assert_eq!(
stale_response
.error
.as_ref()
.and_then(|error| error.data.as_ref())
.and_then(|data| data.get("live_runtime_member_id")),
Some(&json!("rt:review:singleton:0"))
);
Ok(())
}
#[tokio::test]
async fn current_generation_resolvers_ignore_stale_reset_member_rows()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let mut runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-current-generation-target-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let durable_identity = "review:current-generation";
let stale_alias = "rt:review:current-generation:0";
let current_alias = "rt:review:current-generation:1";
for runtime_id in [stale_alias, current_alias] {
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_id.to_string(),
None,
None,
None,
)
.with_labels(BTreeMap::from([(
"agent_identity".to_string(),
durable_identity.to_string(),
)])),
)
.await?;
}
let identity_runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-current-generation-target-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity = AgentIdentity::parse(durable_identity)?;
identity_runtime
.register(
rpc_durable_spec(durable_identity, "worker"),
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(current_alias)?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
}),
None,
)
.await;
let expected_session = runtime
.mob_handle()
.resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(current_alias))
.await
.ok_or("current member must have a bridge session")?;
let params = json!({"identity": durable_identity});
let target = identity_runtime
.member_alias_lifecycle_target(durable_identity)
.await?
.ok_or("durable identity must resolve to lifecycle authority")?;
let operation_runtime = Arc::clone(&identity_runtime);
let handle = runtime.mob_handle();
let resolved_session = IdentityRuntime::run_member_alias_targets_operation_tracked(
vec![target],
move || async move {
resolve_live_target(&handle, Some(&operation_runtime), true, ¶ms).await
},
)
.await?;
assert_eq!(resolved_session, Some(expected_session));
let fallback_error = resolve_live_target(
&runtime.mob_handle(),
None,
false,
&json!({"identity": durable_identity}),
)
.await
.expect_err("two live generations without identity authority must be ambiguous");
assert!(
fallback_error.contains("ambiguous live member alias"),
"unexpected fallback error: {fallback_error}"
);
let cross_mob_fallback = runtime
.local_member_peer_info(durable_identity)
.await
.expect_err("cross-mob fallback without identity authority must be ambiguous");
assert!(
cross_mob_fallback
.to_string()
.contains("ambiguous durable member alias"),
"unexpected cross-mob fallback error: {cross_mob_fallback}"
);
runtime.attach_identity_first_context(Arc::new(
crate::identity_first::IdentityFirstRuntimeContext::new(
Arc::clone(&identity_runtime),
Arc::new(EmptyRosterProvider),
None,
None,
None,
),
));
let (_, comms_name, _) = runtime.local_member_peer_info(durable_identity).await?;
assert!(
comms_name.ends_with(crate::member_comms_id::mob_member_id_str(current_alias).as_ref()),
"peer info must use the current generation: {comms_name}"
);
let remote_comms_name = "remote-mob/worker/peer";
let remote_pubkey = [42_u8; 32];
let remote_peer_id =
meerkat_core::comms::PeerId::from_ed25519_pubkey(&remote_pubkey).to_string();
runtime
.wire_local(
durable_identity,
remote_comms_name,
&remote_peer_id,
&format!("inproc://{remote_comms_name}"),
Some(remote_pubkey),
)
.await?;
let current = runtime
.mob_handle()
.get_member(&crate::member_comms_id::mob_member_id(current_alias))
.await?
.ok_or("current member missing")?;
let stale = runtime
.mob_handle()
.get_member(&crate::member_comms_id::mob_member_id(stale_alias))
.await?
.ok_or("stale member missing")?;
assert!(
current
.wired_to
.iter()
.any(|peer| peer.as_str() == remote_comms_name),
"current generation must receive the wire"
);
assert!(
stale
.wired_to
.iter()
.all(|peer| peer.as_str() != remote_comms_name),
"stale generation must not receive the wire"
);
runtime.shutdown().await;
Ok(())
}
#[tokio::test]
async fn durable_resolution_rejects_hidden_registered_live_binding()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-hidden-bound-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
"rt:review:singleton:0".to_string(),
Some("You are a hidden Review Agent.".into()),
None,
None,
)
.with_labels(BTreeMap::from([
("agent_identity".to_string(), "review:singleton".to_string()),
("role".to_string(), "delegate".to_string()),
("source_mob_id".to_string(), "upstream".to_string()),
])),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-hidden-bound-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse("review:singleton")?;
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity,
agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
None,
)
.await;
for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
let err =
resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
.await
.expect_err("hidden registered live binding must not resolve");
assert!(
err.contains("identity hidden by policy"),
"unexpected error for {requested_identity}: {err}"
);
}
Ok(())
}
#[tokio::test]
async fn live_only_hidden_alias_reports_policy_error()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-hidden-live-only-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
"rt:review:singleton:0".to_string(),
Some("You are a hidden Review Agent.".into()),
None,
None,
)
.with_labels(BTreeMap::from([
("agent_identity".to_string(), "review:singleton".to_string()),
("role".to_string(), "delegate".to_string()),
("source_mob_id".to_string(), "upstream".to_string()),
])),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-hidden-live-only-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
let err =
resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
.await
.expect_err("hidden live-only alias must not collapse into unknown identity");
assert!(
err.contains("identity hidden by policy"),
"unexpected error for {requested_identity}: {err}"
);
}
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
let response: Value = serde_json::from_str(
&handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/status_identity",
"params": { "identity": requested_identity },
})
.to_string(),
Duration::from_secs(1),
None,
Some(&identity_ctx),
)
.await,
)?;
assert_eq!(
response["error"]["data"]["kind"],
json!("identity_hidden_by_policy"),
"unexpected hidden response for {requested_identity}: {response:#?}"
);
}
Ok(())
}
#[tokio::test]
async fn live_only_resolution_rejects_runtime_member_bound_to_other_durable_identity()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-identity-alias-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let mut labels = BTreeMap::new();
labels.insert("agent_identity".to_string(), "other:singleton".to_string());
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
"rt:review:singleton:0".to_string(),
Some("You are a wrong-projected Review Agent.".into()),
None,
None,
)
.with_labels(labels),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-identity-alias-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse("review:singleton")?;
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(record),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let err = resolve_rpc_identity_control_target(&runtime, &identity_rt, "other:singleton")
.await
.expect_err("wrong-projected live alias must not resolve as live-only");
assert!(
err.contains("identity runtime binding belongs to review:singleton"),
"unexpected error: {err}"
);
Ok(())
}
#[tokio::test]
async fn send_message_resolves_bare_durable_identity_through_identity_bridge()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-send-message-identity-bridge-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
for (runtime_id, durable) in [
("rt:atlas-base-001:0", "atlas-base-001"),
("rt:draco-base-001:0", "draco-base-001"),
] {
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_id.to_string(),
Some("You are a swarm base agent.".into()),
None,
None,
)
.with_labels(BTreeMap::from([(
"agent_identity".to_string(),
durable.to_string(),
)])),
)
.await?;
}
runtime
.spawn(SpawnMemberSpec::from_wire(
"worker".to_string(),
"draco-base-001".to_string(),
Some("You are the raw roster member.".into()),
None,
None,
))
.await?;
let session_service = runtime
.mob_runtime()
.session_service()
.cloned()
.expect("test mob spec has a session service");
let bridge: Arc<dyn crate::identity_first::SessionBridge> = Arc::new(
crate::identity_first::MobSessionBridge::with_session_service(
runtime.mob_handle(),
session_service,
),
);
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-send-message-identity-bridge-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(bridge),
default_timeout: None,
})
.with_runtime_services(crate::identity_first::AgentRuntimeServices::new(
runtime.mob_handle(),
));
for (durable, runtime_id) in [
("atlas-base-001", "rt:atlas-base-001:0"),
("draco-base-001", "rt:draco-base-001:0"),
] {
let identity = AgentIdentity::parse(durable)?;
let session_id = runtime
.mob_handle()
.resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(runtime_id))
.await
.unwrap_or_else(meerkat_core::types::SessionId::new);
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(runtime_id)?,
session_id,
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
}
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
let send = |id: u64, params: Value| {
let runtime = &runtime;
let identity_ctx = &identity_ctx;
async move {
let raw = handle_unified_rpc_json(
runtime,
&json!({
"jsonrpc": "2.0",
"id": id,
"method": "mobkit/send_message",
"params": params,
})
.to_string(),
Duration::from_secs(10),
None,
Some(identity_ctx),
)
.await;
serde_json::from_str::<Value>(&raw)
}
};
let response = send(
1,
json!({ "member_id": "atlas-base-001", "message": "status check" }),
)
.await?;
assert!(
response["error"].is_null(),
"bare identity send must bridge-resolve: {response:#?}"
);
assert_eq!(response["result"]["accepted"], json!(true));
assert_eq!(response["result"]["member_id"], json!("atlas-base-001"));
let atlas_session = runtime
.mob_handle()
.resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
"rt:atlas-base-001:0",
))
.await
.expect("atlas runtime member has a bridge session after send")
.to_string();
assert_eq!(response["result"]["session_id"], json!(atlas_session));
let response = send(
2,
json!({
"member_id": "atlas-base-001",
"message": "steer: stand down",
"handling_mode": "steer",
}),
)
.await?;
assert!(
response["error"].is_null(),
"bare identity steer must bridge-resolve: {response:#?}"
);
assert_eq!(response["result"]["accepted"], json!(true));
let response = send(
3,
json!({ "member_id": "draco-base-001", "message": "raw roster delivery" }),
)
.await?;
assert!(
response["error"].is_null(),
"exact roster member send must keep raw semantics: {response:#?}"
);
assert_eq!(response["result"]["accepted"], json!(true));
let draco_raw_session = runtime
.mob_handle()
.resolve_bridge_session_id(&crate::member_comms_id::mob_member_id("draco-base-001"))
.await
.expect("bare draco member has a bridge session after send")
.to_string();
assert_eq!(response["result"]["session_id"], json!(draco_raw_session));
if let Some(draco_rt_session) = runtime
.mob_handle()
.resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
"rt:draco-base-001:0",
))
.await
{
assert_ne!(
response["result"]["session_id"],
json!(draco_rt_session.to_string()),
"exact member id match must not be shadowed by identity resolution"
);
}
let response = send(
4,
json!({ "member_id": "phantom-base-999", "message": "nobody home" }),
)
.await?;
assert_eq!(response["error"]["code"], json!(-32000), "{response:#?}");
let message = response["error"]["message"]
.as_str()
.expect("error message");
assert!(
message.starts_with("send_message failed:"),
"unexpected error message: {message}"
);
Ok(())
}
#[tokio::test]
async fn member_state_wire_vocabulary_is_lowercase_on_both_surfaces()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-state-vocabulary-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
runtime
.spawn(SpawnMemberSpec::from_wire(
"worker".to_string(),
"worker-one".to_string(),
None,
None,
None,
))
.await?;
let raw = handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/get_member",
"params": { "member_id": "worker-one" },
})
.to_string(),
Duration::from_secs(5),
None,
None,
)
.await;
let response: Value = serde_json::from_str(&raw)?;
assert_eq!(
response["result"]["state"],
json!("active"),
"member rows must speak the lowercase SDK vocabulary: {response:#?}"
);
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-state-vocabulary-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse("review:singleton")?;
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
let raw = handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/status_identity",
"params": { "identity": "review:singleton" },
})
.to_string(),
Duration::from_secs(5),
None,
Some(&identity_ctx),
)
.await;
let response: Value = serde_json::from_str(&raw)?;
assert_eq!(
response["result"]["state"],
json!("active"),
"identity status must speak the same lowercase vocabulary: {response:#?}"
);
Ok(())
}
#[tokio::test]
async fn send_message_pins_baseline_member_over_identity_fallback()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-send-message-baseline-pin-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
"rt:draco-base-001:0".to_string(),
None,
None,
None,
),
)
.await?;
runtime
.spawn(
SpawnMemberSpec::from_wire(
"worker".to_string(),
"draco-base-001".to_string(),
None,
None,
None,
)
.with_runtime_mode(meerkat_mob::MobRuntimeMode::TurnDriven),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-send-message-baseline-pin-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse("draco-base-001")?;
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:draco-base-001:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: Arc::new(identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
runtime
.mob_runtime()
.set_baseline_member_specs(vec![SpawnMemberSpec::new(
meerkat_mob::ProfileName::from("worker"),
meerkat_mob::AgentIdentity::from("draco-base-001"),
)])
.await;
runtime
.mob_handle()
.retire(meerkat_mob::AgentIdentity::from("draco-base-001"))
.await?;
let raw = handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/send_message",
"params": { "member_id": "draco-base-001", "message": "mid-reconcile send" },
})
.to_string(),
Duration::from_secs(10),
None,
Some(&identity_ctx),
)
.await;
let response: Value = serde_json::from_str(&raw)?;
assert_eq!(
response["error"]["code"],
json!(-32000),
"transiently-absent baseline member must keep raw member-id semantics \
instead of silently delivering through the identity bridge: {response:#?}"
);
let message = response["error"]["message"]
.as_str()
.expect("error message");
assert!(
message.starts_with("send_message failed:"),
"unexpected error message: {message}"
);
Ok(())
}
#[tokio::test]
async fn separate_identity_context_pins_generated_alias_for_cross_mob_and_fork()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Arc::new(
Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-separate-identity-context-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?,
);
let stale_alias = "rt:review:separate:0";
let current_alias = "rt:review:separate:1";
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
stale_alias.to_string(),
None,
None,
None,
)
.with_labels(BTreeMap::from([(
"agent_identity".to_string(),
"review:separate".to_string(),
)])),
)
.await?;
let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-separate-identity-context-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity = AgentIdentity::parse("review:separate")?;
identity_rt
.register(
rpc_durable_spec(identity.as_str(), "worker"),
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(current_alias)?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: identity_rt,
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
assert!(
runtime.identity_runtime().is_none(),
"the regression requires a context supplied only at dispatch"
);
for (method, params) in [
(
"mobkit/cross_mob/wire",
json!({
"local_member_id": stale_alias,
"remote_member_id": "remote-worker",
"remote_mob_id": "remote-mob",
}),
),
(
"mobkit/fork_helper",
json!({
"source_member_id": stale_alias,
"agent_identity": "fork-probe",
"task": "prove stale source rejection",
}),
),
] {
let response: Value = serde_json::from_str(
&super::handle_unified_rpc_json_arc(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": method,
"method": method,
"params": params,
})
.to_string(),
Duration::from_secs(5),
None,
Some(&identity_ctx),
)
.await,
)?;
assert_eq!(response["error"]["code"], json!(-32000), "{response:#?}");
assert!(
response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("stale runtime alias")),
"{method} must reject through identity authority before raw fallthrough: {response:#?}"
);
}
assert!(
runtime
.mob_handle()
.list_members_including_retiring()
.await
.iter()
.any(|member| {
crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str())
== stale_alias
}),
"the raw stale member remains present, proving rejection was not member-not-found"
);
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn embedded_rpc_uses_runtime_owned_identity_context_when_argument_is_absent()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let mut runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-owned-identity-context-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
let stale_alias = "rt:review:owned-context:0";
let current_alias = "rt:review:owned-context:1";
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
stale_alias.to_string(),
None,
None,
None,
)
.with_labels(BTreeMap::from([(
"agent_identity".to_string(),
"review:owned-context".to_string(),
)])),
)
.await?;
let identity_runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-owned-identity-context-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity = AgentIdentity::parse("review:owned-context")?;
identity_runtime
.register(
rpc_durable_spec(identity.as_str(), "worker"),
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(current_alias)?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
}),
None,
)
.await;
runtime.attach_identity_first_context(Arc::new(
crate::identity_first::IdentityFirstRuntimeContext::new(
identity_runtime,
Arc::new(EmptyRosterProvider),
None,
None,
None,
),
));
let runtime = Arc::new(runtime);
let list: Value = serde_json::from_str(
&super::handle_unified_rpc_json_arc(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/list_members",
"params": {},
})
.to_string(),
Duration::from_secs(5),
None,
None,
)
.await,
)?;
assert_eq!(list["result"], json!([]), "{list:#?}");
let get: Value = serde_json::from_str(
&super::handle_unified_rpc_json_arc(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 2,
"method": "mobkit/get_member",
"params": {"member_id": stale_alias},
})
.to_string(),
Duration::from_secs(5),
None,
None,
)
.await,
)?;
assert_eq!(
get["error"]["data"]["kind"],
json!("stale_identity_runtime_binding"),
"{get:#?}"
);
runtime.shutdown().await;
Ok(())
}
#[tokio::test]
async fn durable_ownership_observed_before_unknown_identity_never_live_fallbacks()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-durable-delete-race-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(1))
.build(),
)
.await?;
let durable_identity = "review:delete-race";
let runtime_alias = "rt:review:delete-race:0";
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_alias.to_string(),
None,
None,
None,
)
.with_labels(BTreeMap::from([(
"agent_identity".to_string(),
durable_identity.to_string(),
)])),
)
.await?;
let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-durable-delete-race-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
});
let identity = AgentIdentity::parse(durable_identity)?;
identity_rt
.register(
rpc_durable_spec(durable_identity, "worker"),
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(runtime_alias)?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity: identity.clone(),
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let target =
resolve_rpc_identity_control_target(&runtime, &identity_rt, durable_identity).await?;
assert!(
target.was_registered,
"resolution must capture durable ownership"
);
assert!(target.live.is_some(), "the race requires a live projection");
identity_rt
.remove(&identity)
.await
.expect("simulate a concurrent durable delete after resolution");
assert!(matches!(
identity_rt.status(&identity).await,
Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_))
));
assert!(
!super::rpc_live_only_fallback_allowed(&target, durable_identity),
"a later UnknownIdentity must preserve the observed durable delete"
);
let mut live_only = target.clone();
live_only.was_registered = false;
assert!(
!super::rpc_live_only_fallback_allowed(&live_only, durable_identity),
"a resolved generated alias must remain identity-owned even when registration vanished"
);
assert!(
!super::rpc_live_only_fallback_allowed(&live_only, runtime_alias),
"generated aliases never enter raw/live-only fallback"
);
live_only
.live
.as_mut()
.expect("live projection")
.runtime_member_id = "legacy-review-member".to_string();
assert!(
super::rpc_live_only_fallback_allowed(&live_only, durable_identity),
"genuine bare live-only identities retain compatibility fallback"
);
let _ = runtime.mob_handle().stop().await;
Ok(())
}
#[tokio::test]
async fn retire_member_current_runtime_alias_uses_identity_authority()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-owned-retire-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
let runtime_alias = "rt:review:owned:0";
spawn_identity_projection_fixture(
&runtime,
SpawnMemberSpec::from_wire(
"worker".to_string(),
runtime_alias.to_string(),
None,
None,
None,
),
)
.await?;
let continuity_store = Arc::new(LocalContinuityStore::in_memory()?);
let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store,
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "rpc-owned-retire-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity = AgentIdentity::parse("review:owned")?;
identity_rt
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(runtime_alias)?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity: identity.clone(),
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
let identity_ctx = IdentityFirstContext {
runtime: Arc::clone(&identity_rt),
roster_provider: Arc::new(EmptyRosterProvider),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: None,
};
let raw = handle_unified_rpc_json(
&runtime,
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/retire_member",
"params": { "member_id": runtime_alias },
})
.to_string(),
Duration::from_secs(5),
None,
Some(&identity_ctx),
)
.await;
let response: Value = serde_json::from_str(&raw)?;
assert!(
response["error"].is_null(),
"identity-owned retire_member failed: {response:#?}"
);
assert_eq!(response["result"]["identity_first"], json!(true));
assert_eq!(
identity_rt.status(&identity).await?.state,
IdentityLifecycleState::Retiring,
"the durable authority, not only the Mob projection, must observe retirement"
);
assert!(
runtime
.mob_handle()
.list_members_including_retiring()
.await
.iter()
.any(|entry| {
crate::member_comms_id::runtime_alias_str(entry.agent_identity.as_str())
== runtime_alias
}),
"a bridge-less identity test proves the raw Mob fallback was not invoked"
);
Ok(())
}
#[tokio::test]
async fn retire_and_respawn_rpcs_succeed_for_idle_member()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let temp_dir = tempfile::tempdir()?;
let runtime = Box::pin(
UnifiedRuntime::builder()
.mob_spec(rpc_test_mob_spec(&temp_dir)?)
.module_config(MobKitConfig {
modules: Vec::new(),
discovery: DiscoverySpec {
namespace: "rpc-idle-lifecycle-test".to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
for member in ["worker-one", "worker-two"] {
runtime
.spawn(SpawnMemberSpec::from_wire(
"worker".to_string(),
member.to_string(),
None,
None,
None,
))
.await?;
}
let send = |id: u64, method: &'static str, params: Value| {
let runtime = &runtime;
async move {
let raw = handle_unified_rpc_json(
runtime,
&json!({
"jsonrpc": "2.0",
"id": id,
"method": method,
"params": params,
})
.to_string(),
Duration::from_secs(10),
None,
None,
)
.await;
serde_json::from_str::<Value>(&raw)
}
};
let response = send(
1,
"mobkit/retire_member",
json!({"member_id": "worker-one"}),
)
.await?;
assert!(
response["error"].is_null(),
"retire_member must succeed for an idle member: {response:#?}"
);
assert_eq!(response["result"]["accepted"], json!(true));
assert!(
!runtime
.mob_handle()
.list_members_including_retiring()
.await
.iter()
.any(|entry| entry.agent_identity.as_str() == "worker-one"),
"retired member must leave the roster"
);
let response = send(
2,
"mobkit/respawn_member",
json!({"member_id": "worker-two"}),
)
.await?;
assert!(
response["error"].is_null(),
"respawn_member must succeed for an idle member: {response:#?}"
);
assert_eq!(response["result"]["accepted"], json!(true));
let members = runtime.mob_handle().list_members_including_retiring().await;
let worker_two = members
.iter()
.find(|entry| entry.agent_identity.as_str() == "worker-two")
.expect("respawned member must remain in the roster");
assert_eq!(
worker_two.status,
meerkat_mob::MobMemberStatus::Active,
"respawned member must be active, not wedged in retiring"
);
Ok(())
}
}