1use std::collections::{BTreeMap, BTreeSet};
4use std::future::Future;
5use std::pin::Pin;
6use std::time::Duration;
7
8use serde::{Deserialize, Serialize};
9use serde_json::Value;
10
11use crate::console_aggregator::is_implicit_delegate_member;
12use crate::mob_handle_runtime::{topology_restore_failed_peer_ids, topology_restore_warning_json};
13use crate::runtime::{
14 BigQuerySessionStoreAdapter, BigQuerySessionStoreError, ConsoleRestJsonRequest,
15 ConsoleRestJsonResponse, DeliveryHistoryRequest, DeliverySendError, DeliverySendRequest,
16 GatingDecideError, GatingDecideRequest, GatingDecision, GatingEvaluateRequest, GatingRiskTier,
17 LocalJsonMemoryStoreError, MemoryIndexError, MemoryIndexRequest, MemoryQueryRequest,
18 MobkitRuntimeHandle, ModuleRouteError, ModuleRouteRequest, ROUTING_RETRY_MAX_CAP,
19 RoutingResolveError, RoutingResolveRequest, RuntimeDecisionState, RuntimeRoute,
20 RuntimeRouteMutationError, ScheduleDefinition, ScheduleValidationError, SessionPersistenceRow,
21 SubscribeError, SubscribeRequest, SubscribeScope, handle_console_rest_json_route,
22 route_module_call, validate_schedules,
23};
24use crate::unified_runtime::{EventQuery, UnifiedRuntime};
25
26mod console_ingress;
27mod gating_methods;
28pub(crate) mod memory_methods;
29pub(crate) mod mob_methods;
30pub(crate) mod params;
31mod routing_delivery_methods;
32mod scheduling_methods;
33mod session_store_methods;
34mod subscribe_methods;
35pub(crate) mod topology_methods;
36pub(crate) mod workgraph_methods;
37
38pub use console_ingress::handle_console_ingress_json;
39
40use gating_methods::{
41 GatingParamsError, parse_gating_audit_params, parse_gating_decide_params,
42 parse_gating_evaluate_params, parse_gating_pending_params,
43};
44use memory_methods::{
45 MemoryParamsError, parse_agent_memory_forget_params, parse_agent_memory_manifest_params,
46 parse_agent_memory_recall_params, parse_agent_memory_remember_params,
47 parse_agent_memory_update_params, parse_memory_index_params, parse_memory_query_params,
48 parse_memory_stores_params,
49};
50use routing_delivery_methods::{
51 RoutingDeliveryParamsError, parse_delivery_history_params, parse_delivery_send_params,
52 parse_routing_resolve_params, parse_routing_route_add_params,
53 parse_routing_route_delete_params, parse_routing_routes_list_params,
54};
55use scheduling_methods::{format_schedule_validation_error, parse_scheduling_params};
56use session_store_methods::{
57 BigQuerySessionStoreRpcError, format_bigquery_store_error, parse_bigquery_session_store_params,
58 run_bigquery_session_store_request,
59};
60use subscribe_methods::{SubscribeParamsError, parse_subscribe_request};
61
62pub const JSONRPC_VERSION: &str = "2.0";
63pub const MOBKIT_CONTRACT_VERSION: &str = "0.4.0";
64pub const MAX_SCHEDULES_PER_REQUEST: usize = 256;
65pub(crate) const MOBPACK_AUTHORING_METHODS: &[&str] = &[
66 "mobkit/mobpacks/schema",
67 "mobkit/mobpacks/catalogs",
68 "mobkit/tools/catalog",
69 "mobkit/skills/catalog",
70 "mobkit/agent_definitions/list",
71 "mobkit/mobpacks/templates",
72 "mobkit/mobpacks/validate",
73 "mobkit/mobpacks/source",
74 "mobkit/mobpacks/export",
75 "mobkit/mobpacks/import",
76 "mobkit/mobpacks/list",
77 "mobkit/mobpacks/get",
78 "mobkit/mobpacks/create",
79 "mobkit/mobpacks/save",
80 "mobkit/mobpacks/delete",
81 "mobkit/mobpacks/undo",
82 "mobkit/mobpacks/redo",
83 "mobkit/mobpacks/apply_operation",
84 "mobkit/mobpacks/graph_projection",
85 "mobkit/mobpacks/graph_to_flow",
86 "mobkit/mobpacks/deploy_command",
87 "mobkit/mobpacks/deploy",
88];
89
90pub(crate) fn mobpack_authoring_capabilities() -> Value {
91 serde_json::json!({
92 "domain": "mobpack_authoring",
93 "runtime_mutation": false,
94 "host_mutation_methods": {
95 "mobkit/mobpacks/deploy": "when execute=true, writes a mobpack archive and runs rkat mob run on the host",
96 "mobkit/mobpacks/validate": "when rkat_validate=true, writes a mobpack archive and runs rkat mob validate on the host"
97 },
98 "deploy_command": "rkat mob run",
99 "methods": MOBPACK_AUTHORING_METHODS,
100 "operations": crate::mobpack::mobpack_authoring_operations(),
101 })
102}
103
104async fn mobpack_runtime_catalog_state(
105 runtime: &UnifiedRuntime,
106) -> crate::mobpack::MobpackRuntimeCatalogState {
107 let loaded_modules = runtime.loaded_modules().await;
108 let runtime_flow_rows = crate::mobpack::runtime_flow_registry_rows_from_definition(
109 runtime.mob_handle().definition(),
110 );
111 let runtime_agent_definition_sources =
112 crate::mobpack::runtime_agent_definition_sources_from_definition(
113 runtime.mob_handle().definition(),
114 );
115 let runtime_skill_realms =
116 crate::mobpack::runtime_skill_realms_from_definition(runtime.mob_handle().definition());
117 let mut runtime_methods = vec![
118 "mobkit/capabilities".to_string(),
119 "mobkit/models/catalog".to_string(),
120 "mobkit/spawn_member".to_string(),
121 "mobkit/list_members".to_string(),
122 "mobkit/get_member".to_string(),
123 "mobkit/run_flow".to_string(),
124 "mobkit/list_flows".to_string(),
125 "mobkit/list_runs".to_string(),
126 ];
127 runtime_methods.extend(
128 MOBPACK_AUTHORING_METHODS
129 .iter()
130 .map(std::string::ToString::to_string),
131 );
132 if runtime.has_contact_directory() {
133 runtime_methods.push("mobkit/cross_mob/directory".to_string());
134 }
135 if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
136 runtime_methods.extend([
137 "mobkit/cross_mob/wire".to_string(),
138 "mobkit/cross_mob/unwire".to_string(),
139 "mobkit/cross_mob/send".to_string(),
140 ]);
141 }
142 crate::mobpack::MobpackRuntimeCatalogState {
143 loaded_modules,
144 runtime_methods,
145 has_contact_directory: runtime.has_contact_directory(),
146 has_peer_mob_handles: runtime.has_peer_mob_handles().await,
147 has_inproc_contacts: runtime.has_inproc_contacts(),
148 runtime_flow_rows,
149 runtime_agent_definition_sources,
150 runtime_skill_realms,
151 }
152}
153
154async fn handle_unified_mobpack_authoring_rpc(
155 runtime: &UnifiedRuntime,
156 method: &str,
157 params: &Value,
158 response_id: Value,
159) -> JsonRpcResponse {
160 let runtime_catalog_state = match method {
161 "mobkit/mobpacks/schema"
162 | "mobkit/mobpacks/catalogs"
163 | "mobkit/tools/catalog"
164 | "mobkit/skills/catalog"
165 | "mobkit/agent_definitions/list"
166 | "mobkit/mobpacks/templates"
167 | "mobkit/mobpacks/list"
168 | "mobkit/mobpacks/get"
169 | "mobkit/mobpacks/apply_operation" => Some(mobpack_runtime_catalog_state(runtime).await),
170 _ => None,
171 };
172 let result = match method {
173 "mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
174 runtime_catalog_state.as_ref(),
175 )),
176 "mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
177 runtime_catalog_state.as_ref(),
178 )),
179 "mobkit/skills/catalog" => Ok(
180 crate::mobpack::mobpack_skills_catalog_response_with_runtime(
181 runtime_catalog_state.as_ref(),
182 ),
183 ),
184 "mobkit/agent_definitions/list" => Ok(
185 crate::mobpack::mobpack_agent_definitions_response_with_runtime(
186 runtime_catalog_state.as_ref(),
187 ),
188 ),
189 "mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
190 runtime_catalog_state.as_ref(),
191 )),
192 _ => {
193 return handle_mobpack_authoring_rpc_with_runtime(
194 method,
195 params,
196 response_id.clone(),
197 runtime_catalog_state.as_ref(),
198 )
199 .unwrap_or_else(|| JsonRpcResponse {
200 jsonrpc: JSONRPC_VERSION.to_string(),
201 id: response_id,
202 result: None,
203 error: Some(JsonRpcError {
204 code: -32601,
205 message: "Method not found".to_string(),
206 data: None,
207 }),
208 });
209 }
210 };
211 match result {
212 Ok(result) => JsonRpcResponse {
213 jsonrpc: JSONRPC_VERSION.to_string(),
214 id: response_id,
215 result: Some(result),
216 error: None,
217 },
218 Err(message) => JsonRpcResponse {
219 jsonrpc: JSONRPC_VERSION.to_string(),
220 id: response_id,
221 result: None,
222 error: Some(JsonRpcError {
223 code: -32602,
224 message,
225 data: None,
226 }),
227 },
228 }
229}
230
231pub(crate) fn handle_mobpack_authoring_rpc(
232 method: &str,
233 params: &Value,
234 response_id: Value,
235) -> Option<JsonRpcResponse> {
236 handle_mobpack_authoring_rpc_with_runtime(method, params, response_id, None)
237}
238
239pub(crate) fn handle_mobpack_authoring_rpc_with_runtime(
240 method: &str,
241 params: &Value,
242 response_id: Value,
243 runtime: Option<&crate::mobpack::MobpackRuntimeCatalogState>,
244) -> Option<JsonRpcResponse> {
245 let result = match method {
246 "mobkit/mobpacks/schema" => Ok(crate::mobpack::mobpack_schema_response_with_runtime(
247 runtime,
248 )),
249 "mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
250 runtime,
251 )),
252 "mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
253 runtime,
254 )),
255 "mobkit/skills/catalog" => {
256 Ok(crate::mobpack::mobpack_skills_catalog_response_with_runtime(runtime))
257 }
258 "mobkit/agent_definitions/list" => {
259 Ok(crate::mobpack::mobpack_agent_definitions_response_with_runtime(runtime))
260 }
261 "mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
262 runtime,
263 )),
264 "mobkit/mobpacks/validate" => crate::mobpack::validate_mobpack(params)
265 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
266 "mobkit/mobpacks/source" => crate::mobpack::source_mobpack(params)
267 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
268 "mobkit/mobpacks/export" => crate::mobpack::export_mobpack(params)
269 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
270 "mobkit/mobpacks/import" => crate::mobpack::import_mobpack(params),
271 "mobkit/mobpacks/list" => crate::mobpack::list_mobpack_drafts_with_runtime(params, runtime),
272 "mobkit/mobpacks/get" => crate::mobpack::get_mobpack_draft_with_runtime(params, runtime),
273 "mobkit/mobpacks/create" => crate::mobpack::create_mobpack_draft(params),
274 "mobkit/mobpacks/save" => crate::mobpack::save_mobpack_draft(params),
275 "mobkit/mobpacks/delete" => crate::mobpack::delete_mobpack_draft(params),
276 "mobkit/mobpacks/undo" => crate::mobpack::undo_mobpack_draft(params),
277 "mobkit/mobpacks/redo" => crate::mobpack::redo_mobpack_draft(params),
278 "mobkit/mobpacks/apply_operation" => {
279 crate::mobpack::apply_mobpack_authoring_operation_with_runtime(params, runtime)
280 }
281 "mobkit/mobpacks/graph_projection" => crate::mobpack::graph_projection_mobpack(params),
282 "mobkit/mobpacks/graph_to_flow" => crate::mobpack::graph_to_flow_mobpack(params),
283 "mobkit/mobpacks/deploy_command" => crate::mobpack::deploy_command_preview(params)
284 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
285 "mobkit/mobpacks/deploy" => crate::mobpack::deploy_mobpack(params)
286 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
287 _ => return None,
288 };
289 Some(match result {
290 Ok(result) => JsonRpcResponse {
291 jsonrpc: JSONRPC_VERSION.to_string(),
292 id: response_id,
293 result: Some(result),
294 error: None,
295 },
296 Err(message) => JsonRpcResponse {
297 jsonrpc: JSONRPC_VERSION.to_string(),
298 id: response_id,
299 result: None,
300 error: Some(JsonRpcError {
301 code: -32602,
302 message,
303 data: None,
304 }),
305 },
306 })
307}
308
309#[derive(Debug, Clone, PartialEq, Eq)]
310pub enum RpcCapabilitiesError {
311 InvalidJson,
312 InvalidSchema,
313 MissingContractVersion,
314 InvalidContractVersion,
315}
316
317impl std::fmt::Display for RpcCapabilitiesError {
318 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
319 match self {
320 Self::InvalidJson => write!(f, "invalid JSON"),
321 Self::InvalidSchema => write!(f, "invalid schema"),
322 Self::MissingContractVersion => write!(f, "missing contract version"),
323 Self::InvalidContractVersion => write!(f, "invalid contract version"),
324 }
325 }
326}
327
328impl std::error::Error for RpcCapabilitiesError {}
329
330#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
331pub struct RpcCapabilities {
332 pub contract_version: String,
333 #[serde(flatten)]
334 pub extra: BTreeMap<String, Value>,
335}
336
337pub fn parse_rpc_capabilities(line: &str) -> Result<RpcCapabilities, RpcCapabilitiesError> {
338 let raw: Value = serde_json::from_str(line).map_err(|_| RpcCapabilitiesError::InvalidJson)?;
339 let object = raw.as_object().ok_or(RpcCapabilitiesError::InvalidSchema)?;
340 let contract = object
341 .get("contract_version")
342 .ok_or(RpcCapabilitiesError::MissingContractVersion)?;
343 let contract_str = contract
344 .as_str()
345 .ok_or(RpcCapabilitiesError::InvalidContractVersion)?;
346 if contract_str.trim().is_empty() {
347 return Err(RpcCapabilitiesError::InvalidContractVersion);
348 }
349 serde_json::from_value(raw).map_err(|_| RpcCapabilitiesError::InvalidSchema)
350}
351
352pub const MOB_EVENTS_STALE_CURSOR_CODE: i64 = -32010;
359
360pub const MEMORY_BACKEND_UNAVAILABLE_CODE: i64 = -32012;
367
368pub const CONSOLE_TIMELINE_REPLAY_UNAVAILABLE_CODE: i64 = -32013;
373
374pub const WORKGRAPH_UNAVAILABLE_CODE: i64 = -32041;
379
380pub const WORKGRAPH_CONFLICT_CODE: i64 = -32042;
385
386pub const WORKGRAPH_ERROR_CODE: i64 = -32000;
389
390#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
391pub struct JsonRpcRequest {
392 pub jsonrpc: String,
393 #[serde(default)]
394 pub id: Option<Value>,
395 pub method: String,
396 #[serde(default)]
397 pub params: Value,
398}
399
400#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
401pub struct JsonRpcError {
402 pub code: i64,
403 pub message: String,
404 #[serde(default, skip_serializing_if = "Option::is_none")]
409 pub data: Option<Value>,
410}
411
412impl JsonRpcError {
413 pub fn new(code: i64, message: impl Into<String>) -> Self {
414 Self {
415 code,
416 message: message.into(),
417 data: None,
418 }
419 }
420
421 #[must_use]
422 pub fn with_data(mut self, data: Value) -> Self {
423 self.data = Some(data);
424 self
425 }
426}
427
428#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
429pub struct JsonRpcResponse {
430 pub jsonrpc: String,
431 pub id: Value,
432 #[serde(skip_serializing_if = "Option::is_none")]
433 pub result: Option<Value>,
434 #[serde(skip_serializing_if = "Option::is_none")]
435 pub error: Option<JsonRpcError>,
436}
437
438pub fn handle_mobkit_rpc_json(
439 runtime: &mut MobkitRuntimeHandle,
440 request_json: &str,
441 timeout: Duration,
442) -> String {
443 let raw_request: Value = match serde_json::from_str(request_json) {
444 Ok(raw_request) => raw_request,
445 Err(_) => {
446 return serialize_response(&JsonRpcResponse {
447 jsonrpc: JSONRPC_VERSION.to_string(),
448 id: Value::Null,
449 result: None,
450 error: Some(JsonRpcError {
451 code: -32700,
452 message: "Parse error".to_string(),
453 data: None,
454 }),
455 });
456 }
457 };
458 let response_id = raw_request
459 .as_object()
460 .and_then(|object| object.get("id"))
461 .cloned()
462 .unwrap_or(Value::Null);
463 let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
464 Ok(request) => request,
465 Err(_) => {
466 return serialize_response(&JsonRpcResponse {
467 jsonrpc: JSONRPC_VERSION.to_string(),
468 id: response_id,
469 result: None,
470 error: Some(JsonRpcError {
471 code: -32600,
472 message: "Invalid Request".to_string(),
473 data: None,
474 }),
475 });
476 }
477 };
478 let is_notification = request.id.is_none();
479 let response_id = request.id.clone().unwrap_or(Value::Null);
480
481 if request.jsonrpc != "2.0" {
482 let response = JsonRpcResponse {
483 jsonrpc: JSONRPC_VERSION.to_string(),
484 id: response_id,
485 result: None,
486 error: Some(JsonRpcError {
487 code: -32600,
488 message: "Invalid Request".to_string(),
489 data: None,
490 }),
491 };
492 return if is_notification {
493 String::new()
494 } else {
495 serialize_response(&response)
496 };
497 }
498
499 let response = match request.method.as_str() {
500 "mobkit/status" => JsonRpcResponse {
501 jsonrpc: JSONRPC_VERSION.to_string(),
502 id: response_id,
503 result: Some(serde_json::json!({
504 "contract_version": MOBKIT_CONTRACT_VERSION,
505 "running": runtime.is_running(),
506 "loaded_modules": runtime.loaded_modules(),
507 })),
508 error: None,
509 },
510 "mobkit/capabilities" => {
511 let mut methods = vec![
512 "mobkit/status",
513 "mobkit/capabilities",
514 "mobkit/reconcile",
515 "mobkit/spawn_member",
516 "mobkit/scheduling/evaluate",
517 "mobkit/scheduling/dispatch",
518 "mobkit/routing/resolve",
519 "mobkit/routing/routes/list",
520 "mobkit/routing/routes/add",
521 "mobkit/routing/routes/delete",
522 "mobkit/delivery/send",
523 "mobkit/delivery/history",
524 "mobkit/events/subscribe",
525 "mobkit/memory/stores",
526 "mobkit/memory/index",
527 "mobkit/memory/query",
528 "mobkit/session_store/bigquery",
529 "mobkit/gating/evaluate",
530 "mobkit/gating/pending",
531 "mobkit/gating/decide",
532 "mobkit/gating/audit",
533 "mobkit/call_tool",
534 "mobkit/models/catalog",
535 ];
536 methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
537 JsonRpcResponse {
538 jsonrpc: JSONRPC_VERSION.to_string(),
539 id: response_id,
540 result: Some(serde_json::json!({
541 "contract_version": MOBKIT_CONTRACT_VERSION,
542 "methods": methods,
543 "loaded_modules": runtime.loaded_modules(),
544 "runtime_capabilities": {
545 "can_spawn_members": false,
546 "can_send_messages": false,
547 "can_wire_members": false,
548 "can_retire_members": false,
549 "available_spawn_modes": ["module"],
550 },
551 "authoring_capabilities": mobpack_authoring_capabilities(),
552 })),
553 error: None,
554 }
555 }
556 "mobkit/models/catalog" => JsonRpcResponse {
557 jsonrpc: JSONRPC_VERSION.to_string(),
558 id: response_id,
559 result: Some(build_models_catalog_result()),
560 error: None,
561 },
562 method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
563 handle_mobpack_authoring_rpc(method, &request.params, response_id.clone())
564 .unwrap_or_else(|| JsonRpcResponse {
565 jsonrpc: JSONRPC_VERSION.to_string(),
566 id: response_id,
567 result: None,
568 error: Some(JsonRpcError {
569 code: -32601,
570 message: "Method not found".to_string(),
571 data: None,
572 }),
573 })
574 }
575 "mobkit/reconcile" => {
576 let modules = match params::required_string_array(&request.params, "modules") {
577 Ok(m) => m,
578 Err(reason) => {
579 return serialize_response(&JsonRpcResponse {
580 jsonrpc: JSONRPC_VERSION.to_string(),
581 id: response_id,
582 result: None,
583 error: Some(JsonRpcError {
584 code: -32602,
585 message: format!("Invalid params: {reason}"),
586 data: None,
587 }),
588 });
589 }
590 };
591
592 match runtime.reconcile_modules(modules.clone(), timeout) {
593 Ok(added) => JsonRpcResponse {
594 jsonrpc: JSONRPC_VERSION.to_string(),
595 id: response_id,
596 result: Some(serde_json::json!({
597 "accepted": true,
598 "reconciled_modules": modules,
599 "added": added
600 })),
601 error: None,
602 },
603 Err(err) => JsonRpcResponse {
604 jsonrpc: JSONRPC_VERSION.to_string(),
605 id: response_id,
606 result: None,
607 error: Some(JsonRpcError {
608 code: -32602,
609 message: format!("Invalid params: {err:?}"),
610 data: None,
611 }),
612 },
613 }
614 }
615 "mobkit/spawn_member" => {
616 let module_id = request
617 .params
618 .get("module_id")
619 .and_then(Value::as_str)
620 .unwrap_or_default()
621 .to_string();
622 if module_id.is_empty() {
623 JsonRpcResponse {
624 jsonrpc: JSONRPC_VERSION.to_string(),
625 id: response_id,
626 result: None,
627 error: Some(JsonRpcError {
628 code: -32602,
629 message: "Invalid params: module_id required".to_string(),
630 data: None,
631 }),
632 }
633 } else {
634 match runtime.spawn_member(&module_id, timeout) {
635 Ok(()) => JsonRpcResponse {
636 jsonrpc: JSONRPC_VERSION.to_string(),
637 id: response_id,
638 result: Some(serde_json::json!({
639 "accepted": true,
640 "module_id": module_id
641 })),
642 error: None,
643 },
644 Err(err) => JsonRpcResponse {
645 jsonrpc: JSONRPC_VERSION.to_string(),
646 id: response_id,
647 result: None,
648 error: Some(JsonRpcError {
649 code: -32602,
650 message: format!("Invalid params: {err:?}"),
651 data: None,
652 }),
653 },
654 }
655 }
656 }
657 "mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
658 Ok((schedules, tick_ms)) => match runtime.evaluate_schedule_tick(&schedules, tick_ms) {
659 Ok(evaluation) => JsonRpcResponse {
660 jsonrpc: JSONRPC_VERSION.to_string(),
661 id: response_id,
662 result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
663 error: None,
664 },
665 Err(err) => JsonRpcResponse {
666 jsonrpc: JSONRPC_VERSION.to_string(),
667 id: response_id,
668 result: None,
669 error: Some(JsonRpcError {
670 code: -32602,
671 message: format!(
672 "Invalid params: {}",
673 format_schedule_validation_error(err)
674 ),
675 data: None,
676 }),
677 },
678 },
679 Err(message) => JsonRpcResponse {
680 jsonrpc: JSONRPC_VERSION.to_string(),
681 id: response_id,
682 result: None,
683 error: Some(JsonRpcError {
684 code: -32602,
685 message: format!("Invalid params: {message}"),
686 data: None,
687 }),
688 },
689 },
690 "mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
691 Ok((schedules, tick_ms)) => match runtime.dispatch_schedule_tick(&schedules, tick_ms) {
692 Ok(dispatch) => JsonRpcResponse {
693 jsonrpc: JSONRPC_VERSION.to_string(),
694 id: response_id,
695 result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
696 error: None,
697 },
698 Err(err) => JsonRpcResponse {
699 jsonrpc: JSONRPC_VERSION.to_string(),
700 id: response_id,
701 result: None,
702 error: Some(JsonRpcError {
703 code: -32602,
704 message: format!(
705 "Invalid params: {}",
706 format_schedule_validation_error(err)
707 ),
708 data: None,
709 }),
710 },
711 },
712 Err(message) => JsonRpcResponse {
713 jsonrpc: JSONRPC_VERSION.to_string(),
714 id: response_id,
715 result: None,
716 error: Some(JsonRpcError {
717 code: -32602,
718 message: format!("Invalid params: {message}"),
719 data: None,
720 }),
721 },
722 },
723 "mobkit/routing/resolve" => {
724 match parse_routing_resolve_params(&request.params).and_then(|resolve_request| {
725 runtime
726 .resolve_routing(resolve_request)
727 .map_err(RoutingDeliveryParamsError::Routing)
728 }) {
729 Ok(resolution) => JsonRpcResponse {
730 jsonrpc: JSONRPC_VERSION.to_string(),
731 id: response_id,
732 result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
733 error: None,
734 },
735 Err(err) => JsonRpcResponse {
736 jsonrpc: JSONRPC_VERSION.to_string(),
737 id: response_id,
738 result: None,
739 error: Some(JsonRpcError {
740 code: -32602,
741 message: format!("Invalid params: {}", err.message()),
742 data: None,
743 }),
744 },
745 }
746 }
747 "mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
748 Ok(()) => JsonRpcResponse {
749 jsonrpc: JSONRPC_VERSION.to_string(),
750 id: response_id,
751 result: Some(serde_json::json!({
752 "routes": runtime.list_runtime_routes()
753 })),
754 error: None,
755 },
756 Err(err) => JsonRpcResponse {
757 jsonrpc: JSONRPC_VERSION.to_string(),
758 id: response_id,
759 result: None,
760 error: Some(JsonRpcError {
761 code: -32602,
762 message: format!("Invalid params: {}", err.message()),
763 data: None,
764 }),
765 },
766 },
767 "mobkit/routing/routes/add" => match parse_routing_route_add_params(&request.params)
768 .and_then(|route| {
769 runtime
770 .add_runtime_route(route)
771 .map_err(RoutingDeliveryParamsError::RouteMutation)
772 }) {
773 Ok(route) => JsonRpcResponse {
774 jsonrpc: JSONRPC_VERSION.to_string(),
775 id: response_id,
776 result: Some(serde_json::json!({ "route": route })),
777 error: None,
778 },
779 Err(err) => JsonRpcResponse {
780 jsonrpc: JSONRPC_VERSION.to_string(),
781 id: response_id,
782 result: None,
783 error: Some(JsonRpcError {
784 code: -32602,
785 message: format!("Invalid params: {}", err.message()),
786 data: None,
787 }),
788 },
789 },
790 "mobkit/routing/routes/delete" => match parse_routing_route_delete_params(&request.params)
791 .and_then(|route_key| {
792 runtime
793 .delete_runtime_route(&route_key)
794 .map_err(RoutingDeliveryParamsError::RouteMutation)
795 }) {
796 Ok(route) => JsonRpcResponse {
797 jsonrpc: JSONRPC_VERSION.to_string(),
798 id: response_id,
799 result: Some(serde_json::json!({ "deleted": route })),
800 error: None,
801 },
802 Err(err) => JsonRpcResponse {
803 jsonrpc: JSONRPC_VERSION.to_string(),
804 id: response_id,
805 result: None,
806 error: Some(JsonRpcError {
807 code: -32602,
808 message: format!("Invalid params: {}", err.message()),
809 data: None,
810 }),
811 },
812 },
813 "mobkit/delivery/send" => {
814 match parse_delivery_send_params(&request.params).and_then(|send_request| {
815 runtime
816 .send_delivery(send_request)
817 .map_err(RoutingDeliveryParamsError::Delivery)
818 }) {
819 Ok(record) => JsonRpcResponse {
820 jsonrpc: JSONRPC_VERSION.to_string(),
821 id: response_id,
822 result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
823 error: None,
824 },
825 Err(err) => JsonRpcResponse {
826 jsonrpc: JSONRPC_VERSION.to_string(),
827 id: response_id,
828 result: None,
829 error: Some(JsonRpcError {
830 code: -32602,
831 message: format!("Invalid params: {}", err.message()),
832 data: None,
833 }),
834 },
835 }
836 }
837 "mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
838 Ok(history_request) => JsonRpcResponse {
839 jsonrpc: JSONRPC_VERSION.to_string(),
840 id: response_id,
841 result: Some(
842 serde_json::to_value(runtime.delivery_history(history_request))
843 .unwrap_or(Value::Null),
844 ),
845 error: None,
846 },
847 Err(err) => JsonRpcResponse {
848 jsonrpc: JSONRPC_VERSION.to_string(),
849 id: response_id,
850 result: None,
851 error: Some(JsonRpcError {
852 code: -32602,
853 message: format!("Invalid params: {}", err.message()),
854 data: None,
855 }),
856 },
857 },
858 "mobkit/events/subscribe" => {
859 match parse_subscribe_request(&request.params).and_then(|subscribe_request| {
860 runtime
861 .subscribe_events(subscribe_request)
862 .map_err(SubscribeParamsError::Runtime)
863 }) {
864 Ok(subscribe_result) => JsonRpcResponse {
865 jsonrpc: JSONRPC_VERSION.to_string(),
866 id: response_id,
867 result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
868 error: None,
869 },
870 Err(err) => JsonRpcResponse {
871 jsonrpc: JSONRPC_VERSION.to_string(),
872 id: response_id,
873 result: None,
874 error: Some(JsonRpcError {
875 code: -32602,
876 message: format!("Invalid params: {}", err.message()),
877 data: None,
878 }),
879 },
880 }
881 }
882 "mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
883 Ok(()) => JsonRpcResponse {
884 jsonrpc: JSONRPC_VERSION.to_string(),
885 id: response_id,
886 result: Some(serde_json::json!({
887 "stores": runtime.memory_stores(),
888 })),
889 error: None,
890 },
891 Err(err) => JsonRpcResponse {
892 jsonrpc: JSONRPC_VERSION.to_string(),
893 id: response_id,
894 result: None,
895 error: Some(JsonRpcError {
896 code: -32602,
897 message: format!("Invalid params: {}", err.message()),
898 data: None,
899 }),
900 },
901 },
902 "mobkit/memory/index" => match parse_memory_index_params(&request.params) {
903 Ok(index_request) => match runtime.memory_index(index_request) {
904 Ok(indexed) => JsonRpcResponse {
905 jsonrpc: JSONRPC_VERSION.to_string(),
906 id: response_id,
907 result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
908 error: None,
909 },
910 Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
911 jsonrpc: JSONRPC_VERSION.to_string(),
912 id: response_id,
913 result: None,
914 error: Some(JsonRpcError {
915 code: MEMORY_BACKEND_UNAVAILABLE_CODE,
916 message: format!(
917 "Memory backend unavailable: {}",
918 MemoryParamsError::backend_message(&error)
919 ),
920 data: None,
921 }),
922 },
923 Err(err) => JsonRpcResponse {
924 jsonrpc: JSONRPC_VERSION.to_string(),
925 id: response_id,
926 result: None,
927 error: Some(JsonRpcError {
928 code: -32602,
929 message: format!(
930 "Invalid params: {}",
931 MemoryParamsError::Index(err).message()
932 ),
933 data: None,
934 }),
935 },
936 },
937 Err(err) => JsonRpcResponse {
938 jsonrpc: JSONRPC_VERSION.to_string(),
939 id: response_id,
940 result: None,
941 error: Some(JsonRpcError {
942 code: -32602,
943 message: format!("Invalid params: {}", err.message()),
944 data: None,
945 }),
946 },
947 },
948 "mobkit/memory/query" => match parse_memory_query_params(&request.params) {
949 Ok(query_request) => JsonRpcResponse {
950 jsonrpc: JSONRPC_VERSION.to_string(),
951 id: response_id,
952 result: Some(
953 serde_json::to_value(runtime.memory_query(query_request))
954 .unwrap_or(Value::Null),
955 ),
956 error: None,
957 },
958 Err(err) => JsonRpcResponse {
959 jsonrpc: JSONRPC_VERSION.to_string(),
960 id: response_id,
961 result: None,
962 error: Some(JsonRpcError {
963 code: -32602,
964 message: format!("Invalid params: {}", err.message()),
965 data: None,
966 }),
967 },
968 },
969 "mobkit/session_store/bigquery" => {
970 match parse_bigquery_session_store_params(&request.params)
971 .and_then(run_bigquery_session_store_request)
972 {
973 Ok(result) => JsonRpcResponse {
974 jsonrpc: JSONRPC_VERSION.to_string(),
975 id: response_id,
976 result: Some(result),
977 error: None,
978 },
979 Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
980 jsonrpc: JSONRPC_VERSION.to_string(),
981 id: response_id,
982 result: None,
983 error: Some(JsonRpcError {
984 code: -32602,
985 message: format!("Invalid params: {message}"),
986 data: None,
987 }),
988 },
989 Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
990 jsonrpc: JSONRPC_VERSION.to_string(),
991 id: response_id,
992 result: None,
993 error: Some(JsonRpcError {
994 code: -32011,
995 message: format!(
996 "BigQuery session store request failed: {}",
997 format_bigquery_store_error(&error)
998 ),
999 data: None,
1000 }),
1001 },
1002 }
1003 }
1004 "mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
1005 Ok(gating_request) => JsonRpcResponse {
1006 jsonrpc: JSONRPC_VERSION.to_string(),
1007 id: response_id,
1008 result: Some(
1009 serde_json::to_value(runtime.evaluate_gating_action(gating_request))
1010 .unwrap_or(Value::Null),
1011 ),
1012 error: None,
1013 },
1014 Err(err) => JsonRpcResponse {
1015 jsonrpc: JSONRPC_VERSION.to_string(),
1016 id: response_id,
1017 result: None,
1018 error: Some(JsonRpcError {
1019 code: -32602,
1020 message: format!("Invalid params: {}", err.message()),
1021 data: None,
1022 }),
1023 },
1024 },
1025 "mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
1026 Ok(()) => JsonRpcResponse {
1027 jsonrpc: JSONRPC_VERSION.to_string(),
1028 id: response_id,
1029 result: Some(serde_json::json!({
1030 "pending": runtime.list_gating_pending(),
1031 })),
1032 error: None,
1033 },
1034 Err(err) => JsonRpcResponse {
1035 jsonrpc: JSONRPC_VERSION.to_string(),
1036 id: response_id,
1037 result: None,
1038 error: Some(JsonRpcError {
1039 code: -32602,
1040 message: format!("Invalid params: {}", err.message()),
1041 data: None,
1042 }),
1043 },
1044 },
1045 "mobkit/gating/decide" => {
1046 match parse_gating_decide_params(&request.params).and_then(|decide_request| {
1047 runtime
1048 .decide_gating_action(decide_request)
1049 .map_err(GatingParamsError::Decision)
1050 }) {
1051 Ok(result) => JsonRpcResponse {
1052 jsonrpc: JSONRPC_VERSION.to_string(),
1053 id: response_id,
1054 result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
1055 error: None,
1056 },
1057 Err(err) => JsonRpcResponse {
1058 jsonrpc: JSONRPC_VERSION.to_string(),
1059 id: response_id,
1060 result: None,
1061 error: Some(JsonRpcError {
1062 code: -32602,
1063 message: format!("Invalid params: {}", err.message()),
1064 data: None,
1065 }),
1066 },
1067 }
1068 }
1069 "mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
1070 Ok(limit) => JsonRpcResponse {
1071 jsonrpc: JSONRPC_VERSION.to_string(),
1072 id: response_id,
1073 result: Some(serde_json::json!({
1074 "entries": runtime.gating_audit_entries(limit),
1075 })),
1076 error: None,
1077 },
1078 Err(err) => JsonRpcResponse {
1079 jsonrpc: JSONRPC_VERSION.to_string(),
1080 id: response_id,
1081 result: None,
1082 error: Some(JsonRpcError {
1083 code: -32602,
1084 message: format!("Invalid params: {}", err.message()),
1085 data: None,
1086 }),
1087 },
1088 },
1089 "mobkit/call_tool" => {
1090 let module_id = request.params.get("module_id").and_then(Value::as_str);
1091 let tool = request.params.get("tool").and_then(Value::as_str);
1092 let arguments = request
1093 .params
1094 .get("arguments")
1095 .cloned()
1096 .unwrap_or(serde_json::json!({}));
1097
1098 match (module_id, tool) {
1099 (Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
1100 let route = route_module_call(
1101 runtime,
1102 &ModuleRouteRequest {
1103 module_id: module_id.to_string(),
1104 method: tool.to_string(),
1105 params: arguments,
1106 },
1107 timeout,
1108 );
1109 match route {
1110 Ok(response) => JsonRpcResponse {
1111 jsonrpc: JSONRPC_VERSION.to_string(),
1112 id: response_id,
1113 result: Some(serde_json::json!({
1114 "module_id": response.module_id,
1115 "tool": response.method,
1116 "result": response.payload
1117 })),
1118 error: None,
1119 },
1120 Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
1121 jsonrpc: JSONRPC_VERSION.to_string(),
1122 id: response_id,
1123 result: None,
1124 error: Some(JsonRpcError {
1125 code: -32601,
1126 message: format!("Module '{mid}' not loaded"),
1127 data: None,
1128 }),
1129 },
1130 Err(err) => JsonRpcResponse {
1131 jsonrpc: JSONRPC_VERSION.to_string(),
1132 id: response_id,
1133 result: None,
1134 error: Some(JsonRpcError {
1135 code: -32000,
1136 message: format!("Tool call failed: {err:?}"),
1137 data: None,
1138 }),
1139 },
1140 }
1141 }
1142 _ => JsonRpcResponse {
1143 jsonrpc: JSONRPC_VERSION.to_string(),
1144 id: response_id,
1145 result: None,
1146 error: Some(JsonRpcError {
1147 code: -32602,
1148 message: "Invalid params: module_id and tool required".to_string(),
1149 data: None,
1150 }),
1151 },
1152 }
1153 }
1154 method if method.contains('/') && !method.starts_with("mobkit/") => {
1155 let module_id = method
1156 .split('/')
1157 .next()
1158 .map(ToString::to_string)
1159 .unwrap_or_default();
1160 let route = route_module_call(
1161 runtime,
1162 &ModuleRouteRequest {
1163 module_id,
1164 method: method.to_string(),
1165 params: request.params,
1166 },
1167 timeout,
1168 );
1169 match route {
1170 Ok(response) => JsonRpcResponse {
1171 jsonrpc: JSONRPC_VERSION.to_string(),
1172 id: response_id,
1173 result: Some(serde_json::json!({
1174 "module_id": response.module_id,
1175 "method": response.method,
1176 "payload": response.payload
1177 })),
1178 error: None,
1179 },
1180 Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
1181 jsonrpc: JSONRPC_VERSION.to_string(),
1182 id: response_id,
1183 result: None,
1184 error: Some(JsonRpcError {
1185 code: -32601,
1186 message: format!("Module '{module_id}' not loaded"),
1187 data: None,
1188 }),
1189 },
1190 Err(err) => JsonRpcResponse {
1191 jsonrpc: JSONRPC_VERSION.to_string(),
1192 id: response_id,
1193 result: None,
1194 error: Some(JsonRpcError {
1195 code: -32000,
1196 message: format!("Module route failed: {err:?}"),
1197 data: None,
1198 }),
1199 },
1200 }
1201 }
1202 _ => JsonRpcResponse {
1203 jsonrpc: JSONRPC_VERSION.to_string(),
1204 id: response_id,
1205 result: None,
1206 error: Some(JsonRpcError {
1207 code: -32601,
1208 message: "Method not found".to_string(),
1209 data: None,
1210 }),
1211 },
1212 };
1213 if is_notification {
1214 String::new()
1215 } else {
1216 serialize_response(&response)
1217 }
1218}
1219
1220pub struct IdentityFirstContext {
1222 pub runtime: std::sync::Arc<crate::identity_first::IdentityRuntime>,
1223 pub roster_provider: std::sync::Arc<dyn crate::identity_first::contracts::RosterProvider>,
1224 pub topology_provider:
1225 Option<std::sync::Arc<dyn crate::identity_first::contracts::TopologyProvider>>,
1226 pub customizer: Option<std::sync::Arc<dyn crate::identity_first::contracts::AgentCustomizer>>,
1227 pub agent_memory_provider:
1228 Option<std::sync::Arc<dyn crate::identity_first::AgentMemoryProvider>>,
1229 pub mob_definition: Option<meerkat_mob::MobDefinition>,
1230}
1231
1232pub fn handle_unified_rpc_json<'a>(
1233 runtime: &'a UnifiedRuntime,
1234 request_json: &'a str,
1235 timeout: Duration,
1236 http_base_url: Option<&'a str>,
1237 identity_ctx: Option<&'a IdentityFirstContext>,
1238) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
1239 handle_unified_rpc_json_with_live(
1240 runtime,
1241 request_json,
1242 timeout,
1243 http_base_url,
1244 identity_ctx,
1245 None,
1246 )
1247}
1248
1249pub fn handle_unified_rpc_json_with_live<'a>(
1254 runtime: &'a UnifiedRuntime,
1255 request_json: &'a str,
1256 timeout: Duration,
1257 http_base_url: Option<&'a str>,
1258 identity_ctx: Option<&'a IdentityFirstContext>,
1259 live: Option<&'a crate::live_wiring::LiveRpcHandler>,
1260) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
1261 Box::pin(handle_unified_rpc_json_inner(
1262 runtime,
1263 request_json,
1264 timeout,
1265 http_base_url,
1266 identity_ctx,
1267 live,
1268 ))
1269}
1270
1271async fn handle_unified_rpc_json_inner(
1272 runtime: &UnifiedRuntime,
1273 request_json: &str,
1274 timeout: Duration,
1275 http_base_url: Option<&str>,
1276 identity_ctx: Option<&IdentityFirstContext>,
1277 live: Option<&crate::live_wiring::LiveRpcHandler>,
1278) -> String {
1279 let raw_request: Value = match serde_json::from_str(request_json) {
1280 Ok(raw_request) => raw_request,
1281 Err(_) => {
1282 return serialize_response(&JsonRpcResponse {
1283 jsonrpc: JSONRPC_VERSION.to_string(),
1284 id: Value::Null,
1285 result: None,
1286 error: Some(JsonRpcError {
1287 code: -32700,
1288 message: "Parse error".to_string(),
1289 data: None,
1290 }),
1291 });
1292 }
1293 };
1294 let response_id = raw_request
1295 .as_object()
1296 .and_then(|object| object.get("id"))
1297 .cloned()
1298 .unwrap_or(Value::Null);
1299 let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
1300 Ok(request) => request,
1301 Err(_) => {
1302 return serialize_response(&JsonRpcResponse {
1303 jsonrpc: JSONRPC_VERSION.to_string(),
1304 id: response_id,
1305 result: None,
1306 error: Some(JsonRpcError {
1307 code: -32600,
1308 message: "Invalid Request".to_string(),
1309 data: None,
1310 }),
1311 });
1312 }
1313 };
1314 let is_notification = request.id.is_none();
1315 let response_id = request.id.clone().unwrap_or(Value::Null);
1316
1317 if request.jsonrpc != "2.0" {
1318 let response = JsonRpcResponse {
1319 jsonrpc: JSONRPC_VERSION.to_string(),
1320 id: response_id,
1321 result: None,
1322 error: Some(JsonRpcError {
1323 code: -32600,
1324 message: "Invalid Request".to_string(),
1325 data: None,
1326 }),
1327 };
1328 return if is_notification {
1329 String::new()
1330 } else {
1331 serialize_response(&response)
1332 };
1333 }
1334
1335 let response = match request.method.as_str() {
1336 "mobkit/status" => {
1337 let mob_state = Some(runtime.mob_handle().status_observation_snapshot());
1338 let is_running = runtime.module_is_running().await;
1339 let loaded = runtime.loaded_modules().await;
1340 let mut result = serde_json::json!({
1341 "contract_version": MOBKIT_CONTRACT_VERSION,
1342 "running": is_running,
1343 "loaded_modules": loaded,
1344 "mob_state": format!("{mob_state:?}"),
1345 });
1346 if let Some(url) = http_base_url {
1347 result["http_base_url"] = Value::String(url.to_string());
1348 }
1349 JsonRpcResponse {
1350 jsonrpc: JSONRPC_VERSION.to_string(),
1351 id: response_id,
1352 result: Some(result),
1353 error: None,
1354 }
1355 }
1356 "mobkit/capabilities" => {
1357 let loaded = runtime.loaded_modules().await;
1358 let mut methods = vec![
1359 "mobkit/init",
1360 "mobkit/status",
1361 "mobkit/capabilities",
1362 "mobkit/reconcile",
1363 "mobkit/spawn_member",
1364 "mobkit/scheduling/evaluate",
1365 "mobkit/scheduling/dispatch",
1366 "mobkit/routing/resolve",
1367 "mobkit/routing/routes/list",
1368 "mobkit/routing/routes/add",
1369 "mobkit/routing/routes/delete",
1370 "mobkit/delivery/send",
1371 "mobkit/delivery/history",
1372 "mobkit/events/subscribe",
1373 "mobkit/query_events",
1374 "mobkit/memory/stores",
1375 "mobkit/memory/index",
1376 "mobkit/memory/query",
1377 "mobkit/session_store/bigquery",
1378 "mobkit/gating/evaluate",
1379 "mobkit/gating/pending",
1380 "mobkit/gating/decide",
1381 "mobkit/gating/audit",
1382 "mobkit/call_tool",
1383 "mobkit/models/catalog",
1384 "mobkit/blob/get",
1385 "mobkit/send_message",
1386 "mobkit/find_members",
1387 "mobkit/ensure_member",
1388 "mobkit/list_members",
1389 "mobkit/get_member",
1390 "mobkit/retire_member",
1391 "mobkit/respawn_member",
1392 "mobkit/reconcile_edges",
1393 "mobkit/rediscover",
1394 "mobkit/mob_events/query",
1395 "mobkit/mob_events/subscribe",
1396 "mobkit/cross_mob/peer_info",
1398 "mobkit/cross_mob/wire_local",
1399 "mobkit/cross_mob/unwire_local",
1400 "mobkit/peer_pubkey",
1401 "mobkit/member_status",
1402 "mobkit/identity/resolved_tools",
1403 "mobkit/force_cancel_member",
1404 "mobkit/spawn_helper",
1405 "mobkit/fork_helper",
1406 "mobkit/attach_existing_session",
1407 "mobkit/cancel_flow",
1408 "mobkit/flow_status",
1409 "mobkit/list_flows",
1410 "mobkit/list_runs",
1411 "mobkit/run_flow",
1412 "mobkit/collect_completed",
1413 "mobkit/wait_ready",
1414 "mobkit/mob_labels/set",
1415 "mobkit/mob_labels/get",
1416 "mobkit/mob_labels/delete",
1417 "mobkit/run_labels/set",
1418 "mobkit/run_labels/get",
1419 "mobkit/run_labels/delete",
1420 ];
1421 methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
1422 let workgraph_configured = runtime.workgraph_service().is_some();
1426 if workgraph_configured {
1427 methods.extend_from_slice(workgraph_methods::WORKGRAPH_READ_METHODS);
1428 methods.extend_from_slice(workgraph_methods::WORKGRAPH_MUTATE_METHODS);
1429 }
1430 if live.is_some() {
1433 methods.extend_from_slice(&[
1434 "mobkit/live/open",
1435 "mobkit/live/status",
1436 "mobkit/live/close",
1437 "mobkit/live/refresh",
1438 "mobkit/live/send_input",
1439 "mobkit/live/commit_input",
1440 "mobkit/live/interrupt",
1441 "mobkit/live/truncate",
1442 ]);
1443 }
1444 if identity_ctx.is_some() {
1445 methods.extend_from_slice(&[
1446 "mobkit/send",
1447 "mobkit/interact",
1448 "mobkit/dispatch",
1449 "mobkit/subscribe",
1450 "mobkit/status_identity",
1451 "mobkit/respawn",
1452 "mobkit/retire",
1453 "mobkit/reset",
1454 "mobkit/delete_identity",
1455 "mobkit/inspect_identity",
1456 "mobkit/reconcile_identity",
1457 ]);
1458 }
1459 if identity_ctx
1460 .and_then(|ctx| ctx.agent_memory_provider.as_ref())
1461 .is_some()
1462 {
1463 methods.push("mobkit/agent_memory/recall");
1464 if identity_ctx
1465 .and_then(|ctx| ctx.agent_memory_provider.as_ref())
1466 .is_some_and(|provider| provider.supports_remember())
1467 {
1468 methods.push("mobkit/agent_memory/remember");
1469 }
1470 if identity_ctx
1471 .and_then(|ctx| ctx.agent_memory_provider.as_ref())
1472 .is_some_and(|provider| provider.supports_forget())
1473 {
1474 methods.push("mobkit/agent_memory/forget");
1475 }
1476 if identity_ctx
1477 .and_then(|ctx| ctx.agent_memory_provider.as_ref())
1478 .is_some_and(|provider| provider.supports_supersede())
1479 {
1480 methods.push("mobkit/agent_memory/update");
1481 }
1482 if identity_ctx
1483 .and_then(|ctx| ctx.agent_memory_provider.as_ref())
1484 .is_some_and(|provider| provider.supports_manifest())
1485 {
1486 methods.push("mobkit/agent_memory/manifest");
1487 }
1488 }
1489 if runtime.has_contact_directory() {
1491 methods.push("mobkit/cross_mob/directory");
1492 }
1493 if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
1497 methods.extend_from_slice(&[
1498 "mobkit/cross_mob/wire",
1499 "mobkit/cross_mob/unwire",
1500 "mobkit/cross_mob/send",
1501 ]);
1502 }
1503 let topology = runtime.topology_runtime_handle();
1504 let (topology_methods, topology_capabilities) =
1505 topology_methods::capability_projection(&topology, None, false);
1506 methods.extend(topology_methods);
1507 JsonRpcResponse {
1508 jsonrpc: JSONRPC_VERSION.to_string(),
1509 id: response_id,
1510 result: Some(serde_json::json!({
1511 "contract_version": MOBKIT_CONTRACT_VERSION,
1512 "runtime_type": "unified",
1513 "methods": methods,
1514 "identity_first": identity_ctx.is_some(),
1518 "workgraph": workgraph_configured,
1521 "loaded_modules": loaded,
1522 "runtime_capabilities": {
1523 "can_spawn_members": true,
1524 "can_send_messages": true,
1525 "can_wire_members": true,
1526 "can_retire_members": true,
1527 "available_spawn_modes": ["module", "profile"],
1528 },
1529 "authoring_capabilities": mobpack_authoring_capabilities(),
1530 "topology_control": topology_capabilities,
1531 })),
1532 error: None,
1533 }
1534 }
1535 topology_methods::TOPOLOGY_QUERY_METHOD => {
1536 let topology = runtime.topology_runtime_handle();
1537 topology_methods::handle_query(&topology, response_id, None, false).await
1538 }
1539 topology_methods::TOPOLOGY_PLAN_METHOD => {
1540 let topology = runtime.topology_runtime_handle();
1541 topology_methods::handle_plan(&topology, response_id, &request.params, None).await
1542 }
1543 topology_methods::TOPOLOGY_APPLY_METHOD => {
1544 let topology = runtime.topology_runtime_handle();
1545 topology_methods::handle_apply(
1546 &topology,
1547 response_id,
1548 &request.params,
1549 None,
1550 Some("local-host"),
1551 )
1552 .await
1553 }
1554 topology_methods::TOPOLOGY_OPERATION_METHOD => {
1555 let topology = runtime.topology_runtime_handle();
1556 topology_methods::handle_operation(&topology, response_id, &request.params, None).await
1557 }
1558 topology_methods::TOPOLOGY_AUDIT_METHOD => {
1559 let topology = runtime.topology_runtime_handle();
1560 topology_methods::handle_audit(&topology, response_id, &request.params, None).await
1561 }
1562 "mobkit/reconcile" => {
1563 let modules = match params::required_string_array(&request.params, "modules") {
1564 Ok(m) => m,
1565 Err(reason) => {
1566 return serialize_response(&JsonRpcResponse {
1567 jsonrpc: JSONRPC_VERSION.to_string(),
1568 id: response_id,
1569 result: None,
1570 error: Some(JsonRpcError {
1571 code: -32602,
1572 message: format!("Invalid params: {reason}"),
1573 data: None,
1574 }),
1575 });
1576 }
1577 };
1578
1579 match runtime.reconcile_modules(modules.clone(), timeout).await {
1580 Ok(added) => JsonRpcResponse {
1581 jsonrpc: JSONRPC_VERSION.to_string(),
1582 id: response_id,
1583 result: Some(serde_json::json!({
1584 "accepted": true,
1585 "reconciled_modules": modules,
1586 "added": added
1587 })),
1588 error: None,
1589 },
1590 Err(err) => JsonRpcResponse {
1591 jsonrpc: JSONRPC_VERSION.to_string(),
1592 id: response_id,
1593 result: None,
1594 error: Some(JsonRpcError {
1595 code: -32602,
1596 message: format!("Invalid params: {err:?}"),
1597 data: None,
1598 }),
1599 },
1600 }
1601 }
1602 "mobkit/spawn_member" => {
1603 let module_id = request.params.get("module_id").and_then(Value::as_str);
1605 let profile = request.params.get("profile").and_then(Value::as_str);
1606 let meerkat_id = request.params.get("meerkat_id").and_then(Value::as_str);
1607
1608 if let Some(module_id) = module_id {
1609 if module_id.is_empty() {
1611 JsonRpcResponse {
1612 jsonrpc: JSONRPC_VERSION.to_string(),
1613 id: response_id,
1614 result: None,
1615 error: Some(JsonRpcError {
1616 code: -32602,
1617 message: "Invalid params: module_id required".to_string(),
1618 data: None,
1619 }),
1620 }
1621 } else {
1622 match runtime.spawn_member(module_id, timeout).await {
1623 Ok(()) => JsonRpcResponse {
1624 jsonrpc: JSONRPC_VERSION.to_string(),
1625 id: response_id,
1626 result: Some(serde_json::json!({
1627 "accepted": true,
1628 "module_id": module_id
1629 })),
1630 error: None,
1631 },
1632 Err(err) => JsonRpcResponse {
1633 jsonrpc: JSONRPC_VERSION.to_string(),
1634 id: response_id,
1635 result: None,
1636 error: Some(JsonRpcError {
1637 code: -32602,
1638 message: format!("Invalid params: {err:?}"),
1639 data: None,
1640 }),
1641 },
1642 }
1643 }
1644 } else if let (Some(profile), Some(meerkat_id)) = (profile, meerkat_id) {
1645 if meerkat_id == crate::console_contracts::SYSTEM_EVENT_IDENTITY {
1650 JsonRpcResponse {
1651 jsonrpc: JSONRPC_VERSION.to_string(),
1652 id: response_id,
1653 result: None,
1654 error: Some(JsonRpcError {
1655 code: -32602,
1656 message: format!(
1657 "Invalid params: '{meerkat_id}' is a reserved runtime identity"
1658 ),
1659 data: None,
1660 }),
1661 }
1662 } else {
1663 let spec = meerkat_mob::SpawnMemberSpec::from_wire(
1665 profile.to_string(),
1666 meerkat_id.to_string(),
1667 request
1668 .params
1669 .get("initial_message")
1670 .and_then(Value::as_str)
1671 .map(|s| meerkat_core::ContentInput::from(s.to_string())),
1672 None,
1673 None,
1674 );
1675 match Box::pin(runtime.spawn(spec)).await {
1676 Ok(_member_ref) => JsonRpcResponse {
1677 jsonrpc: JSONRPC_VERSION.to_string(),
1678 id: response_id,
1679 result: Some(serde_json::json!({
1680 "accepted": true,
1681 "meerkat_id": meerkat_id
1682 })),
1683 error: None,
1684 },
1685 Err(err) => JsonRpcResponse {
1686 jsonrpc: JSONRPC_VERSION.to_string(),
1687 id: response_id,
1688 result: None,
1689 error: Some(JsonRpcError {
1690 code: -32602,
1691 message: format!("Invalid params: {err}"),
1692 data: None,
1693 }),
1694 },
1695 }
1696 }
1697 } else {
1698 JsonRpcResponse {
1699 jsonrpc: JSONRPC_VERSION.to_string(),
1700 id: response_id,
1701 result: None,
1702 error: Some(JsonRpcError {
1703 code: -32602,
1704 message: "Invalid params: module_id or (profile + meerkat_id) required"
1705 .to_string(),
1706 data: None,
1707 }),
1708 }
1709 }
1710 }
1711 "mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
1712 Ok((schedules, tick_ms)) => {
1713 match runtime.evaluate_schedule_tick(&schedules, tick_ms).await {
1714 Ok(evaluation) => JsonRpcResponse {
1715 jsonrpc: JSONRPC_VERSION.to_string(),
1716 id: response_id,
1717 result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
1718 error: None,
1719 },
1720 Err(err) => JsonRpcResponse {
1721 jsonrpc: JSONRPC_VERSION.to_string(),
1722 id: response_id,
1723 result: None,
1724 error: Some(JsonRpcError {
1725 code: -32602,
1726 message: format!(
1727 "Invalid params: {}",
1728 format_schedule_validation_error(err)
1729 ),
1730 data: None,
1731 }),
1732 },
1733 }
1734 }
1735 Err(message) => JsonRpcResponse {
1736 jsonrpc: JSONRPC_VERSION.to_string(),
1737 id: response_id,
1738 result: None,
1739 error: Some(JsonRpcError {
1740 code: -32602,
1741 message: format!("Invalid params: {message}"),
1742 data: None,
1743 }),
1744 },
1745 },
1746 "mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
1747 Ok((schedules, tick_ms)) => {
1748 match runtime.dispatch_schedule_tick(&schedules, tick_ms).await {
1749 Ok(dispatch) => JsonRpcResponse {
1750 jsonrpc: JSONRPC_VERSION.to_string(),
1751 id: response_id,
1752 result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
1753 error: None,
1754 },
1755 Err(err) => JsonRpcResponse {
1756 jsonrpc: JSONRPC_VERSION.to_string(),
1757 id: response_id,
1758 result: None,
1759 error: Some(JsonRpcError {
1760 code: -32602,
1761 message: format!("Invalid params: {err}"),
1762 data: None,
1763 }),
1764 },
1765 }
1766 }
1767 Err(message) => JsonRpcResponse {
1768 jsonrpc: JSONRPC_VERSION.to_string(),
1769 id: response_id,
1770 result: None,
1771 error: Some(JsonRpcError {
1772 code: -32602,
1773 message: format!("Invalid params: {message}"),
1774 data: None,
1775 }),
1776 },
1777 },
1778 "mobkit/routing/resolve" => {
1779 let resolve_result = match parse_routing_resolve_params(&request.params) {
1780 Ok(resolve_request) => runtime
1781 .resolve_routing(resolve_request)
1782 .await
1783 .map_err(RoutingDeliveryParamsError::Routing),
1784 Err(e) => Err(e),
1785 };
1786 match resolve_result {
1787 Ok(resolution) => JsonRpcResponse {
1788 jsonrpc: JSONRPC_VERSION.to_string(),
1789 id: response_id,
1790 result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
1791 error: None,
1792 },
1793 Err(err) => JsonRpcResponse {
1794 jsonrpc: JSONRPC_VERSION.to_string(),
1795 id: response_id,
1796 result: None,
1797 error: Some(JsonRpcError {
1798 code: -32602,
1799 message: format!("Invalid params: {}", err.message()),
1800 data: None,
1801 }),
1802 },
1803 }
1804 }
1805 "mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
1806 Ok(()) => {
1807 let routes = runtime.list_runtime_routes().await;
1808 JsonRpcResponse {
1809 jsonrpc: JSONRPC_VERSION.to_string(),
1810 id: response_id,
1811 result: Some(serde_json::json!({
1812 "routes": routes
1813 })),
1814 error: None,
1815 }
1816 }
1817 Err(err) => JsonRpcResponse {
1818 jsonrpc: JSONRPC_VERSION.to_string(),
1819 id: response_id,
1820 result: None,
1821 error: Some(JsonRpcError {
1822 code: -32602,
1823 message: format!("Invalid params: {}", err.message()),
1824 data: None,
1825 }),
1826 },
1827 },
1828 "mobkit/routing/routes/add" => {
1829 let add_result = match parse_routing_route_add_params(&request.params) {
1830 Ok(route) => runtime
1831 .add_runtime_route(route)
1832 .await
1833 .map_err(RoutingDeliveryParamsError::RouteMutation),
1834 Err(e) => Err(e),
1835 };
1836 match add_result {
1837 Ok(route) => JsonRpcResponse {
1838 jsonrpc: JSONRPC_VERSION.to_string(),
1839 id: response_id,
1840 result: Some(serde_json::json!({ "route": route })),
1841 error: None,
1842 },
1843 Err(err) => JsonRpcResponse {
1844 jsonrpc: JSONRPC_VERSION.to_string(),
1845 id: response_id,
1846 result: None,
1847 error: Some(JsonRpcError {
1848 code: -32602,
1849 message: format!("Invalid params: {}", err.message()),
1850 data: None,
1851 }),
1852 },
1853 }
1854 }
1855 "mobkit/routing/routes/delete" => {
1856 let delete_result = match parse_routing_route_delete_params(&request.params) {
1857 Ok(route_key) => runtime
1858 .delete_runtime_route(&route_key)
1859 .await
1860 .map_err(RoutingDeliveryParamsError::RouteMutation),
1861 Err(e) => Err(e),
1862 };
1863 match delete_result {
1864 Ok(route) => JsonRpcResponse {
1865 jsonrpc: JSONRPC_VERSION.to_string(),
1866 id: response_id,
1867 result: Some(serde_json::json!({ "deleted": route })),
1868 error: None,
1869 },
1870 Err(err) => JsonRpcResponse {
1871 jsonrpc: JSONRPC_VERSION.to_string(),
1872 id: response_id,
1873 result: None,
1874 error: Some(JsonRpcError {
1875 code: -32602,
1876 message: format!("Invalid params: {}", err.message()),
1877 data: None,
1878 }),
1879 },
1880 }
1881 }
1882 "mobkit/delivery/send" => {
1883 let send_result = match parse_delivery_send_params(&request.params) {
1884 Ok(send_request) => runtime
1885 .send_delivery(send_request)
1886 .await
1887 .map_err(RoutingDeliveryParamsError::Delivery),
1888 Err(e) => Err(e),
1889 };
1890 match send_result {
1891 Ok(record) => JsonRpcResponse {
1892 jsonrpc: JSONRPC_VERSION.to_string(),
1893 id: response_id,
1894 result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
1895 error: None,
1896 },
1897 Err(err) => JsonRpcResponse {
1898 jsonrpc: JSONRPC_VERSION.to_string(),
1899 id: response_id,
1900 result: None,
1901 error: Some(JsonRpcError {
1902 code: -32602,
1903 message: format!("Invalid params: {}", err.message()),
1904 data: None,
1905 }),
1906 },
1907 }
1908 }
1909 "mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
1910 Ok(history_request) => {
1911 let history = runtime.delivery_history(history_request).await;
1912 JsonRpcResponse {
1913 jsonrpc: JSONRPC_VERSION.to_string(),
1914 id: response_id,
1915 result: Some(serde_json::to_value(history).unwrap_or(Value::Null)),
1916 error: None,
1917 }
1918 }
1919 Err(err) => JsonRpcResponse {
1920 jsonrpc: JSONRPC_VERSION.to_string(),
1921 id: response_id,
1922 result: None,
1923 error: Some(JsonRpcError {
1924 code: -32602,
1925 message: format!("Invalid params: {}", err.message()),
1926 data: None,
1927 }),
1928 },
1929 },
1930 "mobkit/events/subscribe" => match parse_subscribe_request(&request.params) {
1931 Ok(subscribe_request) => match runtime.subscribe_events(subscribe_request).await {
1932 Ok(subscribe_result) => JsonRpcResponse {
1933 jsonrpc: JSONRPC_VERSION.to_string(),
1934 id: response_id,
1935 result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
1936 error: None,
1937 },
1938 Err(err) => JsonRpcResponse {
1939 jsonrpc: JSONRPC_VERSION.to_string(),
1940 id: response_id,
1941 result: None,
1942 error: Some(JsonRpcError {
1943 code: -32602,
1944 message: format!("Invalid params: {err}"),
1945 data: None,
1946 }),
1947 },
1948 },
1949 Err(err) => JsonRpcResponse {
1950 jsonrpc: JSONRPC_VERSION.to_string(),
1951 id: response_id,
1952 result: None,
1953 error: Some(JsonRpcError {
1954 code: -32602,
1955 message: format!("Invalid params: {}", err.message()),
1956 data: None,
1957 }),
1958 },
1959 },
1960 "mobkit/query_events" => {
1961 let query: EventQuery = if request.params.is_null() {
1962 EventQuery::default()
1963 } else {
1964 match serde_json::from_value(request.params.clone()) {
1965 Ok(query) => query,
1966 Err(err) => {
1967 return serde_json::to_string(&JsonRpcResponse {
1968 jsonrpc: JSONRPC_VERSION.to_string(),
1969 id: response_id,
1970 result: None,
1971 error: Some(JsonRpcError {
1972 code: -32602,
1973 message: format!("Invalid params: invalid query params: {err}"),
1974 data: None,
1975 }),
1976 })
1977 .unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string());
1978 }
1979 }
1980 };
1981 match runtime.event_log_store() {
1982 Some(store) => match store.query(query).await {
1983 Ok(events) => JsonRpcResponse {
1984 jsonrpc: JSONRPC_VERSION.to_string(),
1985 id: response_id,
1986 result: Some(serde_json::to_value(events).unwrap_or(Value::Null)),
1987 error: None,
1988 },
1989 Err(err) => JsonRpcResponse {
1990 jsonrpc: JSONRPC_VERSION.to_string(),
1991 id: response_id,
1992 result: None,
1993 error: Some(JsonRpcError {
1994 code: -32603,
1995 message: format!("query_events failed: {err}"),
1996 data: None,
1997 }),
1998 },
1999 },
2000 None => JsonRpcResponse {
2001 jsonrpc: JSONRPC_VERSION.to_string(),
2002 id: response_id,
2003 result: Some(serde_json::json!({
2004 "status": "no_event_log_configured",
2005 "events": [],
2006 })),
2007 error: None,
2008 },
2009 }
2010 }
2011 "mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
2012 Ok(()) => {
2013 let stores = runtime.memory_stores().await;
2014 JsonRpcResponse {
2015 jsonrpc: JSONRPC_VERSION.to_string(),
2016 id: response_id,
2017 result: Some(serde_json::json!({
2018 "stores": stores,
2019 })),
2020 error: None,
2021 }
2022 }
2023 Err(err) => JsonRpcResponse {
2024 jsonrpc: JSONRPC_VERSION.to_string(),
2025 id: response_id,
2026 result: None,
2027 error: Some(JsonRpcError {
2028 code: -32602,
2029 message: format!("Invalid params: {}", err.message()),
2030 data: None,
2031 }),
2032 },
2033 },
2034 "mobkit/memory/index" => match parse_memory_index_params(&request.params) {
2035 Ok(index_request) => match runtime.memory_index(index_request).await {
2036 Ok(indexed) => JsonRpcResponse {
2037 jsonrpc: JSONRPC_VERSION.to_string(),
2038 id: response_id,
2039 result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
2040 error: None,
2041 },
2042 Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
2043 jsonrpc: JSONRPC_VERSION.to_string(),
2044 id: response_id,
2045 result: None,
2046 error: Some(JsonRpcError {
2047 code: MEMORY_BACKEND_UNAVAILABLE_CODE,
2048 message: format!(
2049 "Memory backend unavailable: {}",
2050 MemoryParamsError::backend_message(&error)
2051 ),
2052 data: None,
2053 }),
2054 },
2055 Err(err) => JsonRpcResponse {
2056 jsonrpc: JSONRPC_VERSION.to_string(),
2057 id: response_id,
2058 result: None,
2059 error: Some(JsonRpcError {
2060 code: -32602,
2061 message: format!(
2062 "Invalid params: {}",
2063 MemoryParamsError::Index(err).message()
2064 ),
2065 data: None,
2066 }),
2067 },
2068 },
2069 Err(err) => JsonRpcResponse {
2070 jsonrpc: JSONRPC_VERSION.to_string(),
2071 id: response_id,
2072 result: None,
2073 error: Some(JsonRpcError {
2074 code: -32602,
2075 message: format!("Invalid params: {}", err.message()),
2076 data: None,
2077 }),
2078 },
2079 },
2080 "mobkit/memory/query" => match parse_memory_query_params(&request.params) {
2081 Ok(query_request) => {
2082 let query_result = runtime.memory_query(query_request).await;
2083 JsonRpcResponse {
2084 jsonrpc: JSONRPC_VERSION.to_string(),
2085 id: response_id,
2086 result: Some(serde_json::to_value(query_result).unwrap_or(Value::Null)),
2087 error: None,
2088 }
2089 }
2090 Err(err) => JsonRpcResponse {
2091 jsonrpc: JSONRPC_VERSION.to_string(),
2092 id: response_id,
2093 result: None,
2094 error: Some(JsonRpcError {
2095 code: -32602,
2096 message: format!("Invalid params: {}", err.message()),
2097 data: None,
2098 }),
2099 },
2100 },
2101 "mobkit/agent_memory/remember" => {
2102 let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
2103 Some(runtime) => runtime,
2104 None => {
2105 return maybe_error_response(
2106 is_notification,
2107 response_id,
2108 -32601,
2109 "agent memory is not configured".to_string(),
2110 );
2111 }
2112 };
2113 match parse_agent_memory_remember_params(&request.params) {
2114 Ok(remember_request) => match runtime
2115 .remember_agent_memory(
2116 &remember_request.realm,
2117 &remember_request.identity,
2118 remember_request.memory,
2119 )
2120 .await
2121 {
2122 Ok(record) => JsonRpcResponse {
2123 jsonrpc: JSONRPC_VERSION.to_string(),
2124 id: response_id,
2125 result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
2126 error: None,
2127 },
2128 Err(err) => JsonRpcResponse {
2129 jsonrpc: JSONRPC_VERSION.to_string(),
2130 id: response_id,
2131 result: None,
2132 error: Some(agent_memory_rpc_error("write", err)),
2133 },
2134 },
2135 Err(err) => JsonRpcResponse {
2136 jsonrpc: JSONRPC_VERSION.to_string(),
2137 id: response_id,
2138 result: None,
2139 error: Some(JsonRpcError {
2140 code: -32602,
2141 message: format!("Invalid params: {}", err.message()),
2142 data: None,
2143 }),
2144 },
2145 }
2146 }
2147 "mobkit/agent_memory/forget" => {
2148 let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
2149 Some(runtime) => runtime,
2150 None => {
2151 return maybe_error_response(
2152 is_notification,
2153 response_id,
2154 -32601,
2155 "agent memory is not configured".to_string(),
2156 );
2157 }
2158 };
2159 match parse_agent_memory_forget_params(&request.params) {
2160 Ok(forget_request) => match runtime
2161 .forget_agent_memory(
2162 &forget_request.realm,
2163 &forget_request.identity,
2164 &forget_request.memory_id,
2165 )
2166 .await
2167 {
2168 Ok(result) => JsonRpcResponse {
2169 jsonrpc: JSONRPC_VERSION.to_string(),
2170 id: response_id,
2171 result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
2172 error: None,
2173 },
2174 Err(err) => JsonRpcResponse {
2175 jsonrpc: JSONRPC_VERSION.to_string(),
2176 id: response_id,
2177 result: None,
2178 error: Some(agent_memory_rpc_error("forget", err)),
2179 },
2180 },
2181 Err(err) => JsonRpcResponse {
2182 jsonrpc: JSONRPC_VERSION.to_string(),
2183 id: response_id,
2184 result: None,
2185 error: Some(JsonRpcError {
2186 code: -32602,
2187 message: format!("Invalid params: {}", err.message()),
2188 data: None,
2189 }),
2190 },
2191 }
2192 }
2193 "mobkit/agent_memory/recall" => {
2194 let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
2195 Some(runtime) => runtime,
2196 None => {
2197 return maybe_error_response(
2198 is_notification,
2199 response_id,
2200 -32601,
2201 "agent memory is not configured".to_string(),
2202 );
2203 }
2204 };
2205 match parse_agent_memory_recall_params(&request.params) {
2206 Ok(recall_request) => {
2207 match runtime.recall_agent_memory(recall_request.request).await {
2208 Ok(records) => JsonRpcResponse {
2209 jsonrpc: JSONRPC_VERSION.to_string(),
2210 id: response_id,
2211 result: Some(serde_json::json!({ "records": records })),
2212 error: None,
2213 },
2214 Err(err) => JsonRpcResponse {
2215 jsonrpc: JSONRPC_VERSION.to_string(),
2216 id: response_id,
2217 result: None,
2218 error: Some(agent_memory_rpc_error("recall", err)),
2219 },
2220 }
2221 }
2222 Err(err) => JsonRpcResponse {
2223 jsonrpc: JSONRPC_VERSION.to_string(),
2224 id: response_id,
2225 result: None,
2226 error: Some(JsonRpcError {
2227 code: -32602,
2228 message: format!("Invalid params: {}", err.message()),
2229 data: None,
2230 }),
2231 },
2232 }
2233 }
2234 "mobkit/agent_memory/update" => {
2235 let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
2236 Some(runtime) => runtime,
2237 None => {
2238 return maybe_error_response(
2239 is_notification,
2240 response_id,
2241 -32601,
2242 "agent memory is not configured".to_string(),
2243 );
2244 }
2245 };
2246 match parse_agent_memory_update_params(&request.params) {
2247 Ok(update_request) => match runtime
2248 .update_agent_memory(
2249 &update_request.realm,
2250 &update_request.identity,
2251 &update_request.memory_id,
2252 update_request.memory,
2253 )
2254 .await
2255 {
2256 Ok(new_id) => JsonRpcResponse {
2257 jsonrpc: JSONRPC_VERSION.to_string(),
2258 id: response_id,
2259 result: Some(serde_json::json!({
2260 "memory_id": new_id,
2261 "supersedes": update_request.memory_id,
2262 })),
2263 error: None,
2264 },
2265 Err(err) => JsonRpcResponse {
2266 jsonrpc: JSONRPC_VERSION.to_string(),
2267 id: response_id,
2268 result: None,
2269 error: Some(agent_memory_rpc_error("update", err)),
2270 },
2271 },
2272 Err(err) => JsonRpcResponse {
2273 jsonrpc: JSONRPC_VERSION.to_string(),
2274 id: response_id,
2275 result: None,
2276 error: Some(JsonRpcError {
2277 code: -32602,
2278 message: format!("Invalid params: {}", err.message()),
2279 data: None,
2280 }),
2281 },
2282 }
2283 }
2284 "mobkit/agent_memory/manifest" => {
2285 let runtime = match identity_ctx.map(|ctx| ctx.runtime.as_ref()) {
2286 Some(runtime) => runtime,
2287 None => {
2288 return maybe_error_response(
2289 is_notification,
2290 response_id,
2291 -32601,
2292 "agent memory is not configured".to_string(),
2293 );
2294 }
2295 };
2296 match parse_agent_memory_manifest_params(&request.params) {
2297 Ok(manifest_request) => match runtime
2298 .manifest_agent_memory(
2299 &manifest_request.realm,
2300 &manifest_request.identity,
2301 manifest_request.tier,
2302 )
2303 .await
2304 {
2305 Ok(records) => JsonRpcResponse {
2306 jsonrpc: JSONRPC_VERSION.to_string(),
2307 id: response_id,
2308 result: Some(serde_json::json!({ "records": records })),
2309 error: None,
2310 },
2311 Err(err) => JsonRpcResponse {
2312 jsonrpc: JSONRPC_VERSION.to_string(),
2313 id: response_id,
2314 result: None,
2315 error: Some(agent_memory_rpc_error("manifest", err)),
2316 },
2317 },
2318 Err(err) => JsonRpcResponse {
2319 jsonrpc: JSONRPC_VERSION.to_string(),
2320 id: response_id,
2321 result: None,
2322 error: Some(JsonRpcError {
2323 code: -32602,
2324 message: format!("Invalid params: {}", err.message()),
2325 data: None,
2326 }),
2327 },
2328 }
2329 }
2330 "mobkit/session_store/bigquery" => {
2331 match parse_bigquery_session_store_params(&request.params)
2332 .and_then(run_bigquery_session_store_request)
2333 {
2334 Ok(result) => JsonRpcResponse {
2335 jsonrpc: JSONRPC_VERSION.to_string(),
2336 id: response_id,
2337 result: Some(result),
2338 error: None,
2339 },
2340 Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
2341 jsonrpc: JSONRPC_VERSION.to_string(),
2342 id: response_id,
2343 result: None,
2344 error: Some(JsonRpcError {
2345 code: -32602,
2346 message: format!("Invalid params: {message}"),
2347 data: None,
2348 }),
2349 },
2350 Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
2351 jsonrpc: JSONRPC_VERSION.to_string(),
2352 id: response_id,
2353 result: None,
2354 error: Some(JsonRpcError {
2355 code: -32011,
2356 message: format!(
2357 "BigQuery session store request failed: {}",
2358 format_bigquery_store_error(&error)
2359 ),
2360 data: None,
2361 }),
2362 },
2363 }
2364 }
2365 "mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
2366 Ok(gating_request) => {
2367 let gating_result = runtime.evaluate_gating_action(gating_request).await;
2368 JsonRpcResponse {
2369 jsonrpc: JSONRPC_VERSION.to_string(),
2370 id: response_id,
2371 result: Some(serde_json::to_value(gating_result).unwrap_or(Value::Null)),
2372 error: None,
2373 }
2374 }
2375 Err(err) => JsonRpcResponse {
2376 jsonrpc: JSONRPC_VERSION.to_string(),
2377 id: response_id,
2378 result: None,
2379 error: Some(JsonRpcError {
2380 code: -32602,
2381 message: format!("Invalid params: {}", err.message()),
2382 data: None,
2383 }),
2384 },
2385 },
2386 "mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
2387 Ok(()) => {
2388 let pending = runtime.list_gating_pending().await;
2389 JsonRpcResponse {
2390 jsonrpc: JSONRPC_VERSION.to_string(),
2391 id: response_id,
2392 result: Some(serde_json::json!({
2393 "pending": pending,
2394 })),
2395 error: None,
2396 }
2397 }
2398 Err(err) => JsonRpcResponse {
2399 jsonrpc: JSONRPC_VERSION.to_string(),
2400 id: response_id,
2401 result: None,
2402 error: Some(JsonRpcError {
2403 code: -32602,
2404 message: format!("Invalid params: {}", err.message()),
2405 data: None,
2406 }),
2407 },
2408 },
2409 "mobkit/gating/decide" => {
2410 let decide_result = match parse_gating_decide_params(&request.params) {
2411 Ok(decide_request) => runtime
2412 .decide_gating_action(decide_request)
2413 .await
2414 .map_err(GatingParamsError::Decision),
2415 Err(e) => Err(e),
2416 };
2417 match decide_result {
2418 Ok(result) => JsonRpcResponse {
2419 jsonrpc: JSONRPC_VERSION.to_string(),
2420 id: response_id,
2421 result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
2422 error: None,
2423 },
2424 Err(err) => JsonRpcResponse {
2425 jsonrpc: JSONRPC_VERSION.to_string(),
2426 id: response_id,
2427 result: None,
2428 error: Some(JsonRpcError {
2429 code: -32602,
2430 message: format!("Invalid params: {}", err.message()),
2431 data: None,
2432 }),
2433 },
2434 }
2435 }
2436 "mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
2437 Ok(limit) => {
2438 let entries = runtime.gating_audit_entries(limit).await;
2439 JsonRpcResponse {
2440 jsonrpc: JSONRPC_VERSION.to_string(),
2441 id: response_id,
2442 result: Some(serde_json::json!({
2443 "entries": entries,
2444 })),
2445 error: None,
2446 }
2447 }
2448 Err(err) => JsonRpcResponse {
2449 jsonrpc: JSONRPC_VERSION.to_string(),
2450 id: response_id,
2451 result: None,
2452 error: Some(JsonRpcError {
2453 code: -32602,
2454 message: format!("Invalid params: {}", err.message()),
2455 data: None,
2456 }),
2457 },
2458 },
2459 "mobkit/call_tool" => {
2460 let module_id = request.params.get("module_id").and_then(Value::as_str);
2461 let tool = request.params.get("tool").and_then(Value::as_str);
2462 let arguments = request
2463 .params
2464 .get("arguments")
2465 .cloned()
2466 .unwrap_or(serde_json::json!({}));
2467
2468 match (module_id, tool) {
2469 (Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
2470 let route = runtime
2471 .route_module_call(
2472 &ModuleRouteRequest {
2473 module_id: module_id.to_string(),
2474 method: tool.to_string(),
2475 params: arguments,
2476 },
2477 timeout,
2478 )
2479 .await;
2480 match route {
2481 Ok(response) => JsonRpcResponse {
2482 jsonrpc: JSONRPC_VERSION.to_string(),
2483 id: response_id,
2484 result: Some(serde_json::json!({
2485 "module_id": response.module_id,
2486 "tool": response.method,
2487 "result": response.payload
2488 })),
2489 error: None,
2490 },
2491 Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
2492 jsonrpc: JSONRPC_VERSION.to_string(),
2493 id: response_id,
2494 result: None,
2495 error: Some(JsonRpcError {
2496 code: -32601,
2497 message: format!("Module '{mid}' not loaded"),
2498 data: None,
2499 }),
2500 },
2501 Err(err) => JsonRpcResponse {
2502 jsonrpc: JSONRPC_VERSION.to_string(),
2503 id: response_id,
2504 result: None,
2505 error: Some(JsonRpcError {
2506 code: -32000,
2507 message: format!("Tool call failed: {err:?}"),
2508 data: None,
2509 }),
2510 },
2511 }
2512 }
2513 _ => JsonRpcResponse {
2514 jsonrpc: JSONRPC_VERSION.to_string(),
2515 id: response_id,
2516 result: None,
2517 error: Some(JsonRpcError {
2518 code: -32602,
2519 message: "Invalid params: module_id and tool required".to_string(),
2520 data: None,
2521 }),
2522 },
2523 }
2524 }
2525 "mobkit/models/catalog" => JsonRpcResponse {
2526 jsonrpc: JSONRPC_VERSION.to_string(),
2527 id: response_id,
2528 result: Some(build_models_catalog_result()),
2529 error: None,
2530 },
2531 method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
2532 handle_unified_mobpack_authoring_rpc(runtime, method, &request.params, response_id)
2533 .await
2534 }
2535 "mobkit/blob/get" => {
2536 mob_methods::handle_blob_get(runtime, response_id, &request.params).await
2537 }
2538 "mobkit/send_message" => {
2539 Box::pin(mob_methods::handle_send_message(
2543 runtime,
2544 identity_ctx.map(|ctx| &ctx.runtime),
2545 response_id,
2546 &request.params,
2547 ))
2548 .await
2549 }
2550 "mobkit/find_members" => {
2551 mob_methods::handle_find_members(
2552 runtime,
2553 identity_ctx.map(|ctx| &ctx.runtime),
2554 response_id,
2555 &request.params,
2556 )
2557 .await
2558 }
2559 "mobkit/ensure_member" => {
2560 Box::pin(mob_methods::handle_ensure_member(
2561 runtime,
2562 response_id,
2563 &request.params,
2564 ))
2565 .await
2566 }
2567 "mobkit/list_members" => {
2568 mob_methods::handle_list_members(
2569 runtime,
2570 identity_ctx.map(|ctx| &ctx.runtime),
2571 response_id,
2572 )
2573 .await
2574 }
2575 "mobkit/get_member" => {
2576 mob_methods::handle_get_member(
2577 runtime,
2578 identity_ctx.map(|ctx| &ctx.runtime),
2579 response_id,
2580 &request.params,
2581 )
2582 .await
2583 }
2584 "mobkit/retire_member" => {
2585 mob_methods::handle_retire_member(
2586 runtime,
2587 identity_ctx.map(|ctx| &ctx.runtime),
2588 response_id,
2589 &request.params,
2590 )
2591 .await
2592 }
2593 "mobkit/respawn_member" => {
2594 Box::pin(mob_methods::handle_respawn_member(
2595 runtime,
2596 identity_ctx.map(|ctx| &ctx.runtime),
2597 response_id,
2598 &request.params,
2599 ))
2600 .await
2601 }
2602 "mobkit/reconcile_edges" => mob_methods::handle_reconcile_edges(runtime, response_id).await,
2603 "mobkit/rediscover" => mob_methods::handle_rediscover(runtime, response_id).await,
2604 "mobkit/mob_events/query" => {
2605 mob_methods::handle_mob_events_query(runtime, response_id, request.params).await
2606 }
2607 "mobkit/mob_events/subscribe" => {
2608 mob_methods::handle_mob_events_subscribe(runtime, response_id, request.params).await
2609 }
2610 "mobkit/cross_mob/wire" => {
2611 Box::pin(mob_methods::handle_cross_mob_wire(
2612 runtime,
2613 response_id,
2614 &request.params,
2615 ))
2616 .await
2617 }
2618 "mobkit/cross_mob/unwire" => {
2619 Box::pin(mob_methods::handle_cross_mob_unwire(
2620 runtime,
2621 response_id,
2622 &request.params,
2623 ))
2624 .await
2625 }
2626 "mobkit/cross_mob/send" => {
2627 mob_methods::handle_cross_mob_send(runtime, response_id, &request.params).await
2628 }
2629 "mobkit/cross_mob/directory" => {
2630 mob_methods::handle_cross_mob_directory(runtime, response_id).await
2631 }
2632 "mobkit/cross_mob/peer_info" => {
2633 mob_methods::handle_cross_mob_peer_info(runtime, response_id, &request.params).await
2634 }
2635 "mobkit/cross_mob/wire_local" => {
2636 mob_methods::handle_cross_mob_wire_local(runtime, response_id, &request.params).await
2637 }
2638 "mobkit/cross_mob/unwire_local" => {
2639 mob_methods::handle_cross_mob_unwire_local(runtime, response_id, &request.params).await
2640 }
2641 "mobkit/peer_pubkey" => mob_methods::handle_peer_pubkey(runtime, response_id).await,
2642 "mobkit/member_status" => {
2643 mob_methods::handle_member_status(
2644 runtime,
2645 identity_ctx.map(|ctx| &ctx.runtime),
2646 response_id,
2647 &request.params,
2648 )
2649 .await
2650 }
2651 "mobkit/identity/resolved_tools" => {
2652 mob_methods::handle_identity_resolved_tools(
2653 runtime,
2654 identity_ctx.map(|ctx| &ctx.runtime),
2655 response_id,
2656 &request.params,
2657 )
2658 .await
2659 }
2660 "mobkit/force_cancel_member" => {
2661 mob_methods::handle_force_cancel_member(
2662 runtime,
2663 identity_ctx.map(|ctx| &ctx.runtime),
2664 response_id,
2665 &request.params,
2666 )
2667 .await
2668 }
2669 "mobkit/spawn_helper" => {
2670 Box::pin(mob_methods::handle_spawn_helper(
2671 runtime,
2672 response_id,
2673 &request.params,
2674 ))
2675 .await
2676 }
2677 "mobkit/fork_helper" => {
2678 Box::pin(mob_methods::handle_fork_helper(
2679 runtime,
2680 response_id,
2681 &request.params,
2682 ))
2683 .await
2684 }
2685 "mobkit/attach_existing_session" => {
2686 Box::pin(mob_methods::handle_attach_existing_session(
2687 runtime,
2688 response_id,
2689 &request.params,
2690 ))
2691 .await
2692 }
2693 "mobkit/cancel_flow" => {
2694 mob_methods::handle_cancel_flow(runtime, response_id, &request.params).await
2695 }
2696 "mobkit/flow_status" => {
2697 mob_methods::handle_flow_status(runtime, response_id, &request.params).await
2698 }
2699 "mobkit/list_flows" => mob_methods::handle_list_flows(runtime, response_id).await,
2700 "mobkit/list_runs" => {
2701 mob_methods::handle_list_runs(runtime, response_id, &request.params).await
2702 }
2703 "mobkit/run_flow" => {
2704 Box::pin(mob_methods::handle_run_flow(
2705 runtime,
2706 response_id,
2707 &request.params,
2708 ))
2709 .await
2710 }
2711 "mobkit/collect_completed" => {
2712 mob_methods::handle_collect_completed(runtime, response_id).await
2713 }
2714 "mobkit/wait_ready" => {
2715 mob_methods::handle_wait_ready(runtime, response_id, &request.params).await
2716 }
2717 "mobkit/mob_labels/set" => {
2718 mob_methods::handle_mob_labels_set(runtime, response_id, &request.params).await
2719 }
2720 "mobkit/mob_labels/get" => mob_methods::handle_mob_labels_get(runtime, response_id).await,
2721 "mobkit/mob_labels/delete" => {
2722 mob_methods::handle_mob_labels_delete(runtime, response_id).await
2723 }
2724 "mobkit/run_labels/set" => {
2725 mob_methods::handle_run_labels_set(runtime, response_id, &request.params).await
2726 }
2727 "mobkit/run_labels/get" => {
2728 mob_methods::handle_run_labels_get(runtime, response_id, &request.params).await
2729 }
2730 "mobkit/run_labels/delete" => {
2731 mob_methods::handle_run_labels_delete(runtime, response_id, &request.params).await
2732 }
2733 "mobkit/send" => {
2735 let identity_rt = match identity_ctx {
2736 Some(ctx) => &*ctx.runtime,
2737 None => return maybe_identity_not_configured(is_notification, response_id),
2738 };
2739 let identity_str = request
2740 .params
2741 .get("identity")
2742 .and_then(|v| v.as_str())
2743 .unwrap_or("");
2744 let target =
2745 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2746 {
2747 Ok(target) => target,
2748 Err(e) => {
2749 return maybe_error_response(
2750 is_notification,
2751 response_id,
2752 -32602,
2753 format!("invalid identity: {e}"),
2754 );
2755 }
2756 };
2757 let identity = target.identity.clone();
2758 if let Some(response) =
2759 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2760 {
2761 return if is_notification {
2762 String::new()
2763 } else {
2764 serialize_response(&response)
2765 };
2766 }
2767 let content_val = request
2768 .params
2769 .get("content")
2770 .cloned()
2771 .unwrap_or(Value::Null);
2772 let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
2773 Ok(content) => content,
2774 Err(err) => {
2775 return maybe_error_response(
2776 is_notification,
2777 response_id,
2778 -32602,
2779 format!("invalid content: {err}"),
2780 );
2781 }
2782 };
2783 match identity_rt.send(&identity, &content).await {
2784 Ok(token) => JsonRpcResponse {
2785 jsonrpc: JSONRPC_VERSION.to_string(),
2786 id: response_id,
2787 result: Some(serde_json::json!({ "fencing_token": token.get() })),
2788 error: None,
2789 },
2790 Err(e) => identity_error_response(response_id, &e),
2791 }
2792 }
2793 "mobkit/interact" => {
2794 let identity_rt = match identity_ctx {
2795 Some(ctx) => &*ctx.runtime,
2796 None => return maybe_identity_not_configured(is_notification, response_id),
2797 };
2798 let identity_str = request
2799 .params
2800 .get("identity")
2801 .and_then(|v| v.as_str())
2802 .unwrap_or("");
2803 let target =
2804 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2805 {
2806 Ok(target) => target,
2807 Err(e) => {
2808 return maybe_error_response(
2809 is_notification,
2810 response_id,
2811 -32602,
2812 format!("invalid identity: {e}"),
2813 );
2814 }
2815 };
2816 let identity = target.identity.clone();
2817 if let Some(response) =
2818 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2819 {
2820 return if is_notification {
2821 String::new()
2822 } else {
2823 serialize_response(&response)
2824 };
2825 }
2826 let content_val = request
2827 .params
2828 .get("content")
2829 .cloned()
2830 .unwrap_or(Value::Null);
2831 let content =
2832 match serde_json::from_value::<meerkat_core::ContentInput>(content_val.clone()) {
2833 Ok(content) => content,
2834 Err(err) => {
2835 return maybe_error_response(
2836 is_notification,
2837 response_id,
2838 -32602,
2839 format!("invalid content: {err}"),
2840 );
2841 }
2842 };
2843 let origin = request
2844 .params
2845 .get("origin")
2846 .and_then(|v| v.as_str())
2847 .unwrap_or("console");
2848 let interaction_id = request
2849 .params
2850 .get("interaction_id")
2851 .and_then(|v| v.as_str())
2852 .map(ToString::to_string)
2853 .unwrap_or_else(|| meerkat_core::types::SessionId::new().to_string());
2854 let runtime_member_id = identity_rt
2855 .status(&identity)
2856 .await
2857 .ok()
2858 .and_then(|status| status.agent_runtime_id.map(|id| id.as_str().to_string()));
2859
2860 if let Err(err) = runtime
2861 .reserve_identity_interaction(
2862 identity.as_str(),
2863 runtime_member_id.as_deref(),
2864 &interaction_id,
2865 origin,
2866 content_val,
2867 )
2868 .await
2869 {
2870 return maybe_error_response(
2871 is_notification,
2872 response_id,
2873 -32003,
2874 format!("failed to reserve interaction: {err}"),
2875 );
2876 }
2877
2878 match identity_rt.send(&identity, &content).await {
2879 Ok(token) => JsonRpcResponse {
2880 jsonrpc: JSONRPC_VERSION.to_string(),
2881 id: response_id,
2882 result: Some(serde_json::json!({
2883 "interaction_id": interaction_id,
2884 "fencing_token": token.get(),
2885 "stream": {
2886 "route": format!("/console/identity/{}/stream", identity.as_str()),
2887 "identity": identity.as_str(),
2888 }
2889 })),
2890 error: None,
2891 },
2892 Err(e) => {
2893 runtime
2894 .record_console_lifecycle(
2895 identity.as_str(),
2896 "interaction_failed",
2897 serde_json::json!({
2898 "interaction_id": interaction_id,
2899 "origin": origin,
2900 "error": e.to_string(),
2901 }),
2902 )
2903 .await;
2904 identity_error_response(response_id, &e)
2905 }
2906 }
2907 }
2908 "mobkit/dispatch" => {
2909 let identity_rt = match identity_ctx {
2910 Some(ctx) => &*ctx.runtime,
2911 None => return maybe_identity_not_configured(is_notification, response_id),
2912 };
2913 let identity_str = request
2914 .params
2915 .get("identity")
2916 .and_then(|v| v.as_str())
2917 .unwrap_or("");
2918 let target =
2919 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2920 {
2921 Ok(target) => target,
2922 Err(e) => {
2923 return maybe_error_response(
2924 is_notification,
2925 response_id,
2926 -32602,
2927 format!("invalid identity: {e}"),
2928 );
2929 }
2930 };
2931 let identity = target.identity.clone();
2932 if let Some(response) =
2933 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2934 {
2935 return if is_notification {
2936 String::new()
2937 } else {
2938 serialize_response(&response)
2939 };
2940 }
2941 let di_val = request
2942 .params
2943 .get("dispatch_input")
2944 .cloned()
2945 .unwrap_or(Value::Null);
2946 let content_val = di_val
2947 .get("content")
2948 .cloned()
2949 .unwrap_or_else(|| Value::String(String::new()));
2950 let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
2951 Ok(content) => content,
2952 Err(err) => {
2953 return maybe_error_response(
2954 is_notification,
2955 response_id,
2956 -32602,
2957 format!("invalid dispatch_input.content: {err}"),
2958 );
2959 }
2960 };
2961 let origin_str = di_val
2962 .get("origin")
2963 .and_then(|v| v.as_str())
2964 .unwrap_or("system");
2965 let origin = match origin_str {
2966 "connector" => crate::identity_first::DispatchOrigin::Connector,
2967 "scheduler" => crate::identity_first::DispatchOrigin::Scheduler,
2968 "policy" => crate::identity_first::DispatchOrigin::Policy,
2969 "flow" => crate::identity_first::DispatchOrigin::Flow,
2970 _ => crate::identity_first::DispatchOrigin::System,
2971 };
2972 let correlation_id = di_val
2973 .get("correlation_id")
2974 .and_then(|v| v.as_str())
2975 .map(crate::identity_first::CorrelationId::new);
2976 let idempotency_key = di_val
2977 .get("idempotency_key")
2978 .and_then(|v| v.as_str())
2979 .map(crate::identity_first::DispatchIdempotencyKey::new);
2980 let dispatch_input = crate::identity_first::DispatchInput {
2981 content,
2982 origin,
2983 correlation_id,
2984 idempotency_key,
2985 };
2986 match identity_rt.dispatch(&identity, &dispatch_input).await {
2987 Ok((token, durable)) => JsonRpcResponse {
2988 jsonrpc: JSONRPC_VERSION.to_string(),
2989 id: response_id,
2990 result: Some(
2991 serde_json::json!({ "fencing_token": token.get(), "durable": durable }),
2992 ),
2993 error: None,
2994 },
2995 Err(e) => identity_error_response(response_id, &e),
2996 }
2997 }
2998 "mobkit/subscribe" => {
2999 let identity_rt = match identity_ctx {
3000 Some(ctx) => &*ctx.runtime,
3001 None => return maybe_identity_not_configured(is_notification, response_id),
3002 };
3003 let identity_str = request
3004 .params
3005 .get("identity")
3006 .and_then(|v| v.as_str())
3007 .unwrap_or("");
3008 let target =
3009 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3010 {
3011 Ok(target) => target,
3012 Err(e) => {
3013 return maybe_error_response(
3014 is_notification,
3015 response_id,
3016 -32602,
3017 format!("invalid identity: {e}"),
3018 );
3019 }
3020 };
3021 let identity = target.identity.clone();
3022 if let Some(response) =
3023 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3024 {
3025 return if is_notification {
3026 String::new()
3027 } else {
3028 serialize_response(&response)
3029 };
3030 }
3031 match identity_rt.subscribe(&identity).await {
3032 Ok(_receiver) => JsonRpcResponse {
3033 jsonrpc: JSONRPC_VERSION.to_string(),
3034 id: response_id,
3035 result: Some(serde_json::json!({
3036 "identity": identity.as_str(),
3037 "stream_id": identity.as_str(),
3038 "subscribed": true,
3039 })),
3040 error: None,
3041 },
3042 Err(e) => identity_error_response(response_id, &e),
3043 }
3044 }
3045 "mobkit/status_identity" => {
3046 let identity_rt = match identity_ctx {
3047 Some(ctx) => &*ctx.runtime,
3048 None => return maybe_identity_not_configured(is_notification, response_id),
3049 };
3050 let identity_str = request
3051 .params
3052 .get("identity")
3053 .and_then(|v| v.as_str())
3054 .unwrap_or("");
3055 let target =
3056 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3057 {
3058 Ok(target) => target,
3059 Err(e) => {
3060 return maybe_error_response(
3061 is_notification,
3062 response_id,
3063 -32602,
3064 format!("invalid identity: {e}"),
3065 );
3066 }
3067 };
3068 let identity = target.identity.clone();
3069 if let Some(response) =
3070 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3071 {
3072 return if is_notification {
3073 String::new()
3074 } else {
3075 serialize_response(&response)
3076 };
3077 }
3078 match identity_rt.status(&identity).await {
3079 Ok(status) => {
3080 let continuity_health =
3081 serde_json::to_value(&status.continuity_health).unwrap_or(Value::Null);
3082 let result = serde_json::json!({
3083 "state": identity_lifecycle_state_json(status.state),
3084 "identity": status.identity.as_str(),
3085 "agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
3086 "session_id": status.session_id.as_ref().map(ToString::to_string),
3087 "profile": status.profile.as_ref().map(meerkat_mob::ProfileName::as_str),
3088 "addressability": addressability_json(status.addressability),
3089 "display_name": status.display_name.as_ref().map(super::identity_first::DisplayName::as_str),
3090 "labels": status.labels,
3091 "generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
3092 "checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
3093 "continuity_health": continuity_health,
3094 "lease_healthy": status.lease.as_ref().map(|lease| lease.healthy),
3095 "lease": status.lease.as_ref().map(|lease| serde_json::json!({
3096 "fencing_token": lease.fencing_token.get(),
3097 "ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
3098 "healthy": lease.healthy,
3099 })),
3100 });
3101 JsonRpcResponse {
3102 jsonrpc: JSONRPC_VERSION.to_string(),
3103 id: response_id,
3104 result: Some(result),
3105 error: None,
3106 }
3107 }
3108 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3109 if let Some(live) = target.live.as_ref() {
3110 JsonRpcResponse {
3111 jsonrpc: JSONRPC_VERSION.to_string(),
3112 id: response_id,
3113 result: Some(rpc_live_identity_status_json(live)),
3114 error: None,
3115 }
3116 } else {
3117 identity_error_response(response_id, &e)
3118 }
3119 }
3120 Err(e) => identity_error_response(response_id, &e),
3121 }
3122 }
3123 "mobkit/respawn" => {
3124 let identity_rt = match identity_ctx {
3125 Some(ctx) => &*ctx.runtime,
3126 None => return maybe_identity_not_configured(is_notification, response_id),
3127 };
3128 let identity_str = request
3129 .params
3130 .get("identity")
3131 .and_then(|v| v.as_str())
3132 .unwrap_or("");
3133 let target =
3134 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3135 {
3136 Ok(target) => target,
3137 Err(e) => {
3138 return maybe_error_response(
3139 is_notification,
3140 response_id,
3141 -32602,
3142 format!("invalid identity: {e}"),
3143 );
3144 }
3145 };
3146 let identity = target.identity.clone();
3147 if let Some(response) =
3148 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3149 {
3150 return if is_notification {
3151 String::new()
3152 } else {
3153 serialize_response(&response)
3154 };
3155 }
3156 let registered_status = match identity_rt.status(&identity).await {
3157 Ok(status) => Some(status),
3158 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => None,
3159 Err(e) => {
3160 let response = identity_error_response(response_id, &e);
3161 return if is_notification {
3162 String::new()
3163 } else {
3164 serialize_response(&response)
3165 };
3166 }
3167 };
3168 match identity_rt.respawn(&identity).await {
3169 Ok(mut record) => {
3170 let live_respawn_warning = match Box::pin(respawn_rpc_runtime_member_id(
3171 runtime,
3172 record.agent_runtime_id.as_str(),
3173 ))
3174 .await
3175 {
3176 Ok(live_result) => {
3177 let topology_restore_warning = live_result
3178 .get("topology_restore_warning")
3179 .filter(|warning| !warning.is_null())
3180 .cloned();
3181 let live_session_id =
3182 live_result.get("session_id").and_then(Value::as_str);
3183 if let Some(live_session_id) = live_session_id {
3184 match meerkat_core::types::SessionId::parse(live_session_id) {
3185 Ok(session_id) => {
3186 match identity_rt
3187 .rebind_session_after_live_respawn(
3188 &identity, session_id,
3189 )
3190 .await
3191 {
3192 Ok(updated_record) => {
3193 record = updated_record;
3194 topology_restore_warning
3195 }
3196 Err(err) => Some(serde_json::json!({
3197 "kind": "identity_rebind_failed_after_member_respawn",
3198 "message": err.to_string(),
3199 "identity": identity.as_str(),
3200 "agent_runtime_id": record.agent_runtime_id.as_str(),
3201 "live_session_id": live_session_id,
3202 })),
3203 }
3204 }
3205 Err(err) => Some(serde_json::json!({
3206 "kind": "member_respawn_session_id_invalid",
3207 "message": err.to_string(),
3208 "identity": identity.as_str(),
3209 "agent_runtime_id": record.agent_runtime_id.as_str(),
3210 "live_session_id": live_session_id,
3211 })),
3212 }
3213 } else {
3214 topology_restore_warning
3215 }
3216 }
3217 Err(err) => Some(serde_json::json!({
3218 "kind": "member_respawn_failed_after_identity_refresh",
3219 "message": err,
3220 "identity": identity.as_str(),
3221 "agent_runtime_id": record.agent_runtime_id.as_str(),
3222 })),
3223 };
3224 let cleanup_warning = if registered_status.is_some()
3225 && let Err(err) = retire_stale_rpc_members_for_identity(
3226 runtime,
3227 identity.as_str(),
3228 Some(record.agent_runtime_id.as_str()),
3229 )
3230 .await
3231 {
3232 Some(serde_json::json!({
3233 "kind": "stale_member_cleanup_failed_after_identity_respawn",
3234 "message": err,
3235 "identity": identity.as_str(),
3236 "agent_runtime_id": record.agent_runtime_id.as_str(),
3237 }))
3238 } else {
3239 None
3240 };
3241 runtime
3242 .record_console_lifecycle(
3243 identity.as_str(),
3244 "identity_respawned",
3245 serde_json::json!({
3246 "generation": record.generation.get(),
3247 "checkpoint_version": record.checkpoint_version.get(),
3248 "live_respawn_warning": live_respawn_warning.clone(),
3249 "cleanup_warning": cleanup_warning.clone(),
3250 }),
3251 )
3252 .await;
3253 JsonRpcResponse {
3254 jsonrpc: JSONRPC_VERSION.to_string(),
3255 id: response_id,
3256 result: Some(serde_json::json!({
3257 "identity": record.identity.as_str(),
3258 "agent_runtime_id": record.agent_runtime_id.as_str(),
3259 "session_id": record.session_id.to_string(),
3260 "generation": record.generation.get(),
3261 "checkpoint_version": record.checkpoint_version.get(),
3262 "live_respawn_warning": live_respawn_warning,
3263 "cleanup_warning": cleanup_warning,
3264 })),
3265 error: None,
3266 }
3267 }
3268 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3269 if let Some(live) = target.live.as_ref() {
3270 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
3271 Ok(result) => {
3272 runtime
3273 .record_console_lifecycle(
3274 live.identity.as_str(),
3275 "identity_respawned",
3276 serde_json::json!({}),
3277 )
3278 .await;
3279 JsonRpcResponse {
3280 jsonrpc: JSONRPC_VERSION.to_string(),
3281 id: response_id,
3282 result: Some(result),
3283 error: None,
3284 }
3285 }
3286 Err(err) => JsonRpcResponse {
3287 jsonrpc: JSONRPC_VERSION.to_string(),
3288 id: response_id,
3289 result: None,
3290 error: Some(JsonRpcError {
3291 code: -32000,
3292 message: format!("respawn failed: {err}"),
3293 data: None,
3294 }),
3295 },
3296 }
3297 } else {
3298 identity_error_response(response_id, &e)
3299 }
3300 }
3301 Err(e) => identity_error_response(response_id, &e),
3302 }
3303 }
3304 "mobkit/retire" => {
3305 let identity_rt = match identity_ctx {
3306 Some(ctx) => &*ctx.runtime,
3307 None => return maybe_identity_not_configured(is_notification, response_id),
3308 };
3309 let identity_str = request
3310 .params
3311 .get("identity")
3312 .and_then(|v| v.as_str())
3313 .unwrap_or("");
3314 let target =
3315 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3316 {
3317 Ok(target) => target,
3318 Err(e) => {
3319 return maybe_error_response(
3320 is_notification,
3321 response_id,
3322 -32602,
3323 format!("invalid identity: {e}"),
3324 );
3325 }
3326 };
3327 let identity = target.identity.clone();
3328 if let Some(response) =
3329 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3330 {
3331 return if is_notification {
3332 String::new()
3333 } else {
3334 serialize_response(&response)
3335 };
3336 }
3337 let registered_status = match identity_rt.status(&identity).await {
3338 Ok(status) => Some(status),
3339 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => None,
3340 Err(e) => {
3341 let response = identity_error_response(response_id, &e);
3342 return if is_notification {
3343 String::new()
3344 } else {
3345 serialize_response(&response)
3346 };
3347 }
3348 };
3349 match identity_rt.retire(&identity).await {
3350 Ok(token) => {
3351 let keep_runtime_member_id = registered_status
3352 .as_ref()
3353 .and_then(|status| status.agent_runtime_id.as_ref())
3354 .filter(|_| identity_rt.has_session_bridge())
3355 .map(crate::identity_first::AgentRuntimeId::as_str);
3356 let cleanup_warning = if registered_status.is_some()
3357 && let Err(err) = retire_stale_rpc_members_for_identity(
3358 runtime,
3359 identity.as_str(),
3360 keep_runtime_member_id,
3361 )
3362 .await
3363 {
3364 Some(serde_json::json!({
3365 "kind": "stale_member_cleanup_failed_after_identity_retire",
3366 "message": err,
3367 "identity": identity.as_str(),
3368 }))
3369 } else {
3370 None
3371 };
3372 runtime
3373 .record_console_lifecycle(
3374 identity.as_str(),
3375 "identity_retired",
3376 serde_json::json!({
3377 "fencing_token": token.get(),
3378 "cleanup_warning": cleanup_warning.clone(),
3379 }),
3380 )
3381 .await;
3382 JsonRpcResponse {
3383 jsonrpc: JSONRPC_VERSION.to_string(),
3384 id: response_id,
3385 result: Some(serde_json::json!({
3386 "fencing_token": token.get(),
3387 "cleanup_warning": cleanup_warning,
3388 })),
3389 error: None,
3390 }
3391 }
3392 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3393 if let Some(live) = target.live.as_ref() {
3394 match retire_rpc_live_identity(runtime, live).await {
3395 Ok(()) => {
3396 runtime
3397 .record_console_lifecycle(
3398 live.identity.as_str(),
3399 "identity_retired",
3400 serde_json::json!({}),
3401 )
3402 .await;
3403 JsonRpcResponse {
3404 jsonrpc: JSONRPC_VERSION.to_string(),
3405 id: response_id,
3406 result: Some(
3407 serde_json::json!({ "identity": live.identity.as_str() }),
3408 ),
3409 error: None,
3410 }
3411 }
3412 Err(err) => JsonRpcResponse {
3413 jsonrpc: JSONRPC_VERSION.to_string(),
3414 id: response_id,
3415 result: None,
3416 error: Some(JsonRpcError {
3417 code: -32000,
3418 message: format!("retire failed: {err}"),
3419 data: None,
3420 }),
3421 },
3422 }
3423 } else {
3424 identity_error_response(response_id, &e)
3425 }
3426 }
3427 Err(e) => identity_error_response(response_id, &e),
3428 }
3429 }
3430 "mobkit/reset" => {
3431 let identity_reset_ctx = match identity_ctx {
3432 Some(ctx) => ctx,
3433 None => return maybe_identity_not_configured(is_notification, response_id),
3434 };
3435 let identity_rt = &*identity_reset_ctx.runtime;
3436 let identity_str = request
3437 .params
3438 .get("identity")
3439 .and_then(|v| v.as_str())
3440 .unwrap_or("");
3441 let target =
3442 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3443 {
3444 Ok(target) => target,
3445 Err(e) => {
3446 return maybe_error_response(
3447 is_notification,
3448 response_id,
3449 -32602,
3450 format!("invalid identity: {e}"),
3451 );
3452 }
3453 };
3454 let identity = target.identity.clone();
3455 if let Some(response) =
3456 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3457 {
3458 return if is_notification {
3459 String::new()
3460 } else {
3461 serialize_response(&response)
3462 };
3463 }
3464 let _registered_status = match identity_rt.status(&identity).await {
3465 Ok(status) => {
3466 if !identity_rt.has_session_bridge() {
3467 let response = rpc_reset_requires_session_bridge_response(response_id);
3468 return if is_notification {
3469 String::new()
3470 } else {
3471 serialize_response(&response)
3472 };
3473 }
3474 status
3475 }
3476 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3477 if let Some(live) = target.live.as_ref() {
3478 let response =
3479 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
3480 Ok(result) => {
3481 runtime
3482 .record_console_lifecycle(
3483 live.identity.as_str(),
3484 "identity_reset",
3485 serde_json::json!({}),
3486 )
3487 .await;
3488 JsonRpcResponse {
3489 jsonrpc: JSONRPC_VERSION.to_string(),
3490 id: response_id,
3491 result: Some(result),
3492 error: None,
3493 }
3494 }
3495 Err(err) => JsonRpcResponse {
3496 jsonrpc: JSONRPC_VERSION.to_string(),
3497 id: response_id,
3498 result: None,
3499 error: Some(JsonRpcError {
3500 code: -32000,
3501 message: format!("reset failed: {err}"),
3502 data: None,
3503 }),
3504 },
3505 };
3506 return if is_notification {
3507 String::new()
3508 } else {
3509 serialize_response(&response)
3510 };
3511 }
3512 let response = identity_error_response(response_id, &e);
3513 return if is_notification {
3514 String::new()
3515 } else {
3516 serialize_response(&response)
3517 };
3518 }
3519 Err(e) => {
3520 let response = identity_error_response(response_id, &e);
3521 return if is_notification {
3522 String::new()
3523 } else {
3524 serialize_response(&response)
3525 };
3526 }
3527 };
3528 identity_rt.set_reset_roster_provider_context(
3529 Some(identity_reset_ctx.roster_provider.clone()),
3530 identity_reset_ctx.mob_definition.clone(),
3531 );
3532 match identity_rt.reset(&identity).await {
3533 Ok(record) => {
3534 let cleanup_warning = Some(serde_json::json!({
3535 "kind": "stale_member_cleanup_skipped_after_identity_reset",
3536 "message": "reset published the new generation without retiring stale live mob members; identity control calls reject stale runtime ids",
3537 "identity": identity.as_str(),
3538 "agent_runtime_id": record.agent_runtime_id.as_str(),
3539 }));
3540 runtime
3541 .record_console_lifecycle(
3542 identity.as_str(),
3543 "identity_reset",
3544 serde_json::json!({
3545 "generation": record.generation.get(),
3546 "checkpoint_version": record.checkpoint_version.get(),
3547 "cleanup_warning": cleanup_warning.clone(),
3548 }),
3549 )
3550 .await;
3551 JsonRpcResponse {
3552 jsonrpc: JSONRPC_VERSION.to_string(),
3553 id: response_id,
3554 result: Some(serde_json::json!({
3555 "identity": record.identity.as_str(),
3556 "agent_runtime_id": record.agent_runtime_id.as_str(),
3557 "session_id": record.session_id.to_string(),
3558 "generation": record.generation.get(),
3559 "checkpoint_version": record.checkpoint_version.get(),
3560 "cleanup_warning": cleanup_warning,
3561 })),
3562 error: None,
3563 }
3564 }
3565 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3566 if let Some(live) = target.live.as_ref() {
3567 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
3568 Ok(result) => {
3569 runtime
3570 .record_console_lifecycle(
3571 live.identity.as_str(),
3572 "identity_reset",
3573 serde_json::json!({}),
3574 )
3575 .await;
3576 JsonRpcResponse {
3577 jsonrpc: JSONRPC_VERSION.to_string(),
3578 id: response_id,
3579 result: Some(result),
3580 error: None,
3581 }
3582 }
3583 Err(err) => JsonRpcResponse {
3584 jsonrpc: JSONRPC_VERSION.to_string(),
3585 id: response_id,
3586 result: None,
3587 error: Some(JsonRpcError {
3588 code: -32000,
3589 message: format!("reset failed: {err}"),
3590 data: None,
3591 }),
3592 },
3593 }
3594 } else {
3595 identity_error_response(response_id, &e)
3596 }
3597 }
3598 Err(e) => identity_error_response(response_id, &e),
3599 }
3600 }
3601 "mobkit/delete_identity" => {
3602 let identity_rt = match identity_ctx {
3603 Some(ctx) => &*ctx.runtime,
3604 None => return maybe_identity_not_configured(is_notification, response_id),
3605 };
3606 let identity_str = request
3607 .params
3608 .get("identity")
3609 .and_then(|v| v.as_str())
3610 .unwrap_or("");
3611 let target =
3612 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3613 {
3614 Ok(target) => target,
3615 Err(e) => {
3616 return maybe_error_response(
3617 is_notification,
3618 response_id,
3619 -32602,
3620 format!("invalid identity: {e}"),
3621 );
3622 }
3623 };
3624 let identity = target.identity.clone();
3625 if let Some(response) =
3626 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3627 {
3628 return if is_notification {
3629 String::new()
3630 } else {
3631 serialize_response(&response)
3632 };
3633 }
3634 let registered_status = match identity_rt.status(&identity).await {
3635 Ok(status) => status,
3636 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3637 if target.live.is_some() {
3638 let response = JsonRpcResponse {
3639 jsonrpc: JSONRPC_VERSION.to_string(),
3640 id: response_id,
3641 result: None,
3642 error: Some(JsonRpcError {
3643 code: -32602,
3644 message: format!(
3645 "delete_identity requires durable identity: {} is live-only",
3646 identity.as_str()
3647 ),
3648 data: Some(serde_json::json!({
3649 "kind": "live_only_identity_delete_unsupported",
3650 "identity": identity.as_str(),
3651 })),
3652 }),
3653 };
3654 return if is_notification {
3655 String::new()
3656 } else {
3657 serialize_response(&response)
3658 };
3659 }
3660 let response = identity_error_response(response_id, &e);
3661 return if is_notification {
3662 String::new()
3663 } else {
3664 serialize_response(&response)
3665 };
3666 }
3667 Err(e) => {
3668 let response = identity_error_response(response_id, &e);
3669 return if is_notification {
3670 String::new()
3671 } else {
3672 serialize_response(&response)
3673 };
3674 }
3675 };
3676 let keep_runtime_member_id = registered_status
3677 .agent_runtime_id
3678 .as_ref()
3679 .filter(|_| identity_rt.has_session_bridge())
3680 .map(crate::identity_first::AgentRuntimeId::as_str);
3681 match identity_rt.delete_identity(&identity).await {
3682 Ok(()) => {
3683 let cleanup_warning = if let Err(err) = retire_stale_rpc_members_for_identity(
3684 runtime,
3685 identity.as_str(),
3686 keep_runtime_member_id,
3687 )
3688 .await
3689 {
3690 Some(serde_json::json!({
3691 "kind": "stale_member_cleanup_failed_after_identity_delete",
3692 "identity": identity.as_str(),
3693 "message": err,
3694 }))
3695 } else {
3696 None
3697 };
3698 runtime
3699 .record_console_lifecycle(
3700 identity.as_str(),
3701 "identity_deleted",
3702 serde_json::json!({
3703 "cleanup_warning": cleanup_warning,
3704 }),
3705 )
3706 .await;
3707 JsonRpcResponse {
3708 jsonrpc: JSONRPC_VERSION.to_string(),
3709 id: response_id,
3710 result: Some(serde_json::json!({
3711 "identity": identity.as_str(),
3712 "cleanup_warning": cleanup_warning,
3713 })),
3714 error: None,
3715 }
3716 }
3717 Err(e) => identity_error_response(response_id, &e),
3718 }
3719 }
3720 "mobkit/inspect_identity" => {
3721 let identity_rt = match identity_ctx {
3722 Some(ctx) => &*ctx.runtime,
3723 None => return maybe_identity_not_configured(is_notification, response_id),
3724 };
3725 let identity_str = request
3726 .params
3727 .get("identity")
3728 .and_then(|v| v.as_str())
3729 .unwrap_or("");
3730 let target =
3731 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3732 {
3733 Ok(target) => target,
3734 Err(e) => {
3735 return maybe_error_response(
3736 is_notification,
3737 response_id,
3738 -32602,
3739 format!("invalid identity: {e}"),
3740 );
3741 }
3742 };
3743 let identity = target.identity.clone();
3744 let status = identity_rt.status(&identity).await;
3745 if let Some(response) =
3746 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3747 {
3748 return if is_notification {
3749 String::new()
3750 } else {
3751 serialize_response(&response)
3752 };
3753 }
3754 match identity_rt.inspect(&identity).await {
3755 Ok(inspection) => {
3756 let status = status.ok();
3757 JsonRpcResponse {
3758 jsonrpc: JSONRPC_VERSION.to_string(),
3759 id: response_id,
3760 result: Some(serde_json::json!({
3761 "identity": identity.as_str(),
3762 "state": status.as_ref().map(|status| identity_lifecycle_state_json(status.state)),
3763 "profile": status.as_ref().and_then(|status| status.profile.as_ref().map(meerkat_mob::ProfileName::as_str)),
3764 "addressability": status.as_ref().map(|status| addressability_json(status.addressability)),
3765 "display_name": status.as_ref().and_then(|status| status.display_name.as_ref().map(super::identity_first::DisplayName::as_str)),
3766 "labels": status.as_ref().map(|status| status.labels.clone()).unwrap_or_default(),
3767 "generation": status.as_ref().and_then(|status| status.generation.map(super::identity_first::ContinuityGeneration::get)),
3768 "checkpoint_version": status.as_ref().and_then(|status| status.checkpoint_version.map(super::identity_first::CheckpointVersion::get)),
3769 "continuity_health": status.as_ref().and_then(|status| serde_json::to_value(&status.continuity_health).ok()).unwrap_or(Value::Null),
3770 "lease_healthy": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| lease.healthy)),
3771 "continuity": status.as_ref().map(|status| serde_json::json!({
3772 "generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
3773 "checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
3774 "session_id": status.session_id.as_ref().map(ToString::to_string),
3775 "agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
3776 })).unwrap_or_else(|| serde_json::json!({})),
3777 "lease": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| serde_json::json!({
3778 "fencing_token": lease.fencing_token.get(),
3779 "ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
3780 "healthy": lease.healthy,
3781 }))),
3782 "output_preview": inspection.output_preview,
3783 "is_final": inspection.is_final,
3784 "peer_reachable_count": inspection.peer_reachable_count,
3785 })),
3786 error: None,
3787 }
3788 }
3789 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3790 if let Some(live) = target.live.as_ref() {
3791 JsonRpcResponse {
3792 jsonrpc: JSONRPC_VERSION.to_string(),
3793 id: response_id,
3794 result: Some(rpc_live_identity_inspect_json(runtime, live).await),
3795 error: None,
3796 }
3797 } else {
3798 identity_error_response(response_id, &e)
3799 }
3800 }
3801 Err(e) => identity_error_response(response_id, &e),
3802 }
3803 }
3804 "mobkit/reconcile_identity" => {
3805 let ctx = match identity_ctx {
3806 Some(ctx) => ctx,
3807 None => return maybe_identity_not_configured(is_notification, response_id),
3808 };
3809 let roster_specs = match ctx
3811 .roster_provider
3812 .roster(&crate::identity_first::RosterContext {
3813 mob_definition: ctx.mob_definition.clone(),
3814 previous_identities: Vec::new(),
3815 })
3816 .await
3817 {
3818 Ok(specs) => specs,
3819 Err(e) => {
3820 return maybe_error_response(
3821 is_notification,
3822 response_id,
3823 -32603,
3824 format!("roster provider failed: {e}"),
3825 );
3826 }
3827 };
3828 match crate::identity_first::restore_flow(
3829 &ctx.runtime,
3830 &roster_specs,
3831 ctx.topology_provider.as_deref(),
3832 ctx.customizer.as_deref(),
3833 )
3834 .await
3835 {
3836 Ok(result) => {
3837 let outcomes: serde_json::Map<String, Value> = result
3838 .outcomes
3839 .iter()
3840 .map(|(id, outcome)| {
3841 let val = match outcome {
3842 crate::identity_first::RestoreOutcome::Created {
3843 record, ..
3844 } => {
3845 serde_json::json!({
3846 "outcome": "created",
3847 "identity": record.identity.as_str(),
3848 "agent_runtime_id": record.agent_runtime_id.as_str(),
3849 "session_id": record.session_id.to_string(),
3850 "generation": record.generation.get(),
3851 })
3852 }
3853 crate::identity_first::RestoreOutcome::Dormant {
3854 record, ..
3855 } => {
3856 serde_json::json!({
3857 "outcome": "dormant",
3858 "identity": id.as_str(),
3859 "agent_runtime_id": record.as_ref().map(|record| record.agent_runtime_id.as_str()),
3860 "session_id": record.as_ref().map(|record| record.session_id.to_string()),
3861 "generation": record.as_ref().map(|record| record.generation.get()),
3862 })
3863 }
3864 crate::identity_first::RestoreOutcome::Resumed {
3865 record, ..
3866 } => {
3867 serde_json::json!({
3868 "outcome": "resumed",
3869 "identity": record.identity.as_str(),
3870 "agent_runtime_id": record.agent_runtime_id.as_str(),
3871 "session_id": record.session_id.to_string(),
3872 "generation": record.generation.get(),
3873 })
3874 }
3875 crate::identity_first::RestoreOutcome::Broken(failure) => {
3876 serde_json::json!({
3877 "outcome": "broken",
3878 "identity": failure.identity.as_str(),
3879 "detail": failure.detail,
3880 })
3881 }
3882 };
3883 (id.to_string(), val)
3884 })
3885 .collect();
3886 JsonRpcResponse {
3887 jsonrpc: JSONRPC_VERSION.to_string(),
3888 id: response_id,
3889 result: Some(serde_json::json!({
3890 "outcomes": outcomes,
3891 "managed_edges": result.managed_edges.len(),
3892 })),
3893 error: None,
3894 }
3895 }
3896 Err(e) => identity_error_response(response_id, &e),
3897 }
3898 }
3899 method if method.contains('/') && !method.starts_with("mobkit/") => {
3900 let module_id = method
3901 .split('/')
3902 .next()
3903 .map(ToString::to_string)
3904 .unwrap_or_default();
3905 let route = runtime
3906 .route_module_call(
3907 &ModuleRouteRequest {
3908 module_id: module_id.clone(),
3909 method: method.to_string(),
3910 params: request.params,
3911 },
3912 timeout,
3913 )
3914 .await;
3915 match route {
3916 Ok(response) => JsonRpcResponse {
3917 jsonrpc: JSONRPC_VERSION.to_string(),
3918 id: response_id,
3919 result: Some(serde_json::json!({
3920 "module_id": response.module_id,
3921 "method": response.method,
3922 "payload": response.payload
3923 })),
3924 error: None,
3925 },
3926 Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
3927 jsonrpc: JSONRPC_VERSION.to_string(),
3928 id: response_id,
3929 result: None,
3930 error: Some(JsonRpcError {
3931 code: -32601,
3932 message: format!("Module '{module_id}' not loaded"),
3933 data: None,
3934 }),
3935 },
3936 Err(err) => JsonRpcResponse {
3937 jsonrpc: JSONRPC_VERSION.to_string(),
3938 id: response_id,
3939 result: None,
3940 error: Some(JsonRpcError {
3941 code: -32000,
3942 message: format!("Module route failed: {err:?}"),
3943 data: None,
3944 }),
3945 },
3946 }
3947 }
3948 method if method.starts_with("mobkit/live/") => match live {
3949 None => crate::live_wiring::live_unavailable_response(response_id),
3950 Some(live) => {
3951 let resolved_session = resolve_live_target(runtime, &request.params).await;
3958 live(
3959 resolved_session,
3960 method.to_string(),
3961 request.params.clone(),
3962 response_id,
3963 )
3964 .await
3965 }
3966 },
3967 method if workgraph_methods::is_workgraph_method(method) => {
3968 let service = runtime.workgraph_service();
3969 let admission = runtime.workgraph_admission();
3970 match workgraph_methods::handle_workgraph_method(
3973 service.as_ref(),
3974 &admission,
3975 workgraph_methods::WorkgraphSurface::HostStdin,
3976 method,
3977 &request.params,
3978 )
3979 .await
3980 {
3981 Ok(result) => JsonRpcResponse {
3982 jsonrpc: JSONRPC_VERSION.to_string(),
3983 id: response_id,
3984 result: Some(result),
3985 error: None,
3986 },
3987 Err(error) => JsonRpcResponse {
3988 jsonrpc: JSONRPC_VERSION.to_string(),
3989 id: response_id,
3990 result: None,
3991 error: Some(error),
3992 },
3993 }
3994 }
3995 _ => JsonRpcResponse {
3996 jsonrpc: JSONRPC_VERSION.to_string(),
3997 id: response_id,
3998 result: None,
3999 error: Some(JsonRpcError {
4000 code: -32601,
4001 message: "Method not found".to_string(),
4002 data: None,
4003 }),
4004 },
4005 };
4006 if is_notification {
4007 String::new()
4008 } else {
4009 serialize_response(&response)
4010 }
4011}
4012
4013fn build_models_catalog_result() -> Value {
4014 let entries: Vec<Value> = meerkat_models::catalog()
4015 .iter()
4016 .filter_map(|e| {
4017 let mut val = serde_json::to_value(e).ok()?;
4018 if let Some(provider) = meerkat_core::Provider::parse_strict(e.provider)
4019 && let Some(profile) = meerkat_models::profile_for(provider, e.id)
4020 && let Ok(p) = serde_json::to_value(&profile)
4021 {
4022 val["profile"] = p;
4023 }
4024 Some(val)
4025 })
4026 .collect();
4027 let defaults: Vec<Value> = meerkat_models::provider_defaults()
4028 .iter()
4029 .filter_map(|d| serde_json::to_value(d).ok())
4030 .collect();
4031 serde_json::json!({
4032 "models": entries,
4033 "provider_defaults": defaults,
4034 })
4035}
4036
4037#[derive(Debug, Clone)]
4038struct RpcLiveIdentityAlias {
4039 identity: crate::identity_first::AgentIdentity,
4040 runtime_member_id: String,
4041 member: meerkat_mob::runtime::MobMemberListEntry,
4042 session_id: Option<String>,
4043}
4044
4045#[derive(Debug, Clone)]
4046struct RpcIdentityControlTarget {
4047 identity: crate::identity_first::AgentIdentity,
4048 live: Option<RpcLiveIdentityAlias>,
4049}
4050
4051fn rpc_member_durable_identity(member: &meerkat_mob::runtime::MobMemberListEntry) -> String {
4052 member
4053 .labels
4054 .get("agent_identity")
4055 .filter(|value| !value.trim().is_empty())
4056 .cloned()
4057 .unwrap_or_else(|| {
4060 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
4061 })
4062}
4063
4064async fn resolve_rpc_live_identity_alias(
4065 runtime: &UnifiedRuntime,
4066 requested_identity: &str,
4067) -> Result<Option<RpcLiveIdentityAlias>, String> {
4068 let matches = resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
4069 if matches.len() > 1 {
4070 return Err(format!(
4071 "ambiguous live identity alias {requested_identity}: candidates [{}]",
4072 matches
4073 .iter()
4074 .map(|entry| entry.runtime_member_id.clone())
4075 .collect::<Vec<_>>()
4076 .join(", ")
4077 ));
4078 }
4079 Ok(matches.into_iter().next())
4080}
4081
4082async fn resolve_rpc_live_runtime_member_alias(
4083 runtime: &UnifiedRuntime,
4084 runtime_member_id: &str,
4085) -> Result<Option<RpcLiveIdentityAlias>, String> {
4086 let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
4087 let handle = runtime.mob_handle();
4088 let Some(member) = handle
4089 .list_members_including_retiring()
4090 .await
4091 .into_iter()
4092 .find(|entry| entry.agent_identity == requested_member_id)
4093 else {
4094 return Ok(None);
4095 };
4096 if !rpc_live_identity_alias_member_visible(&member) {
4097 return Ok(None);
4098 }
4099 let durable_identity = rpc_member_durable_identity(&member);
4100 let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
4101 .map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
4102 let session_id = handle
4103 .resolve_bridge_session_id_observation(&member.agent_identity)
4104 .await
4105 .map(|session_id| session_id.to_string());
4106 Ok(Some(RpcLiveIdentityAlias {
4107 identity,
4108 runtime_member_id: crate::member_comms_id::runtime_alias_str(
4109 member.agent_identity.as_str(),
4110 )
4111 .into_owned(),
4112 member,
4113 session_id,
4114 }))
4115}
4116
4117async fn rpc_runtime_member_alias_exists_hidden(
4118 runtime: &UnifiedRuntime,
4119 runtime_member_id: &str,
4120) -> bool {
4121 let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
4122 runtime
4123 .mob_handle()
4124 .list_members_including_retiring()
4125 .await
4126 .into_iter()
4127 .find(|entry| entry.agent_identity == requested_member_id)
4128 .is_some_and(|member| !rpc_live_identity_alias_member_visible(&member))
4129}
4130
4131async fn rpc_live_identity_alias_exists_hidden(
4132 runtime: &UnifiedRuntime,
4133 requested_identity: &str,
4134) -> bool {
4135 let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
4136 runtime
4137 .mob_handle()
4138 .list_members_including_retiring()
4139 .await
4140 .into_iter()
4141 .any(|member| {
4142 (member.agent_identity == requested_member_id
4143 || member
4144 .labels
4145 .get("agent_identity")
4146 .is_some_and(|identity| identity == requested_identity))
4147 && !rpc_live_identity_alias_member_visible(&member)
4148 })
4149}
4150
4151async fn resolve_rpc_live_identity_alias_candidates(
4152 runtime: &UnifiedRuntime,
4153 requested_identity: &str,
4154) -> Result<Vec<RpcLiveIdentityAlias>, String> {
4155 let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
4156 let handle = runtime.mob_handle();
4157 let members = handle.list_members_including_retiring().await;
4158 let exact_matches = members
4159 .iter()
4160 .filter(|entry| entry.agent_identity == requested_member_id)
4161 .cloned()
4162 .collect::<Vec<_>>();
4163 let label_matches = members
4164 .iter()
4165 .filter(|entry| {
4166 entry
4167 .labels
4168 .get("agent_identity")
4169 .is_some_and(|identity| identity == requested_identity)
4170 })
4171 .cloned()
4172 .collect::<Vec<_>>();
4173 let mut matches = exact_matches;
4174 matches.extend(label_matches);
4175 let mut seen_member_ids = BTreeSet::new();
4176 matches.retain(|entry| seen_member_ids.insert(entry.agent_identity.to_string()));
4177 let mut aliases = Vec::with_capacity(matches.len());
4178 for member in matches {
4179 if !rpc_live_identity_alias_member_visible(&member) {
4180 continue;
4181 }
4182 let durable_identity = rpc_member_durable_identity(&member);
4183 let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
4184 .map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
4185 let session_id = handle
4186 .resolve_bridge_session_id_observation(&member.agent_identity)
4187 .await
4188 .map(|session_id| session_id.to_string());
4189 aliases.push(RpcLiveIdentityAlias {
4190 identity,
4191 runtime_member_id: crate::member_comms_id::runtime_alias_str(
4192 member.agent_identity.as_str(),
4193 )
4194 .into_owned(),
4195 member,
4196 session_id,
4197 });
4198 }
4199 Ok(aliases)
4200}
4201
4202fn rpc_live_identity_alias_member_visible(
4203 member: &meerkat_mob::runtime::MobMemberListEntry,
4204) -> bool {
4205 rpc_live_identity_alias_visible(member.role.as_str(), &member.labels)
4206}
4207
4208fn rpc_live_identity_alias_visible(
4209 member_role: &str,
4210 labels: &std::collections::BTreeMap<String, String>,
4211) -> bool {
4212 let projected_role = labels
4213 .get("role")
4214 .map(String::as_str)
4215 .unwrap_or(member_role);
4216 !is_implicit_delegate_member(member_role, labels)
4217 && !is_implicit_delegate_member(projected_role, labels)
4218}
4219
4220async fn resolve_rpc_identity_control_target(
4221 runtime: &UnifiedRuntime,
4222 identity_rt: &crate::identity_first::IdentityRuntime,
4223 requested_identity: &str,
4224) -> Result<RpcIdentityControlTarget, String> {
4225 if requested_identity.starts_with("rt:") {
4226 for status in identity_rt.statuses().await {
4227 if status
4228 .agent_runtime_id
4229 .as_ref()
4230 .is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
4231 {
4232 let identity = status.identity;
4233 let registered_live =
4234 resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
4235 if let Some(registered) = registered_live {
4236 return Ok(RpcIdentityControlTarget {
4237 identity,
4238 live: Some(registered),
4239 });
4240 }
4241 if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
4242 return Err(format!("identity hidden by policy: {requested_identity}"));
4243 }
4244 let durable_live_candidates =
4245 resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
4246 let durable_live = if durable_live_candidates.len() > 1 {
4247 return Err(format!(
4248 "ambiguous live identity alias {}: candidates [{}]",
4249 identity.as_str(),
4250 durable_live_candidates
4251 .iter()
4252 .map(|alias| alias.runtime_member_id.clone())
4253 .collect::<Vec<_>>()
4254 .join(", ")
4255 ));
4256 } else {
4257 durable_live_candidates.into_iter().next()
4258 };
4259 return Ok(RpcIdentityControlTarget {
4260 identity,
4261 live: durable_live,
4262 });
4263 }
4264 }
4265 let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
4266 if let Some(live_alias) = live {
4267 if let Ok(status) = identity_rt.status(&live_alias.identity).await
4268 && !rpc_live_alias_matches_status_runtime(Some(&live_alias), &status)
4269 {
4270 return Ok(RpcIdentityControlTarget {
4271 identity: live_alias.identity.clone(),
4272 live: Some(live_alias),
4273 });
4274 }
4275 let live_identity_candidates =
4276 resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
4277 .await?;
4278 if live_identity_candidates.len() > 1 {
4279 return Err(format!(
4280 "ambiguous live identity alias {}: candidates [{}]",
4281 live_alias.identity.as_str(),
4282 live_identity_candidates
4283 .iter()
4284 .map(|alias| alias.runtime_member_id.clone())
4285 .collect::<Vec<_>>()
4286 .join(", ")
4287 ));
4288 }
4289 return Ok(RpcIdentityControlTarget {
4290 identity: live_alias.identity.clone(),
4291 live: Some(live_alias),
4292 });
4293 }
4294 if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
4295 return Err(format!("identity hidden by policy: {requested_identity}"));
4296 }
4297 return Err(format!("runtime identity not found: {requested_identity}"));
4298 }
4299 if let Ok(identity) = crate::identity_first::AgentIdentity::parse(requested_identity) {
4300 match identity_rt.status(&identity).await {
4301 Ok(status) => {
4302 let registered_live = match status.agent_runtime_id.as_ref() {
4303 Some(runtime_id) => {
4304 resolve_rpc_live_runtime_member_alias(runtime, runtime_id.as_str()).await?
4305 }
4306 None => None,
4307 };
4308 if let Some(registered) = registered_live {
4309 return Ok(RpcIdentityControlTarget {
4310 identity,
4311 live: Some(registered),
4312 });
4313 }
4314 if let Some(runtime_id) = status.agent_runtime_id.as_ref()
4315 && rpc_runtime_member_alias_exists_hidden(runtime, runtime_id.as_str()).await
4316 {
4317 return Err(format!("identity hidden by policy: {requested_identity}"));
4318 }
4319 let requested_live_candidates =
4320 resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
4321 let requested_live = if requested_live_candidates.len() > 1 {
4322 return Err(format!(
4323 "ambiguous live identity alias {requested_identity}: candidates [{}]",
4324 requested_live_candidates
4325 .iter()
4326 .map(|alias| alias.runtime_member_id.clone())
4327 .collect::<Vec<_>>()
4328 .join(", ")
4329 ));
4330 } else {
4331 requested_live_candidates.into_iter().next()
4332 };
4333 return Ok(RpcIdentityControlTarget {
4334 identity,
4335 live: requested_live,
4336 });
4337 }
4338 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {}
4339 Err(err) => return Err(err.to_string()),
4340 }
4341 }
4342 for status in identity_rt.statuses().await {
4343 if status
4344 .agent_runtime_id
4345 .as_ref()
4346 .is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
4347 {
4348 let identity = status.identity;
4349 let registered_live =
4350 resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
4351 let durable_live_candidates =
4352 resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
4353 let durable_live = if durable_live_candidates.len() > 1 {
4354 return Err(format!(
4355 "ambiguous live identity alias {}: candidates [{}]",
4356 identity.as_str(),
4357 durable_live_candidates
4358 .iter()
4359 .map(|alias| alias.runtime_member_id.clone())
4360 .collect::<Vec<_>>()
4361 .join(", ")
4362 ));
4363 } else {
4364 durable_live_candidates.into_iter().next()
4365 };
4366 let live = match (registered_live, durable_live) {
4367 (Some(registered), Some(durable))
4368 if registered.runtime_member_id == durable.runtime_member_id =>
4369 {
4370 Some(registered)
4371 }
4372 (Some(registered), None) => Some(registered),
4373 (Some(_registered), Some(durable)) => Some(durable),
4374 (None, durable) => durable,
4375 };
4376 return Ok(RpcIdentityControlTarget { identity, live });
4377 }
4378 }
4379 let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
4380 if let Some(live_alias) = live {
4381 if let Some(bound_status) = identity_rt.statuses().await.into_iter().find(|status| {
4382 status
4383 .agent_runtime_id
4384 .as_ref()
4385 .is_some_and(|runtime_id| runtime_id.as_str() == live_alias.runtime_member_id)
4386 }) && bound_status.identity != live_alias.identity
4387 {
4388 return Err(format!(
4389 "stale live identity alias: live console alias {} resolves to {}, but identity runtime binding belongs to {}",
4390 live_alias.identity.as_str(),
4391 live_alias.runtime_member_id,
4392 bound_status.identity.as_str(),
4393 ));
4394 }
4395 let live_identity_candidates =
4396 resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
4397 .await?;
4398 if live_identity_candidates.len() > 1 {
4399 return Err(format!(
4400 "ambiguous live identity alias {}: candidates [{}]",
4401 live_alias.identity.as_str(),
4402 live_identity_candidates
4403 .iter()
4404 .map(|alias| alias.runtime_member_id.clone())
4405 .collect::<Vec<_>>()
4406 .join(", ")
4407 ));
4408 }
4409 return Ok(RpcIdentityControlTarget {
4410 identity: live_alias.identity.clone(),
4411 live: Some(live_alias),
4412 });
4413 }
4414 if rpc_live_identity_alias_exists_hidden(runtime, requested_identity).await {
4415 return Err(format!("identity hidden by policy: {requested_identity}"));
4416 }
4417 let identity = crate::identity_first::AgentIdentity::parse(requested_identity)
4418 .map_err(|err| err.to_string())?;
4419 Ok(RpcIdentityControlTarget {
4420 identity,
4421 live: None,
4422 })
4423}
4424
4425fn rpc_reset_requires_session_bridge_response(response_id: Value) -> JsonRpcResponse {
4426 JsonRpcResponse {
4427 jsonrpc: JSONRPC_VERSION.to_string(),
4428 id: response_id,
4429 result: None,
4430 error: Some(JsonRpcError {
4431 code: -32602,
4432 message: "reset requires an identity runtime with a session bridge".to_string(),
4433 data: Some(serde_json::json!({
4434 "kind": "identity_reset_requires_session_bridge",
4435 })),
4436 }),
4437 }
4438}
4439
4440fn rpc_live_alias_matches_status_runtime(
4441 alias: Option<&RpcLiveIdentityAlias>,
4442 status: &crate::identity_first::IdentityStatus,
4443) -> bool {
4444 let Some(alias) = alias else {
4445 return true;
4446 };
4447 let session_matches = match (
4448 status.session_id.as_ref().map(ToString::to_string),
4449 alias.session_id.as_deref(),
4450 ) {
4451 (Some(status_session), Some(live_session)) => status_session == live_session,
4452 _ => true,
4453 };
4454 status
4455 .agent_runtime_id
4456 .as_ref()
4457 .is_some_and(|runtime_id| runtime_id.as_str() == alias.runtime_member_id)
4458 && alias.identity == status.identity
4459 && session_matches
4460}
4461
4462async fn rpc_stale_live_alias_error_response(
4463 identity_rt: &crate::identity_first::IdentityRuntime,
4464 target: &RpcIdentityControlTarget,
4465 response_id: Value,
4466) -> Option<JsonRpcResponse> {
4467 let live = target.live.as_ref()?;
4468 let Ok(status) = identity_rt.status(&target.identity).await else {
4469 return None;
4470 };
4471 if rpc_live_alias_matches_status_runtime(Some(live), &status) {
4472 return None;
4473 }
4474 Some(JsonRpcResponse {
4475 jsonrpc: JSONRPC_VERSION.to_string(),
4476 id: response_id,
4477 result: None,
4478 error: Some(JsonRpcError {
4479 code: -32000,
4480 message: format!(
4481 "identity runtime binding for {} points at {}, but requested live member is {}",
4482 target.identity.as_str(),
4483 status
4484 .agent_runtime_id
4485 .as_ref()
4486 .map(crate::identity_first::AgentRuntimeId::as_str)
4487 .unwrap_or("<none>"),
4488 live.runtime_member_id
4489 ),
4490 data: Some(serde_json::json!({
4491 "kind": "stale_identity_runtime_binding",
4492 "identity": target.identity.as_str(),
4493 "registered_runtime_member_id": status.agent_runtime_id.as_ref().map(crate::identity_first::AgentRuntimeId::as_str),
4494 "live_runtime_member_id": live.runtime_member_id,
4495 "registered_session_id": status.session_id.as_ref().map(ToString::to_string),
4496 "live_session_id": live.session_id,
4497 })),
4498 }),
4499 })
4500}
4501
4502fn rpc_member_is_addressable(member: &meerkat_mob::runtime::MobMemberListEntry) -> bool {
4503 member
4504 .labels
4505 .get("addressable")
4506 .map(|value| !value.eq_ignore_ascii_case("false"))
4507 .unwrap_or(true)
4508}
4509
4510fn rpc_live_identity_status_json(alias: &RpcLiveIdentityAlias) -> Value {
4511 serde_json::json!({
4512 "state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
4513 "identity": alias.identity.as_str(),
4514 "agent_runtime_id": alias.runtime_member_id,
4515 "session_id": alias.session_id,
4516 "profile": alias.member.role.to_string(),
4517 "addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
4518 "display_name": alias.member.labels.get("display_name"),
4519 "labels": alias.member.labels,
4520 "generation": Value::Null,
4521 "checkpoint_version": Value::Null,
4522 "continuity_health": Value::Null,
4523 "lease_healthy": Value::Null,
4524 "lease": Value::Null,
4525 })
4526}
4527
4528async fn rpc_live_identity_inspect_json(
4529 runtime: &UnifiedRuntime,
4530 alias: &RpcLiveIdentityAlias,
4531) -> Value {
4532 let snapshot = runtime
4533 .mob_handle()
4534 .member_status(&crate::member_comms_id::mob_member_id(
4535 alias.runtime_member_id.as_str(),
4536 ))
4537 .await
4538 .ok();
4539 serde_json::json!({
4540 "identity": alias.identity.as_str(),
4541 "state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
4542 "profile": alias.member.role.to_string(),
4543 "addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
4544 "display_name": alias.member.labels.get("display_name"),
4545 "labels": alias.member.labels,
4546 "generation": Value::Null,
4547 "checkpoint_version": Value::Null,
4548 "continuity_health": Value::Null,
4549 "lease_healthy": Value::Null,
4550 "continuity": {
4551 "generation": Value::Null,
4552 "checkpoint_version": Value::Null,
4553 "session_id": alias.session_id,
4554 "agent_runtime_id": alias.runtime_member_id,
4555 },
4556 "lease": Value::Null,
4557 "output_preview": snapshot.as_ref().and_then(|snapshot| snapshot.output_preview.clone()),
4558 "is_final": snapshot.as_ref().map(|snapshot| snapshot.is_final).unwrap_or(false),
4559 "peer_reachable_count": alias.member.wired_to.len(),
4560 "progress": snapshot.as_ref().and_then(|snapshot| snapshot.progress.clone()),
4563 })
4564}
4565
4566async fn retire_rpc_live_identity(
4567 runtime: &UnifiedRuntime,
4568 alias: &RpcLiveIdentityAlias,
4569) -> Result<(), String> {
4570 retire_rpc_runtime_member_id(runtime, alias.runtime_member_id.as_str()).await
4571}
4572
4573async fn retire_rpc_runtime_member_id(
4574 runtime: &UnifiedRuntime,
4575 runtime_member_id: &str,
4576) -> Result<(), String> {
4577 match runtime
4578 .mob_handle()
4579 .retire(crate::member_comms_id::mob_member_id(runtime_member_id))
4580 .await
4581 {
4582 Ok(()) => Ok(()),
4583 Err(err) if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) => Ok(()),
4584 Err(err) => Err(err.to_string()),
4585 }
4586}
4587
4588fn rpc_member_id_matches_durable_identity(member_id: &str, durable_identity: &str) -> bool {
4589 crate::member_comms_id::runtime_alias_str(member_id) == durable_identity
4592}
4593
4594async fn retire_stale_rpc_members_for_identity(
4595 runtime: &UnifiedRuntime,
4596 durable_identity: &str,
4597 keep_runtime_member_id: Option<&str>,
4598) -> Result<(), String> {
4599 let stale_members = runtime
4600 .mob_handle()
4601 .list_members_including_retiring()
4602 .await
4603 .into_iter()
4604 .filter(|member| {
4605 if !rpc_live_identity_alias_member_visible(member) {
4606 return false;
4607 }
4608 (rpc_member_id_matches_durable_identity(
4609 member.agent_identity.as_str(),
4610 durable_identity,
4611 ) || member
4612 .labels
4613 .get("agent_identity")
4614 .is_some_and(|identity| identity == durable_identity))
4615 && keep_runtime_member_id
4616 .map(|keep| {
4617 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str())
4619 != keep
4620 })
4621 .unwrap_or(true)
4622 })
4623 .map(|member| {
4625 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
4626 })
4627 .collect::<Vec<_>>();
4628 for member_id in stale_members {
4629 retire_rpc_runtime_member_id(runtime, &member_id).await?;
4630 }
4631 Ok(())
4632}
4633
4634async fn respawn_rpc_live_identity(
4635 runtime: &UnifiedRuntime,
4636 alias: &RpcLiveIdentityAlias,
4637) -> Result<Value, String> {
4638 let mut result = Box::pin(respawn_rpc_runtime_member_id(
4639 runtime,
4640 alias.runtime_member_id.as_str(),
4641 ))
4642 .await?;
4643 result["identity"] = serde_json::json!(alias.identity.as_str());
4644 Ok(result)
4645}
4646
4647async fn respawn_rpc_runtime_member_id(
4648 runtime: &UnifiedRuntime,
4649 runtime_member_id: &str,
4650) -> Result<Value, String> {
4651 let handle = runtime.mob_handle();
4652 let member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
4653 let entry_before_respawn = handle.get_member(&member_id).await.ok().flatten();
4656 let mut topology_restore_warning = None;
4657 match handle.respawn(member_id.clone(), None).await {
4658 Ok(_receipt) => {}
4659 Err(err) => {
4660 if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(&err) {
4661 tracing::warn!(
4662 member_id = %member_id,
4663 failed_peer_count = failed_peer_ids.len(),
4664 failed_peer_ids = ?failed_peer_ids,
4665 "rpc member respawn restored member with isolated peer edges; continuing degraded respawn"
4666 );
4667 topology_restore_warning = Some(topology_restore_warning_json(&failed_peer_ids));
4668 } else if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) {
4669 if handle
4672 .get_member(&member_id)
4673 .await
4674 .map_err(|lookup_err| lookup_err.to_string())?
4675 .is_none()
4676 && let Some(entry) = entry_before_respawn
4677 {
4678 let mut spec =
4679 meerkat_mob::SpawnMemberSpec::new(entry.role.clone(), member_id.clone());
4680 if !entry.labels.is_empty() {
4681 spec = spec.with_labels(entry.labels.clone());
4682 }
4683 handle
4684 .ensure_member(spec)
4685 .await
4686 .map_err(|ensure_err| ensure_err.to_string())?;
4687 }
4688 } else {
4689 return Err(err.to_string());
4690 }
4691 }
4692 }
4693 let session_id = handle
4694 .resolve_bridge_session_id_observation(&member_id)
4695 .await
4696 .map(|session_id| session_id.to_string());
4697 Ok(serde_json::json!({
4698 "agent_runtime_id": runtime_member_id,
4699 "session_id": session_id,
4700 "generation": Value::Null,
4701 "checkpoint_version": Value::Null,
4702 "topology_restore_warning": topology_restore_warning,
4703 }))
4704}
4705
4706fn identity_not_configured(response_id: Value) -> String {
4707 error_response(response_id, -32601, "identity-first runtime not configured")
4708}
4709
4710fn maybe_identity_not_configured(is_notification: bool, response_id: Value) -> String {
4711 if is_notification {
4712 String::new()
4713 } else {
4714 identity_not_configured(response_id)
4715 }
4716}
4717
4718fn addressability_json(addressability: crate::identity_first::AgentAddressability) -> &'static str {
4719 match addressability {
4720 crate::identity_first::AgentAddressability::Addressable => "addressable",
4721 crate::identity_first::AgentAddressability::InternalOnly => "internal_only",
4722 }
4723}
4724
4725fn identity_lifecycle_state_json(
4728 state: crate::identity_first::IdentityLifecycleState,
4729) -> &'static str {
4730 state.wire_str()
4731}
4732
4733fn identity_error_response(
4734 response_id: Value,
4735 err: &crate::identity_first::IdentityRuntimeError,
4736) -> JsonRpcResponse {
4737 use crate::identity_first::IdentityRuntimeError;
4738 let (code, message) = match err {
4739 IdentityRuntimeError::UnknownIdentity(id) => (-32001, format!("unknown identity: {id}")),
4740 IdentityRuntimeError::NotAddressable(na) => {
4741 (-32002, format!("not addressable: {}", na.identity))
4742 }
4743 IdentityRuntimeError::NoActiveLease(id) => (-32003, format!("no active lease: {id}")),
4744 IdentityRuntimeError::LeaseLost(id) => (-32005, format!("lease lost: {id}")),
4750 _ => (-32603, format!("{err}")),
4751 };
4752 JsonRpcResponse {
4753 jsonrpc: JSONRPC_VERSION.to_string(),
4754 id: response_id,
4755 result: None,
4756 error: Some(JsonRpcError {
4757 code,
4758 message,
4759 data: None,
4760 }),
4761 }
4762}
4763
4764fn error_response(response_id: Value, code: i64, message: impl Into<String>) -> String {
4765 let message = message.into();
4766 let ambiguous_alias_rest = message
4767 .strip_prefix("ambiguous live identity alias ")
4768 .or_else(|| message.strip_prefix("invalid identity: ambiguous live identity alias "));
4769 let stale_live_alias_rest = message
4770 .strip_prefix("stale live identity alias: live console alias ")
4771 .or_else(|| {
4772 message.strip_prefix("invalid identity: stale live identity alias: live console alias ")
4773 });
4774 let hidden_policy_identity = message
4775 .strip_prefix("identity hidden by policy: ")
4776 .or_else(|| message.strip_prefix("invalid identity: identity hidden by policy: "));
4777 let data = if let Some(rest) = ambiguous_alias_rest {
4778 let (identity, candidates) = rest
4779 .split_once(": candidates [")
4780 .map(|(identity, candidates)| {
4781 (
4782 identity.to_string(),
4783 candidates
4784 .trim_end_matches(']')
4785 .split(',')
4786 .map(str::trim)
4787 .filter(|value| !value.is_empty())
4788 .map(str::to_string)
4789 .collect::<Vec<_>>(),
4790 )
4791 })
4792 .unwrap_or_else(|| (rest.to_string(), Vec::new()));
4793 Some(serde_json::json!({
4794 "kind": "ambiguous_live_identity_alias",
4795 "identity": identity,
4796 "candidates": candidates,
4797 }))
4798 } else if let Some(rest) = stale_live_alias_rest {
4799 let (identity, rest) = rest.split_once(" resolves to ").unwrap_or((rest, ""));
4800 let (runtime_member_id, bound_identity) = rest
4801 .split_once(", but identity runtime binding belongs to ")
4802 .unwrap_or((rest, ""));
4803 Some(serde_json::json!({
4804 "kind": "stale_live_identity_alias",
4805 "identity": identity,
4806 "live_runtime_member_id": runtime_member_id,
4807 "bound_identity": bound_identity,
4808 }))
4809 } else {
4810 hidden_policy_identity.map(|identity| {
4811 serde_json::json!({
4812 "kind": "identity_hidden_by_policy",
4813 "identity": identity,
4814 })
4815 })
4816 };
4817 serialize_response(&JsonRpcResponse {
4818 jsonrpc: JSONRPC_VERSION.to_string(),
4819 id: response_id,
4820 result: None,
4821 error: Some(JsonRpcError {
4822 code,
4823 message,
4824 data,
4825 }),
4826 })
4827}
4828
4829async fn resolve_live_target(
4835 runtime: &UnifiedRuntime,
4836 params: &Value,
4837) -> Option<meerkat_core::types::SessionId> {
4838 if let Some(raw) = params.get("session_id").and_then(Value::as_str) {
4839 return meerkat_core::types::SessionId::parse(raw).ok();
4840 }
4841 let raw = params
4842 .get("identity")
4843 .and_then(Value::as_str)
4844 .or_else(|| params.get("member_id").and_then(Value::as_str))?;
4845 let handle = runtime.mob_handle();
4846 let direct = crate::member_comms_id::mob_member_id(raw);
4847 let member_id = if handle.get_member(&direct).await.ok().flatten().is_some() {
4848 direct
4849 } else if let Some(identity) = handle
4850 .roster()
4851 .await
4852 .find_by_label("agent_identity", raw)
4853 .map(|entry| entry.agent_identity.clone())
4854 {
4855 identity
4856 } else {
4857 direct
4858 };
4859 handle.resolve_bridge_session_id(&member_id).await
4860}
4861
4862fn maybe_error_response(
4863 is_notification: bool,
4864 response_id: Value,
4865 code: i64,
4866 message: impl Into<String>,
4867) -> String {
4868 if is_notification {
4869 String::new()
4870 } else {
4871 error_response(response_id, code, message)
4872 }
4873}
4874
4875pub(crate) fn agent_memory_rpc_error(
4876 operation: &str,
4877 err: crate::identity_first::AgentMemoryError,
4878) -> JsonRpcError {
4879 let code = match &err {
4880 crate::identity_first::AgentMemoryError::InvalidConfig(_)
4881 | crate::identity_first::AgentMemoryError::InvalidRecord(_) => -32602,
4882 crate::identity_first::AgentMemoryError::Unsupported(_) => -32601,
4883 crate::identity_first::AgentMemoryError::Io(_)
4884 | crate::identity_first::AgentMemoryError::Parse(_)
4885 | crate::identity_first::AgentMemoryError::Timeout(_) => -32603,
4886 };
4887 JsonRpcError {
4888 code,
4889 message: format!("agent memory {operation} failed: {err}"),
4890 data: None,
4891 }
4892}
4893
4894fn serialize_response(response: &JsonRpcResponse) -> String {
4895 serde_json::to_string(response).unwrap_or_else(|_| {
4896 r#"{"jsonrpc":"2.0","id":null,"error":{"code":-32603,"message":"Internal error"}}"#
4897 .to_string()
4898 })
4899}
4900
4901#[cfg(test)]
4902#[allow(clippy::expect_used)]
4903mod tests {
4904 use super::{
4905 error_response, handle_unified_rpc_json, identity_error_response,
4906 resolve_rpc_identity_control_target, rpc_live_identity_alias_visible,
4907 rpc_member_id_matches_durable_identity,
4908 };
4909 use crate::identity_first::contracts::{ContinuityStore, LeaseProvider, RosterProvider};
4910 use crate::identity_first::{
4911 AgentAddressability, AgentBuildDraft, AgentIdentity, AgentRuntimeId, BridgeError,
4912 CheckpointVersion, ContinuityGeneration, ContinuityRecord, DurabilityPolicy,
4913 DurableAgentSpec, FencingToken, IdentityLifecycleState, IdentityRuntime,
4914 IdentityRuntimeConfig, LeaseAcquireResult, LeaseGrant, LocalContinuityStore,
4915 LocalLeaseProvider, MarkdownAgentMemoryStore, ResumeSessionOutcome, RosterContext,
4916 RosterError, SessionBridge, SessionSnapshot,
4917 };
4918 use crate::{
4919 DiscoverySpec, IdentityFirstContext, MobBootstrapOptions, MobBootstrapSpec, MobKitConfig,
4920 UnifiedRuntime,
4921 };
4922 use async_trait::async_trait;
4923 use meerkat::{AgentFactory, Config, build_ephemeral_service};
4924 use meerkat_client::TestClient;
4925 use meerkat_mob::{MobDefinition, MobStorage, SpawnMemberSpec};
4926 use serde_json::{Value, json};
4927 use std::collections::BTreeMap;
4928 use std::sync::Arc;
4929 use std::time::Duration;
4930
4931 #[derive(Debug, Default)]
4932 struct EmptyRosterProvider;
4933
4934 #[async_trait]
4935 impl RosterProvider for EmptyRosterProvider {
4936 async fn roster(
4937 &self,
4938 _context: &RosterContext,
4939 ) -> Result<Vec<DurableAgentSpec>, RosterError> {
4940 Ok(Vec::new())
4941 }
4942 }
4943
4944 #[derive(Debug, Default)]
4945 struct ContextRequiredEmptyRosterProvider {
4946 missing_definition_calls: std::sync::atomic::AtomicUsize,
4947 }
4948
4949 impl ContextRequiredEmptyRosterProvider {
4950 fn missing_definition_calls(&self) -> usize {
4951 self.missing_definition_calls
4952 .load(std::sync::atomic::Ordering::SeqCst)
4953 }
4954 }
4955
4956 #[async_trait]
4957 impl RosterProvider for ContextRequiredEmptyRosterProvider {
4958 async fn roster(
4959 &self,
4960 context: &RosterContext,
4961 ) -> Result<Vec<DurableAgentSpec>, RosterError> {
4962 if context.mob_definition.is_none() {
4963 self.missing_definition_calls
4964 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
4965 return Err(RosterError::ProviderUnavailable(
4966 "mob definition required".to_string(),
4967 ));
4968 }
4969 Ok(Vec::new())
4970 }
4971 }
4972
4973 #[derive(Debug)]
4974 struct ContextRequiredStaticRosterProvider {
4975 specs: Vec<DurableAgentSpec>,
4976 missing_definition_calls: std::sync::atomic::AtomicUsize,
4977 }
4978
4979 impl ContextRequiredStaticRosterProvider {
4980 fn new(specs: Vec<DurableAgentSpec>) -> Self {
4981 Self {
4982 specs,
4983 missing_definition_calls: std::sync::atomic::AtomicUsize::new(0),
4984 }
4985 }
4986
4987 fn missing_definition_calls(&self) -> usize {
4988 self.missing_definition_calls
4989 .load(std::sync::atomic::Ordering::SeqCst)
4990 }
4991 }
4992
4993 #[async_trait]
4994 impl RosterProvider for ContextRequiredStaticRosterProvider {
4995 async fn roster(
4996 &self,
4997 context: &RosterContext,
4998 ) -> Result<Vec<DurableAgentSpec>, RosterError> {
4999 if context.mob_definition.is_none() {
5000 self.missing_definition_calls
5001 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
5002 return Err(RosterError::ProviderUnavailable(
5003 "mob definition required".to_string(),
5004 ));
5005 }
5006 Ok(self.specs.clone())
5007 }
5008 }
5009
5010 #[derive(Debug, Default)]
5011 struct RpcResetTestBridge {
5012 create_calls: std::sync::atomic::AtomicUsize,
5013 last_create_spec: tokio::sync::Mutex<Option<DurableAgentSpec>>,
5014 }
5015
5016 impl RpcResetTestBridge {
5017 async fn last_create_spec(&self) -> Option<DurableAgentSpec> {
5018 self.last_create_spec.lock().await.clone()
5019 }
5020 }
5021
5022 #[async_trait]
5023 impl SessionBridge for RpcResetTestBridge {
5024 async fn create_session(
5025 &self,
5026 _identity: &AgentIdentity,
5027 _runtime_id: &AgentRuntimeId,
5028 spec: &DurableAgentSpec,
5029 _draft: &AgentBuildDraft,
5030 session_id: &meerkat_core::types::SessionId,
5031 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
5032 self.create_calls
5033 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
5034 *self.last_create_spec.lock().await = Some(spec.clone());
5035 Ok(session_id.clone())
5036 }
5037
5038 async fn resume_session(
5039 &self,
5040 _identity: &AgentIdentity,
5041 _runtime_id: &AgentRuntimeId,
5042 _spec: &DurableAgentSpec,
5043 _draft: &AgentBuildDraft,
5044 session_id: &meerkat_core::types::SessionId,
5045 _snapshot: &SessionSnapshot,
5046 ) -> Result<ResumeSessionOutcome, BridgeError> {
5047 Ok(ResumeSessionOutcome::Resumed {
5048 session_id: session_id.clone(),
5049 })
5050 }
5051
5052 async fn deliver(
5053 &self,
5054 _runtime_id: &AgentRuntimeId,
5055 _content: &meerkat_core::ContentInput,
5056 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
5057 Ok(meerkat_core::types::SessionId::new())
5058 }
5059
5060 async fn checkpoint_session(
5061 &self,
5062 _runtime_id: &AgentRuntimeId,
5063 _session_id: &meerkat_core::types::SessionId,
5064 ) -> Result<SessionSnapshot, BridgeError> {
5065 Ok(SessionSnapshot { data: Vec::new() })
5066 }
5067
5068 async fn retire_member(&self, _runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
5069 Ok(())
5070 }
5071 }
5072
5073 struct ReadOnlyAgentMemoryProvider;
5074
5075 #[async_trait]
5076 impl crate::identity_first::AgentMemoryProvider for ReadOnlyAgentMemoryProvider {
5077 async fn recall(
5078 &self,
5079 _request: crate::identity_first::AgentMemoryRecallRequest,
5080 ) -> Result<
5081 Vec<crate::identity_first::AgentMemoryRecord>,
5082 crate::identity_first::AgentMemoryError,
5083 > {
5084 Ok(Vec::new())
5085 }
5086 }
5087
5088 fn rpc_test_mob_spec(
5089 temp_dir: &tempfile::TempDir,
5090 ) -> Result<MobBootstrapSpec, Box<dyn std::error::Error + Send + Sync>> {
5091 let session_path = temp_dir.path().join("sessions");
5092 std::fs::create_dir_all(&session_path)?;
5093 let factory = AgentFactory::new(&session_path).comms(true);
5094 let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
5095 let definition = MobDefinition::from_toml(
5096 r#"
5097[mob]
5098id = "rpc-identity-alias-test"
5099
5100[profiles.worker]
5101model = "gpt-5.5"
5102external_addressable = true
5103
5104[profiles.worker.tools]
5105comms = true
5106"#,
5107 )?;
5108 Ok(
5109 MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
5110 .with_options(MobBootstrapOptions {
5111 allow_ephemeral_sessions: true,
5112 notify_orchestrator_on_resume: true,
5113 default_llm_client: Some(Arc::new(TestClient::default())),
5114 }),
5115 )
5116 }
5117
5118 fn rpc_reprofile_mob_spec(
5119 temp_dir: &tempfile::TempDir,
5120 ) -> Result<MobBootstrapSpec, Box<dyn std::error::Error + Send + Sync>> {
5121 let session_path = temp_dir.path().join("sessions");
5122 std::fs::create_dir_all(&session_path)?;
5123 let factory = AgentFactory::new(&session_path).comms(true);
5124 let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
5125 let definition = MobDefinition::from_toml(
5126 r#"
5127[mob]
5128id = "rpc-reset-reprofile-test"
5129
5130[profiles.domain]
5131model = "gpt-5.5"
5132external_addressable = true
5133
5134[profiles.domain.tools]
5135comms = true
5136
5137[profiles.security]
5138model = "gpt-5.5"
5139external_addressable = true
5140
5141[profiles.security.tools]
5142comms = true
5143shell = true
5144"#,
5145 )?;
5146 Ok(
5147 MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
5148 .with_options(MobBootstrapOptions {
5149 allow_ephemeral_sessions: true,
5150 notify_orchestrator_on_resume: true,
5151 default_llm_client: Some(Arc::new(TestClient::default())),
5152 }),
5153 )
5154 }
5155
5156 fn rpc_durable_spec(identity: &str, profile: &str) -> DurableAgentSpec {
5157 DurableAgentSpec {
5158 identity: AgentIdentity::parse(identity).expect("valid identity"),
5159 profile: meerkat_mob::ProfileName::from(profile),
5160 addressability: AgentAddressability::Addressable,
5161 display_name: None,
5162 labels: BTreeMap::new(),
5163 context: None,
5164 additional_instructions: Vec::new(),
5165 initial_message: None,
5166 runtime_mode_override: None,
5167 backend: None,
5168 binding: None,
5169 }
5170 }
5171
5172 #[tokio::test]
5179 async fn unified_rpc_spawn_rejects_reserved_system_identity()
5180 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5181 let temp_dir = tempfile::tempdir()?;
5182 let runtime = Box::pin(
5183 UnifiedRuntime::builder()
5184 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5185 .module_config(MobKitConfig {
5186 modules: Vec::new(),
5187 discovery: DiscoverySpec {
5188 namespace: "rpc-reserved-identity-test".to_string(),
5189 modules: Vec::new(),
5190 },
5191 pre_spawn: Vec::new(),
5192 })
5193 .timeout(Duration::from_secs(1))
5194 .build(),
5195 )
5196 .await?;
5197
5198 let response: Value = serde_json::from_str(
5199 &handle_unified_rpc_json(
5200 &runtime,
5201 &json!({
5202 "jsonrpc": "2.0",
5203 "id": 1,
5204 "method": "mobkit/spawn_member",
5205 "params": {
5206 "profile": "worker",
5207 "meerkat_id": crate::console_contracts::SYSTEM_EVENT_IDENTITY,
5208 },
5209 })
5210 .to_string(),
5211 Duration::from_secs(1),
5212 None,
5213 None,
5214 )
5215 .await,
5216 )?;
5217
5218 assert_eq!(response["error"]["code"], json!(-32602), "{response:#?}");
5219 assert!(
5220 response["error"]["message"]
5221 .as_str()
5222 .is_some_and(|message| message.contains("reserved")),
5223 "rejection must name the reservation: {response:#?}"
5224 );
5225 let _ = runtime.mob_handle().stop().await;
5226 Ok(())
5227 }
5228
5229 #[tokio::test]
5230 async fn unified_capabilities_separate_mobpack_authoring_from_runtime_controls()
5231 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5232 let temp_dir = tempfile::tempdir()?;
5233 let runtime = Box::pin(
5234 UnifiedRuntime::builder()
5235 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5236 .module_config(MobKitConfig {
5237 modules: Vec::new(),
5238 discovery: DiscoverySpec {
5239 namespace: "rpc-authoring-capabilities-test".to_string(),
5240 modules: Vec::new(),
5241 },
5242 pre_spawn: Vec::new(),
5243 })
5244 .timeout(Duration::from_secs(1))
5245 .build(),
5246 )
5247 .await?;
5248
5249 let response: Value = serde_json::from_str(
5250 &handle_unified_rpc_json(
5251 &runtime,
5252 &json!({
5253 "jsonrpc": "2.0",
5254 "id": 1,
5255 "method": "mobkit/capabilities",
5256 })
5257 .to_string(),
5258 Duration::from_secs(1),
5259 None,
5260 None,
5261 )
5262 .await,
5263 )?;
5264
5265 assert!(response["error"].is_null(), "{response:#?}");
5266 let methods = response["result"]["methods"]
5267 .as_array()
5268 .expect("methods array")
5269 .iter()
5270 .filter_map(Value::as_str)
5271 .collect::<Vec<_>>();
5272 for method in super::MOBPACK_AUTHORING_METHODS {
5273 assert!(
5274 methods.contains(method),
5275 "missing authoring method {method}"
5276 );
5277 }
5278 assert_eq!(
5279 response["result"]["authoring_capabilities"]["domain"],
5280 json!("mobpack_authoring")
5281 );
5282 assert_eq!(
5283 response["result"]["authoring_capabilities"]["runtime_mutation"],
5284 json!(false)
5285 );
5286 assert_eq!(
5287 response["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
5288 json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
5289 );
5290 assert_eq!(
5291 response["result"]["authoring_capabilities"]["methods"]
5292 .as_array()
5293 .expect("authoring methods")
5294 .iter()
5295 .filter_map(Value::as_str)
5296 .collect::<Vec<_>>(),
5297 super::MOBPACK_AUTHORING_METHODS
5298 );
5299 assert_eq!(
5300 response["result"]["authoring_capabilities"]["deploy_command"],
5301 json!("rkat mob run")
5302 );
5303
5304 Ok(())
5305 }
5306
5307 #[tokio::test]
5308 async fn unified_rpc_dispatches_mobpack_authoring_methods()
5309 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5310 let temp_dir = tempfile::tempdir()?;
5311 let runtime = Box::pin(
5312 UnifiedRuntime::builder()
5313 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5314 .module_config(MobKitConfig {
5315 modules: Vec::new(),
5316 discovery: DiscoverySpec {
5317 namespace: "rpc-authoring-dispatch-test".to_string(),
5318 modules: Vec::new(),
5319 },
5320 pre_spawn: Vec::new(),
5321 })
5322 .timeout(Duration::from_secs(1))
5323 .build(),
5324 )
5325 .await?;
5326
5327 let response: Value = serde_json::from_str(
5328 &handle_unified_rpc_json(
5329 &runtime,
5330 &json!({
5331 "jsonrpc": "2.0",
5332 "id": 1,
5333 "method": "mobkit/mobpacks/schema",
5334 })
5335 .to_string(),
5336 Duration::from_secs(1),
5337 None,
5338 None,
5339 )
5340 .await,
5341 )?;
5342
5343 assert!(response["error"].is_null(), "{response:#?}");
5344 assert_eq!(
5345 response["result"]["media_type"],
5346 json!("application/vnd.meerkat.mobpack")
5347 );
5348 assert_eq!(
5349 response["result"]["commands"]["deploy_rpc"],
5350 json!("mobkit/mobpacks/deploy")
5351 );
5352 assert_eq!(
5353 response["result"]["deploy_settings"]["runtime_backed"],
5354 json!(true)
5355 );
5356 assert_eq!(
5357 response["result"]["deploy_settings"]["authoring_provider"]["runtime_binding"],
5358 json!("bound")
5359 );
5360 assert_eq!(
5361 response["result"]["deploy_settings"]["provenance"]["source"],
5362 json!("UnifiedRuntime.authoring_provider.deploy_target")
5363 );
5364 assert!(response["result"]["sample_mobpacks"].is_null());
5365 assert!(response["result"]["agent_definitions"].is_null());
5366
5367 let catalogs: Value = serde_json::from_str(
5368 &handle_unified_rpc_json(
5369 &runtime,
5370 &json!({
5371 "jsonrpc": "2.0",
5372 "id": 2,
5373 "method": "mobkit/mobpacks/catalogs",
5374 })
5375 .to_string(),
5376 Duration::from_secs(1),
5377 None,
5378 None,
5379 )
5380 .await,
5381 )?;
5382 assert!(catalogs["error"].is_null(), "{catalogs:#?}");
5383 assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
5384 assert_eq!(
5385 catalogs["result"]["authoring_provider"]["id"],
5386 json!("unified_runtime")
5387 );
5388 assert_eq!(
5389 catalogs["result"]["authoring_provider"]["runtime_binding"],
5390 json!("bound")
5391 );
5392 assert_eq!(
5393 catalogs["result"]["sources"]["runtime"],
5394 json!("unified_runtime")
5395 );
5396 assert_eq!(
5397 catalogs["result"]["sources"]["runtime_binding"],
5398 json!("bound")
5399 );
5400 assert!(
5401 catalogs["result"]["runtime_unavailable_reason"].is_null(),
5402 "{catalogs:#?}"
5403 );
5404 assert_eq!(
5405 catalogs["result"]["catalog_snapshot"]["runtime_backed"],
5406 json!(true)
5407 );
5408 assert_eq!(
5409 catalogs["result"]["authoring_provider"]["deploy_target"]["command"],
5410 json!("rkat mob run")
5411 );
5412 assert!(
5413 catalogs["result"]["authoring_provider"]["runtime_methods"]
5414 .as_array()
5415 .is_some_and(|methods| methods.contains(&json!("mobkit/mobpacks/deploy"))),
5416 "{catalogs:#?}"
5417 );
5418
5419 Ok(())
5420 }
5421
5422 #[tokio::test]
5423 async fn reconcile_identity_passes_mob_definition_to_roster_provider()
5424 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5425 let temp_dir = tempfile::tempdir()?;
5426 let runtime = Box::pin(
5427 UnifiedRuntime::builder()
5428 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5429 .module_config(MobKitConfig {
5430 modules: Vec::new(),
5431 discovery: DiscoverySpec {
5432 namespace: "rpc-reconcile-roster-context-test".to_string(),
5433 modules: Vec::new(),
5434 },
5435 pre_spawn: Vec::new(),
5436 })
5437 .timeout(Duration::from_secs(1))
5438 .build(),
5439 )
5440 .await?;
5441 let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
5442 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5443 lease_provider: Arc::new(LocalLeaseProvider::new()),
5444 runtime_instance_id: "rpc-reconcile-roster-context-test".to_string(),
5445 has_runtime_store: true,
5446 durability_policy: DurabilityPolicy::SyncWriteThrough,
5447 bridge: None,
5448 default_timeout: None,
5449 }));
5450 let roster_provider = Arc::new(ContextRequiredEmptyRosterProvider::default());
5451 let identity_ctx = IdentityFirstContext {
5452 runtime: identity_rt,
5453 roster_provider: roster_provider.clone(),
5454 topology_provider: None,
5455 customizer: None,
5456 agent_memory_provider: None,
5457 mob_definition: Some(runtime.mob_handle().definition().clone()),
5458 };
5459
5460 let response: Value = serde_json::from_str(
5461 &handle_unified_rpc_json(
5462 &runtime,
5463 &json!({
5464 "jsonrpc": "2.0",
5465 "id": 1,
5466 "method": "mobkit/reconcile_identity",
5467 "params": {},
5468 })
5469 .to_string(),
5470 Duration::from_secs(1),
5471 None,
5472 Some(&identity_ctx),
5473 )
5474 .await,
5475 )?;
5476
5477 assert!(response["error"].is_null(), "{response:#?}");
5478 assert_eq!(
5479 roster_provider.missing_definition_calls(),
5480 0,
5481 "reconcile_identity must preserve mob_definition in roster context"
5482 );
5483 Ok(())
5484 }
5485
5486 #[tokio::test]
5487 async fn reset_identity_preserves_mob_definition_for_roster_reprofile()
5488 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5489 let temp_dir = tempfile::tempdir()?;
5490 let runtime = Box::pin(
5491 UnifiedRuntime::builder()
5492 .mob_spec(rpc_reprofile_mob_spec(&temp_dir)?)
5493 .module_config(MobKitConfig {
5494 modules: Vec::new(),
5495 discovery: DiscoverySpec {
5496 namespace: "rpc-reset-roster-context-test".to_string(),
5497 modules: Vec::new(),
5498 },
5499 pre_spawn: Vec::new(),
5500 })
5501 .timeout(Duration::from_secs(1))
5502 .build(),
5503 )
5504 .await?;
5505
5506 let continuity_store = Arc::new(LocalContinuityStore::in_memory()?);
5507 let lease_provider = Arc::new(LocalLeaseProvider::new());
5508 let bridge = Arc::new(RpcResetTestBridge::default());
5509 let identity_rt = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
5510 continuity_store: continuity_store.clone(),
5511 lease_provider: lease_provider.clone(),
5512 runtime_instance_id: "rpc-reset-roster-context-test".to_string(),
5513 has_runtime_store: true,
5514 durability_policy: DurabilityPolicy::SyncWriteThrough,
5515 bridge: Some(bridge.clone()),
5516 default_timeout: None,
5517 }));
5518
5519 let identity = AgentIdentity::parse("domain:security")?;
5520 let initial_grants = lease_provider
5521 .acquire_leases(
5522 std::slice::from_ref(&identity),
5523 "rpc-reset-roster-context-test",
5524 )
5525 .await?;
5526 let initial_grant = match initial_grants.get(&identity) {
5527 Some(LeaseAcquireResult::Acquired(grant)) => grant.clone(),
5528 other => return Err(format!("expected acquired lease, got {other:?}").into()),
5529 };
5530 let initial_record = ContinuityRecord {
5531 identity: identity.clone(),
5532 agent_runtime_id: AgentRuntimeId::parse("rt:domain:security:0")?,
5533 session_id: meerkat_core::types::SessionId::new(),
5534 generation: ContinuityGeneration::new(0),
5535 checkpoint_version: CheckpointVersion::new(0),
5536 };
5537 continuity_store
5538 .upsert_continuity_record(&initial_record, initial_grant.fencing_token)
5539 .await?;
5540 identity_rt
5541 .register(
5542 rpc_durable_spec(identity.as_str(), "domain"),
5543 IdentityLifecycleState::Active,
5544 Some(initial_record),
5545 Some(initial_grant),
5546 )
5547 .await;
5548
5549 let roster_provider = Arc::new(ContextRequiredStaticRosterProvider::new(vec![
5550 rpc_durable_spec(identity.as_str(), "security"),
5551 ]));
5552 let identity_ctx = IdentityFirstContext {
5553 runtime: identity_rt.clone(),
5554 roster_provider: roster_provider.clone(),
5555 topology_provider: None,
5556 customizer: None,
5557 agent_memory_provider: None,
5558 mob_definition: Some(runtime.mob_handle().definition().clone()),
5559 };
5560
5561 let response: Value = serde_json::from_str(
5562 &handle_unified_rpc_json(
5563 &runtime,
5564 &json!({
5565 "jsonrpc": "2.0",
5566 "id": 1,
5567 "method": "mobkit/reset",
5568 "params": { "identity": identity.as_str() },
5569 })
5570 .to_string(),
5571 Duration::from_secs(1),
5572 None,
5573 Some(&identity_ctx),
5574 )
5575 .await,
5576 )?;
5577
5578 assert!(response["error"].is_null(), "{response:#?}");
5579 assert_eq!(
5580 roster_provider.missing_definition_calls(),
5581 0,
5582 "mobkit/reset must preserve mob_definition when installing the reset roster provider"
5583 );
5584 assert_eq!(
5585 bridge
5586 .create_calls
5587 .load(std::sync::atomic::Ordering::SeqCst),
5588 1,
5589 "reset should rebuild through the bridge"
5590 );
5591 let created = bridge
5592 .last_create_spec()
5593 .await
5594 .expect("reset should record created spec");
5595 assert_eq!(created.profile.as_str(), "security");
5596 assert_eq!(
5597 identity_rt
5598 .status(&identity)
5599 .await?
5600 .profile
5601 .expect("identity should keep profile")
5602 .as_str(),
5603 "security"
5604 );
5605 Ok(())
5606 }
5607
5608 #[tokio::test]
5609 async fn unified_rpc_agent_memory_remember_writes_identity_scoped_record()
5610 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5611 let temp_dir = tempfile::tempdir()?;
5612 let runtime = Box::pin(
5613 UnifiedRuntime::builder()
5614 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5615 .module_config(MobKitConfig {
5616 modules: Vec::new(),
5617 discovery: DiscoverySpec {
5618 namespace: "rpc-agent-memory-remember-test".to_string(),
5619 modules: Vec::new(),
5620 },
5621 pre_spawn: Vec::new(),
5622 })
5623 .timeout(Duration::from_secs(1))
5624 .build(),
5625 )
5626 .await?;
5627 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5628 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5629 lease_provider: Arc::new(LocalLeaseProvider::new()),
5630 runtime_instance_id: "rpc-agent-memory-remember-test".to_string(),
5631 has_runtime_store: true,
5632 durability_policy: DurabilityPolicy::SyncWriteThrough,
5633 bridge: None,
5634 default_timeout: None,
5635 });
5636 let store = Arc::new(MarkdownAgentMemoryStore::open(
5637 temp_dir.path().join("agent-memory"),
5638 )?);
5639 let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> = store.clone();
5640 identity_rt
5641 .set_agent_memory(Some(
5642 crate::identity_first::AgentMemoryRuntimeInjector::new(
5643 provider.clone(),
5644 crate::identity_first::AgentMemoryConfig::default(),
5645 ),
5646 ))
5647 .await;
5648 let memory_identity = AgentIdentity::parse("identity:luka")?;
5649 identity_rt
5650 .register(
5651 DurableAgentSpec {
5652 identity: memory_identity.clone(),
5653 profile: meerkat_mob::ProfileName::from("worker"),
5654 addressability: AgentAddressability::Addressable,
5655 display_name: None,
5656 labels: Default::default(),
5657 context: None,
5658 additional_instructions: Vec::new(),
5659 initial_message: None,
5660 runtime_mode_override: None,
5661 backend: None,
5662 binding: None,
5663 },
5664 IdentityLifecycleState::Active,
5665 None,
5666 None,
5667 )
5668 .await;
5669 let identity_ctx = IdentityFirstContext {
5670 runtime: Arc::new(identity_rt),
5671 roster_provider: Arc::new(EmptyRosterProvider),
5672 topology_provider: None,
5673 customizer: None,
5674 agent_memory_provider: Some(provider),
5675 mob_definition: None,
5676 };
5677
5678 let capabilities: Value = serde_json::from_str(
5679 &handle_unified_rpc_json(
5680 &runtime,
5681 &json!({
5682 "jsonrpc": "2.0",
5683 "id": 1,
5684 "method": "mobkit/capabilities",
5685 })
5686 .to_string(),
5687 Duration::from_secs(1),
5688 None,
5689 Some(&identity_ctx),
5690 )
5691 .await,
5692 )?;
5693 assert!(
5694 capabilities["result"]["methods"]
5695 .as_array()
5696 .is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/remember"))),
5697 "{capabilities:#?}"
5698 );
5699 assert!(
5700 capabilities["result"]["methods"]
5701 .as_array()
5702 .is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/recall"))),
5703 "{capabilities:#?}"
5704 );
5705 assert!(
5706 capabilities["result"]["methods"]
5707 .as_array()
5708 .is_some_and(|methods| methods.contains(&json!("mobkit/agent_memory/forget"))),
5709 "{capabilities:#?}"
5710 );
5711 assert!(
5714 capabilities["result"]["methods"]
5715 .as_array()
5716 .is_some_and(|methods| !methods.contains(&json!("mobkit/agent_memory/update"))),
5717 "{capabilities:#?}"
5718 );
5719 assert!(
5720 capabilities["result"]["methods"]
5721 .as_array()
5722 .is_some_and(|methods| !methods.contains(&json!("mobkit/agent_memory/manifest"))),
5723 "{capabilities:#?}"
5724 );
5725
5726 let response: Value = serde_json::from_str(
5727 &handle_unified_rpc_json(
5728 &runtime,
5729 &json!({
5730 "jsonrpc": "2.0",
5731 "id": 2,
5732 "method": "mobkit/agent_memory/remember",
5733 "params": {
5734 "identity": "identity:luka",
5735 "realm": "family",
5736 "title": "School pickup",
5737 "body": "Pickup is before calendar planning.",
5738 "tags": ["family", "calendar"]
5739 },
5740 })
5741 .to_string(),
5742 Duration::from_secs(1),
5743 None,
5744 Some(&identity_ctx),
5745 )
5746 .await,
5747 )?;
5748
5749 assert!(response["error"].is_null(), "{response:#?}");
5750 let memory_id = response["result"]["memory_id"]
5751 .as_str()
5752 .ok_or("memory_id should be present")?
5753 .to_string();
5754 assert_eq!(response["result"]["title"], json!("School pickup"));
5755 assert_eq!(
5756 response["result"]["body"],
5757 json!("Pickup is before calendar planning.")
5758 );
5759 assert_eq!(response["result"]["tags"], json!(["calendar", "family"]));
5760
5761 let records = store.read_records("family", &memory_identity)?;
5762 assert_eq!(records.len(), 1);
5763 assert_eq!(records[0].title, "School pickup");
5764 assert_eq!(records[0].tags, vec!["calendar", "family"]);
5765
5766 let recall_response: Value = serde_json::from_str(
5767 &handle_unified_rpc_json(
5768 &runtime,
5769 &json!({
5770 "jsonrpc": "2.0",
5771 "id": 3,
5772 "method": "mobkit/agent_memory/recall",
5773 "params": {
5774 "identity": "identity:luka",
5775 "realm": "family",
5776 "selection": "contextual",
5777 "query_terms": ["pickup"],
5778 "max_entries": 4
5779 },
5780 })
5781 .to_string(),
5782 Duration::from_secs(1),
5783 None,
5784 Some(&identity_ctx),
5785 )
5786 .await,
5787 )?;
5788
5789 assert!(recall_response["error"].is_null(), "{recall_response:#?}");
5790 assert_eq!(
5791 recall_response["result"]["records"]
5792 .as_array()
5793 .map(Vec::len),
5794 Some(1)
5795 );
5796 assert_eq!(
5797 recall_response["result"]["records"][0]["body"],
5798 json!("Pickup is before calendar planning.")
5799 );
5800
5801 let forget_response: Value = serde_json::from_str(
5802 &handle_unified_rpc_json(
5803 &runtime,
5804 &json!({
5805 "jsonrpc": "2.0",
5806 "id": 4,
5807 "method": "mobkit/agent_memory/forget",
5808 "params": {
5809 "identity": "identity:luka",
5810 "realm": "family",
5811 "memory_id": memory_id.clone()
5812 },
5813 })
5814 .to_string(),
5815 Duration::from_secs(1),
5816 None,
5817 Some(&identity_ctx),
5818 )
5819 .await,
5820 )?;
5821 assert!(forget_response["error"].is_null(), "{forget_response:#?}");
5822 assert_eq!(forget_response["result"]["memory_id"], json!(memory_id));
5823 assert_eq!(forget_response["result"]["deleted"], json!(true));
5824 assert!(store.read_records("family", &memory_identity)?.is_empty());
5825
5826 let recall_after_forget_response: Value = serde_json::from_str(
5827 &handle_unified_rpc_json(
5828 &runtime,
5829 &json!({
5830 "jsonrpc": "2.0",
5831 "id": 5,
5832 "method": "mobkit/agent_memory/recall",
5833 "params": {
5834 "identity": "identity:luka",
5835 "realm": "family",
5836 "selection": "always"
5837 },
5838 })
5839 .to_string(),
5840 Duration::from_secs(1),
5841 None,
5842 Some(&identity_ctx),
5843 )
5844 .await,
5845 )?;
5846 assert!(
5847 recall_after_forget_response["error"].is_null(),
5848 "{recall_after_forget_response:#?}"
5849 );
5850 assert_eq!(
5851 recall_after_forget_response["result"]["records"]
5852 .as_array()
5853 .map(Vec::len),
5854 Some(0)
5855 );
5856
5857 let unknown_response: Value = serde_json::from_str(
5858 &handle_unified_rpc_json(
5859 &runtime,
5860 &json!({
5861 "jsonrpc": "2.0",
5862 "id": 6,
5863 "method": "mobkit/agent_memory/remember",
5864 "params": {
5865 "identity": "identity:unknown",
5866 "title": "Orphan",
5867 "body": "This should not be written."
5868 },
5869 })
5870 .to_string(),
5871 Duration::from_secs(1),
5872 None,
5873 Some(&identity_ctx),
5874 )
5875 .await,
5876 )?;
5877 assert_eq!(unknown_response["error"]["code"], json!(-32602));
5878 assert!(
5879 unknown_response["error"]["message"]
5880 .as_str()
5881 .is_some_and(|message| message.contains("unknown identity")),
5882 "{unknown_response:#?}"
5883 );
5884
5885 let unknown_forget_response: Value = serde_json::from_str(
5886 &handle_unified_rpc_json(
5887 &runtime,
5888 &json!({
5889 "jsonrpc": "2.0",
5890 "id": 7,
5891 "method": "mobkit/agent_memory/forget",
5892 "params": {
5893 "identity": "identity:unknown",
5894 "memory_id": "mem-missing"
5895 },
5896 })
5897 .to_string(),
5898 Duration::from_secs(1),
5899 None,
5900 Some(&identity_ctx),
5901 )
5902 .await,
5903 )?;
5904 assert_eq!(unknown_forget_response["error"]["code"], json!(-32602));
5905 assert!(
5906 unknown_forget_response["error"]["message"]
5907 .as_str()
5908 .is_some_and(|message| message.contains("unknown identity")),
5909 "{unknown_forget_response:#?}"
5910 );
5911
5912 Ok(())
5913 }
5914
5915 #[tokio::test]
5916 async fn unified_rpc_agent_memory_update_and_manifest_over_sqlite_store()
5917 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5918 let temp_dir = tempfile::tempdir()?;
5919 let runtime = Box::pin(
5920 UnifiedRuntime::builder()
5921 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5922 .module_config(MobKitConfig {
5923 modules: Vec::new(),
5924 discovery: DiscoverySpec {
5925 namespace: "rpc-agent-memory-sqlite-test".to_string(),
5926 modules: Vec::new(),
5927 },
5928 pre_spawn: Vec::new(),
5929 })
5930 .timeout(Duration::from_secs(1))
5931 .build(),
5932 )
5933 .await?;
5934 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5935 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5936 lease_provider: Arc::new(LocalLeaseProvider::new()),
5937 runtime_instance_id: "rpc-agent-memory-sqlite-test".to_string(),
5938 has_runtime_store: true,
5939 durability_policy: DurabilityPolicy::SyncWriteThrough,
5940 bridge: None,
5941 default_timeout: None,
5942 });
5943 let store = Arc::new(crate::memory::SqliteAgentMemoryStore::open(
5944 temp_dir.path().join("agent-memory"),
5945 )?);
5946 let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> = store.clone();
5947 identity_rt
5948 .set_agent_memory(Some(
5949 crate::identity_first::AgentMemoryRuntimeInjector::new(
5950 provider.clone(),
5951 crate::identity_first::AgentMemoryConfig::default(),
5952 ),
5953 ))
5954 .await;
5955 let memory_identity = AgentIdentity::parse("identity:luka")?;
5956 identity_rt
5957 .register(
5958 DurableAgentSpec {
5959 identity: memory_identity.clone(),
5960 profile: meerkat_mob::ProfileName::from("worker"),
5961 addressability: AgentAddressability::Addressable,
5962 display_name: None,
5963 labels: Default::default(),
5964 context: None,
5965 additional_instructions: Vec::new(),
5966 initial_message: None,
5967 runtime_mode_override: None,
5968 backend: None,
5969 binding: None,
5970 },
5971 IdentityLifecycleState::Active,
5972 None,
5973 None,
5974 )
5975 .await;
5976 let identity_ctx = IdentityFirstContext {
5977 runtime: Arc::new(identity_rt),
5978 roster_provider: Arc::new(EmptyRosterProvider),
5979 topology_provider: None,
5980 customizer: None,
5981 agent_memory_provider: Some(provider),
5982 mob_definition: None,
5983 };
5984
5985 let capabilities: Value = serde_json::from_str(
5988 &handle_unified_rpc_json(
5989 &runtime,
5990 &json!({
5991 "jsonrpc": "2.0",
5992 "id": 1,
5993 "method": "mobkit/capabilities",
5994 })
5995 .to_string(),
5996 Duration::from_secs(1),
5997 None,
5998 Some(&identity_ctx),
5999 )
6000 .await,
6001 )?;
6002 for method in [
6003 "mobkit/agent_memory/recall",
6004 "mobkit/agent_memory/remember",
6005 "mobkit/agent_memory/forget",
6006 "mobkit/agent_memory/update",
6007 "mobkit/agent_memory/manifest",
6008 ] {
6009 assert!(
6010 capabilities["result"]["methods"]
6011 .as_array()
6012 .is_some_and(|methods| methods.contains(&json!(method))),
6013 "missing {method}: {capabilities:#?}"
6014 );
6015 }
6016
6017 let remember_response: Value = serde_json::from_str(
6018 &handle_unified_rpc_json(
6019 &runtime,
6020 &json!({
6021 "jsonrpc": "2.0",
6022 "id": 2,
6023 "method": "mobkit/agent_memory/remember",
6024 "params": {
6025 "identity": "identity:luka",
6026 "realm": "family",
6027 "title": "School pickup",
6028 "body": "Pickup is before calendar planning.",
6029 },
6030 })
6031 .to_string(),
6032 Duration::from_secs(1),
6033 None,
6034 Some(&identity_ctx),
6035 )
6036 .await,
6037 )?;
6038 assert!(
6039 remember_response["error"].is_null(),
6040 "{remember_response:#?}"
6041 );
6042 let memory_id = remember_response["result"]["memory_id"]
6043 .as_str()
6044 .ok_or("memory_id should be present")?
6045 .to_string();
6046
6047 let update_response: Value = serde_json::from_str(
6048 &handle_unified_rpc_json(
6049 &runtime,
6050 &json!({
6051 "jsonrpc": "2.0",
6052 "id": 3,
6053 "method": "mobkit/agent_memory/update",
6054 "params": {
6055 "identity": "identity:luka",
6056 "realm": "family",
6057 "memory_id": memory_id.clone(),
6058 "title": "School pickup",
6059 "body": "Pickup moved to after calendar planning.",
6060 "tags": ["family"]
6061 },
6062 })
6063 .to_string(),
6064 Duration::from_secs(1),
6065 None,
6066 Some(&identity_ctx),
6067 )
6068 .await,
6069 )?;
6070 assert!(update_response["error"].is_null(), "{update_response:#?}");
6071 let new_id = update_response["result"]["memory_id"]
6072 .as_str()
6073 .ok_or("updated memory_id should be present")?
6074 .to_string();
6075 assert_ne!(new_id, memory_id);
6076 assert_eq!(update_response["result"]["supersedes"], json!(memory_id));
6077
6078 let recall_response: Value = serde_json::from_str(
6080 &handle_unified_rpc_json(
6081 &runtime,
6082 &json!({
6083 "jsonrpc": "2.0",
6084 "id": 4,
6085 "method": "mobkit/agent_memory/recall",
6086 "params": {
6087 "identity": "identity:luka",
6088 "realm": "family",
6089 "selection": "always"
6090 },
6091 })
6092 .to_string(),
6093 Duration::from_secs(1),
6094 None,
6095 Some(&identity_ctx),
6096 )
6097 .await,
6098 )?;
6099 assert!(recall_response["error"].is_null(), "{recall_response:#?}");
6100 assert_eq!(
6101 recall_response["result"]["records"]
6102 .as_array()
6103 .map(Vec::len),
6104 Some(1)
6105 );
6106 assert_eq!(
6107 recall_response["result"]["records"][0]["memory_id"],
6108 json!(new_id)
6109 );
6110
6111 let manifest_response: Value = serde_json::from_str(
6112 &handle_unified_rpc_json(
6113 &runtime,
6114 &json!({
6115 "jsonrpc": "2.0",
6116 "id": 5,
6117 "method": "mobkit/agent_memory/manifest",
6118 "params": {
6119 "identity": "identity:luka",
6120 "realm": "family",
6121 "tier": "working_set",
6122 "k": 4
6123 },
6124 })
6125 .to_string(),
6126 Duration::from_secs(1),
6127 None,
6128 Some(&identity_ctx),
6129 )
6130 .await,
6131 )?;
6132 assert!(
6133 manifest_response["error"].is_null(),
6134 "{manifest_response:#?}"
6135 );
6136 let records = manifest_response["result"]["records"]
6137 .as_array()
6138 .ok_or("manifest records array")?;
6139 assert_eq!(records.len(), 1, "{manifest_response:#?}");
6140 assert_eq!(records[0]["id"], json!(new_id));
6141 assert_eq!(records[0]["kind"], json!("fact"));
6142 assert_eq!(records[0]["age_days"], json!(0));
6143 assert!(
6144 records[0].get("body").is_none(),
6145 "manifest is an index, never a dump: {manifest_response:#?}"
6146 );
6147
6148 let bad_tier_response: Value = serde_json::from_str(
6149 &handle_unified_rpc_json(
6150 &runtime,
6151 &json!({
6152 "jsonrpc": "2.0",
6153 "id": 6,
6154 "method": "mobkit/agent_memory/manifest",
6155 "params": {
6156 "identity": "identity:luka",
6157 "tier": "everything"
6158 },
6159 })
6160 .to_string(),
6161 Duration::from_secs(1),
6162 None,
6163 Some(&identity_ctx),
6164 )
6165 .await,
6166 )?;
6167 assert_eq!(bad_tier_response["error"]["code"], json!(-32602));
6168
6169 Ok(())
6170 }
6171
6172 #[tokio::test]
6173 async fn unified_rpc_agent_memory_capabilities_do_not_advertise_read_only_writes()
6174 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6175 let temp_dir = tempfile::tempdir()?;
6176 let runtime = Box::pin(
6177 UnifiedRuntime::builder()
6178 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6179 .module_config(MobKitConfig {
6180 modules: Vec::new(),
6181 discovery: DiscoverySpec {
6182 namespace: "rpc-agent-memory-read-only-test".to_string(),
6183 modules: Vec::new(),
6184 },
6185 pre_spawn: Vec::new(),
6186 })
6187 .timeout(Duration::from_secs(1))
6188 .build(),
6189 )
6190 .await?;
6191 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6192 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6193 lease_provider: Arc::new(LocalLeaseProvider::new()),
6194 runtime_instance_id: "rpc-agent-memory-read-only-test".to_string(),
6195 has_runtime_store: true,
6196 durability_policy: DurabilityPolicy::SyncWriteThrough,
6197 bridge: None,
6198 default_timeout: None,
6199 });
6200 let provider: Arc<dyn crate::identity_first::AgentMemoryProvider> =
6201 Arc::new(ReadOnlyAgentMemoryProvider);
6202 let identity_ctx = IdentityFirstContext {
6203 runtime: Arc::new(identity_rt),
6204 roster_provider: Arc::new(EmptyRosterProvider),
6205 topology_provider: None,
6206 customizer: None,
6207 agent_memory_provider: Some(provider),
6208 mob_definition: None,
6209 };
6210
6211 let capabilities: Value = serde_json::from_str(
6212 &handle_unified_rpc_json(
6213 &runtime,
6214 &json!({
6215 "jsonrpc": "2.0",
6216 "id": 1,
6217 "method": "mobkit/capabilities",
6218 })
6219 .to_string(),
6220 Duration::from_secs(1),
6221 None,
6222 Some(&identity_ctx),
6223 )
6224 .await,
6225 )?;
6226 let methods = capabilities["result"]["methods"]
6227 .as_array()
6228 .ok_or("methods should be an array")?;
6229
6230 assert!(methods.contains(&json!("mobkit/agent_memory/recall")));
6231 assert!(!methods.contains(&json!("mobkit/agent_memory/remember")));
6232 assert!(!methods.contains(&json!("mobkit/agent_memory/forget")));
6233 Ok(())
6234 }
6235
6236 #[test]
6237 fn identity_lease_lost_maps_off_capability_unavailable_code() {
6238 let identity = AgentIdentity::parse("review:singleton").expect("valid identity");
6244 let err = crate::identity_first::IdentityRuntimeError::LeaseLost(identity);
6245 let response = identity_error_response(json!("req-1"), &err);
6246 let error = response.error.expect("lease-lost must surface an error");
6247 assert_ne!(
6248 error.code, -32004,
6249 "LeaseLost must not use the capability code"
6250 );
6251 assert_eq!(
6252 error.code, -32005,
6253 "LeaseLost has its own identity-plane code"
6254 );
6255 assert!(error.message.contains("lease lost"));
6256 }
6257
6258 #[test]
6259 fn mobpack_authoring_rpc_helper_preserves_runtime_catalog_binding()
6260 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6261 let runtime = crate::mobpack::MobpackRuntimeCatalogState {
6262 loaded_modules: vec!["editor-host".to_string()],
6263 runtime_methods: vec![
6264 "mobkit/mobpacks/catalogs".to_string(),
6265 "mobkit/mobpacks/apply_operation".to_string(),
6266 "mobkit/mobpacks/deploy".to_string(),
6267 ],
6268 has_contact_directory: true,
6269 has_peer_mob_handles: false,
6270 has_inproc_contacts: false,
6271 runtime_flow_rows: vec![json!({
6272 "id": "runtime_rpc_main",
6273 "source": "mobkit/runtime/flow_projection",
6274 "document": {
6275 "mob_id": "runtime_rpc",
6276 "flow": { "name": "main", "steps": [] },
6277 "members": []
6278 },
6279 "validation": { "ok": true }
6280 })],
6281 runtime_agent_definition_sources: vec![json!({
6282 "id": "runtime_profiles_rpc",
6283 "name": "Runtime RPC profiles",
6284 "source": "mobkit/runtime/agent-definitions",
6285 "document": {
6286 "mob_id": "runtime_rpc",
6287 "members": [{
6288 "id": "m_runtime_reviewer",
6289 "name": "Runtime reviewer",
6290 "role": "runtime_reviewer",
6291 "profileBinding": "inline",
6292 "model": "gpt-5.5",
6293 "runtimeMode": "turn_driven",
6294 "tools": ["builtins"],
6295 "skills": ["mob.runtime.review"],
6296 "schema": ""
6297 }],
6298 "schemas": []
6299 }
6300 })],
6301 runtime_skill_realms: vec![json!({
6302 "id": "runtime_rpc",
6303 "label": "Runtime RPC",
6304 "source": "mobkit/runtime/skills",
6305 "skills": [{
6306 "id": "mob.runtime.review",
6307 "label": "Runtime review",
6308 "source": "inline",
6309 "content": "Review runtime work."
6310 }]
6311 })],
6312 };
6313
6314 let catalogs = super::handle_mobpack_authoring_rpc_with_runtime(
6315 "mobkit/mobpacks/catalogs",
6316 &json!({}),
6317 json!(1),
6318 Some(&runtime),
6319 )
6320 .expect("catalogs method");
6321 let catalogs: Value = serde_json::to_value(catalogs)?;
6322 assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
6323 assert_eq!(
6324 catalogs["result"]["authoring_provider"]["runtime_binding"],
6325 json!("bound")
6326 );
6327 assert_eq!(
6328 catalogs["result"]["runtime_flows"][0]["id"],
6329 json!("runtime_rpc_main")
6330 );
6331 let listed = super::handle_mobpack_authoring_rpc_with_runtime(
6332 "mobkit/mobpacks/list",
6333 &json!({}),
6334 json!(2),
6335 Some(&runtime),
6336 )
6337 .expect("list method");
6338 let listed: Value = serde_json::to_value(listed)?;
6339 assert_eq!(listed["result"]["runtime_backed"], json!(true));
6340 assert_eq!(listed["result"]["rows"][0]["id"], json!("runtime_rpc_main"));
6341
6342 let fetched = super::handle_mobpack_authoring_rpc_with_runtime(
6343 "mobkit/mobpacks/get",
6344 &json!({ "id": "runtime_rpc_main" }),
6345 json!(3),
6346 Some(&runtime),
6347 )
6348 .expect("get method");
6349 let fetched: Value = serde_json::to_value(fetched)?;
6350 assert_eq!(fetched["result"]["runtime_backed"], json!(true));
6351 assert_eq!(fetched["result"]["row"]["id"], json!("runtime_rpc_main"));
6352
6353 let definitions = super::handle_mobpack_authoring_rpc_with_runtime(
6354 "mobkit/agent_definitions/list",
6355 &json!({}),
6356 json!(4),
6357 Some(&runtime),
6358 )
6359 .expect("agent definitions method");
6360 let definitions: Value = serde_json::to_value(definitions)?;
6361 let runtime_definition = definitions["result"]["agent_definitions"]
6362 .as_array()
6363 .and_then(|rows| {
6364 rows.iter()
6365 .find(|row| row["sourceOrigin"] == "mobkit/runtime/agent-definitions")
6366 })
6367 .expect("runtime profile definition");
6368 assert_eq!(runtime_definition["role"], json!("runtime_reviewer"));
6369 assert_eq!(
6370 runtime_definition["toolDefinitions"][0]["id"],
6371 json!("builtins")
6372 );
6373 assert_eq!(
6374 runtime_definition["skillDefinitions"][0]["id"],
6375 json!("mob.runtime.review")
6376 );
6377 assert_eq!(
6378 definitions["result"]["catalog_snapshot"]["runtime_backed"],
6379 json!(true)
6380 );
6381
6382 let tools = super::handle_mobpack_authoring_rpc_with_runtime(
6383 "mobkit/tools/catalog",
6384 &json!({}),
6385 json!(5),
6386 Some(&runtime),
6387 )
6388 .expect("tools catalog method");
6389 let tools: Value = serde_json::to_value(tools)?;
6390 let mob_tool = tools["result"]["tool_catalog"]
6391 .as_array()
6392 .expect("tools")
6393 .iter()
6394 .find(|tool| tool["id"] == "mob")
6395 .expect("mob tool");
6396 assert_eq!(mob_tool["runtime_availability"]["available"], json!(false));
6397
6398 let agents = super::handle_mobpack_authoring_rpc_with_runtime(
6399 "mobkit/agent_definitions/list",
6400 &json!({}),
6401 json!(3),
6402 Some(&runtime),
6403 )
6404 .expect("agent definitions method");
6405 let agents: Value = serde_json::to_value(agents)?;
6406 assert_eq!(agents["result"]["runtime_backed"], json!(true));
6407 let planner = agents["result"]["agent_definitions"]
6408 .as_array()
6409 .expect("agent definitions")
6410 .iter()
6411 .find(|definition| definition["role"] == "planner")
6412 .expect("planner definition");
6413 let planner_mob_tool = planner["toolDefinitions"]
6414 .as_array()
6415 .expect("planner tools")
6416 .iter()
6417 .find(|tool| tool["id"] == "mob")
6418 .expect("planner mob tool");
6419 assert_eq!(
6420 planner_mob_tool["runtimeAvailability"]["state"],
6421 json!("unavailable")
6422 );
6423
6424 Ok(())
6425 }
6426
6427 #[test]
6428 fn module_rpc_dispatches_mobpack_authoring_methods()
6429 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6430 let config = MobKitConfig {
6431 modules: Vec::new(),
6432 discovery: DiscoverySpec {
6433 namespace: "module-rpc-authoring-dispatch-test".to_string(),
6434 modules: Vec::new(),
6435 },
6436 pre_spawn: Vec::new(),
6437 };
6438 let mut runtime = crate::start_mobkit_runtime(config, Vec::new(), Duration::from_secs(1))?;
6439
6440 let capabilities: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
6441 &mut runtime,
6442 &json!({
6443 "jsonrpc": "2.0",
6444 "id": 1,
6445 "method": "mobkit/capabilities",
6446 })
6447 .to_string(),
6448 Duration::from_secs(1),
6449 ))?;
6450 let methods = capabilities["result"]["methods"]
6451 .as_array()
6452 .expect("methods array")
6453 .iter()
6454 .filter_map(Value::as_str)
6455 .collect::<Vec<_>>();
6456 for method in super::MOBPACK_AUTHORING_METHODS {
6457 assert!(
6458 methods.contains(method),
6459 "missing authoring method {method}"
6460 );
6461 }
6462 assert_eq!(
6463 capabilities["result"]["authoring_capabilities"]["runtime_mutation"],
6464 json!(false)
6465 );
6466 assert_eq!(
6467 capabilities["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
6468 json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
6469 );
6470
6471 let schema: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
6472 &mut runtime,
6473 &json!({
6474 "jsonrpc": "2.0",
6475 "id": 2,
6476 "method": "mobkit/mobpacks/schema",
6477 })
6478 .to_string(),
6479 Duration::from_secs(1),
6480 ))?;
6481 assert!(schema["error"].is_null(), "{schema:#?}");
6482 assert_eq!(
6483 schema["result"]["commands"]["deploy_rpc"],
6484 json!("mobkit/mobpacks/deploy")
6485 );
6486 assert!(schema["result"]["agent_definitions"].is_null());
6487
6488 let _ = runtime.shutdown();
6489 Ok(())
6490 }
6491
6492 #[test]
6493 fn generated_runtime_ids_match_their_durable_identity_prefix() {
6494 assert!(!rpc_member_id_matches_durable_identity(
6495 "rt:review:singleton:0",
6496 "review:singleton",
6497 ));
6498 assert!(!rpc_member_id_matches_durable_identity(
6499 "review:singleton:gen1",
6500 "review:singleton",
6501 ));
6502 assert!(!rpc_member_id_matches_durable_identity(
6503 "review:singleton:1",
6504 "review:singleton",
6505 ));
6506 assert!(!rpc_member_id_matches_durable_identity(
6507 "rt:reviewer:singleton:0",
6508 "review:singleton",
6509 ));
6510 assert!(!rpc_member_id_matches_durable_identity(
6511 "rt:review:singleton:qa:0",
6512 "review:singleton",
6513 ));
6514 assert!(!rpc_member_id_matches_durable_identity(
6515 "review:singleton:qa",
6516 "review:singleton",
6517 ));
6518 }
6519
6520 #[test]
6521 fn rpc_live_identity_visibility_matches_delegate_projection_labels() {
6522 assert!(rpc_live_identity_alias_visible("worker", &BTreeMap::new()));
6523
6524 let mut labels = BTreeMap::new();
6525 labels.insert("role".to_string(), "delegate".to_string());
6526 labels.insert("source_mob_id".to_string(), "mob-a".to_string());
6527 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
6528 assert!(!rpc_live_identity_alias_visible("worker", &labels));
6529 assert!(!rpc_live_identity_alias_visible("delegate", &labels));
6530 }
6531
6532 #[test]
6533 fn ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error> {
6534 let response: Value = serde_json::from_str(&error_response(
6535 json!(1),
6536 -32602,
6537 "ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
6538 ))?;
6539
6540 assert_eq!(
6541 response["error"]["data"]["kind"],
6542 json!("ambiguous_live_identity_alias")
6543 );
6544 assert_eq!(
6545 response["error"]["data"]["identity"],
6546 json!("review:singleton")
6547 );
6548 assert_eq!(
6549 response["error"]["data"]["candidates"],
6550 json!(["rt:review:singleton:0", "rt:review:singleton:1"])
6551 );
6552 Ok(())
6553 }
6554
6555 #[test]
6556 fn wrapped_ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error>
6557 {
6558 let response: Value = serde_json::from_str(&error_response(
6559 json!(1),
6560 -32602,
6561 "invalid identity: ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
6562 ))?;
6563
6564 assert_eq!(
6565 response["error"]["data"]["kind"],
6566 json!("ambiguous_live_identity_alias")
6567 );
6568 assert_eq!(
6569 response["error"]["data"]["identity"],
6570 json!("review:singleton")
6571 );
6572 assert_eq!(
6573 response["error"]["data"]["candidates"],
6574 json!(["rt:review:singleton:0", "rt:review:singleton:1"])
6575 );
6576 Ok(())
6577 }
6578
6579 #[tokio::test]
6580 async fn runtime_id_live_only_resolution_rejects_duplicate_projected_identity()
6581 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6582 let temp_dir = tempfile::tempdir()?;
6583 let runtime = Box::pin(
6584 UnifiedRuntime::builder()
6585 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6586 .module_config(MobKitConfig {
6587 modules: Vec::new(),
6588 discovery: DiscoverySpec {
6589 namespace: "rpc-identity-alias-test".to_string(),
6590 modules: Vec::new(),
6591 },
6592 pre_spawn: Vec::new(),
6593 })
6594 .timeout(Duration::from_secs(1))
6595 .build(),
6596 )
6597 .await?;
6598 for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
6599 let mut labels = BTreeMap::new();
6600 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
6601 runtime
6602 .spawn(
6603 SpawnMemberSpec::from_wire(
6604 "worker".to_string(),
6605 runtime_id.to_string(),
6606 Some("You are a duplicate Review Agent.".into()),
6607 None,
6608 None,
6609 )
6610 .with_labels(labels),
6611 )
6612 .await?;
6613 }
6614 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6615 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6616 lease_provider: Arc::new(LocalLeaseProvider::new()),
6617 runtime_instance_id: "rpc-identity-alias-test".to_string(),
6618 has_runtime_store: true,
6619 durability_policy: DurabilityPolicy::SyncWriteThrough,
6620 bridge: None,
6621 default_timeout: None,
6622 });
6623
6624 let err =
6625 resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:0")
6626 .await
6627 .expect_err("runtime-id live-only fallback should reject duplicate durable alias");
6628 assert!(
6629 err.contains("ambiguous live identity alias review:singleton"),
6630 "unexpected error: {err}"
6631 );
6632
6633 Ok(())
6634 }
6635
6636 #[tokio::test]
6637 async fn durable_resolution_prefers_registered_live_binding_over_stale_duplicates()
6638 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6639 let temp_dir = tempfile::tempdir()?;
6640 let runtime = Box::pin(
6641 UnifiedRuntime::builder()
6642 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6643 .module_config(MobKitConfig {
6644 modules: Vec::new(),
6645 discovery: DiscoverySpec {
6646 namespace: "rpc-identity-alias-test".to_string(),
6647 modules: Vec::new(),
6648 },
6649 pre_spawn: Vec::new(),
6650 })
6651 .timeout(Duration::from_secs(1))
6652 .build(),
6653 )
6654 .await?;
6655 for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
6656 let mut labels = BTreeMap::new();
6657 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
6658 runtime
6659 .spawn(
6660 SpawnMemberSpec::from_wire(
6661 "worker".to_string(),
6662 runtime_id.to_string(),
6663 Some("You are a duplicate Review Agent.".into()),
6664 None,
6665 None,
6666 )
6667 .with_labels(labels),
6668 )
6669 .await?;
6670 }
6671 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6672 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6673 lease_provider: Arc::new(LocalLeaseProvider::new()),
6674 runtime_instance_id: "rpc-identity-alias-test".to_string(),
6675 has_runtime_store: true,
6676 durability_policy: DurabilityPolicy::SyncWriteThrough,
6677 bridge: None,
6678 default_timeout: None,
6679 });
6680 let identity = AgentIdentity::parse("review:singleton")?;
6681 let record = ContinuityRecord {
6682 identity: identity.clone(),
6683 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:1")?,
6684 session_id: meerkat_core::types::SessionId::new(),
6685 generation: ContinuityGeneration::new(1),
6686 checkpoint_version: CheckpointVersion::new(0),
6687 };
6688 identity_rt
6689 .register(
6690 DurableAgentSpec {
6691 identity,
6692 profile: meerkat_mob::ProfileName::from("worker"),
6693 addressability: AgentAddressability::Addressable,
6694 display_name: None,
6695 labels: BTreeMap::new(),
6696 context: None,
6697 additional_instructions: Vec::new(),
6698 initial_message: None,
6699 runtime_mode_override: None,
6700 backend: None,
6701 binding: None,
6702 },
6703 IdentityLifecycleState::Active,
6704 Some(record),
6705 None,
6706 )
6707 .await;
6708
6709 let target =
6710 resolve_rpc_identity_control_target(&runtime, &identity_rt, "review:singleton").await?;
6711 assert_eq!(target.identity.as_str(), "review:singleton");
6712 assert_eq!(
6713 target
6714 .live
6715 .as_ref()
6716 .map(|alias| alias.runtime_member_id.as_str()),
6717 Some("rt:review:singleton:1")
6718 );
6719
6720 let target =
6721 resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:1")
6722 .await?;
6723 assert_eq!(
6724 target
6725 .live
6726 .as_ref()
6727 .map(|alias| alias.runtime_member_id.as_str()),
6728 Some("rt:review:singleton:1")
6729 );
6730
6731 let stale_target =
6732 resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:0")
6733 .await?;
6734 assert_eq!(
6735 stale_target
6736 .live
6737 .as_ref()
6738 .map(|alias| alias.runtime_member_id.as_str()),
6739 Some("rt:review:singleton:0")
6740 );
6741 let stale_response =
6742 super::rpc_stale_live_alias_error_response(&identity_rt, &stale_target, json!(99))
6743 .await
6744 .expect("old reset generation should be rejected as stale");
6745 assert_eq!(
6746 stale_response
6747 .error
6748 .as_ref()
6749 .and_then(|error| error.data.as_ref())
6750 .and_then(|data| data.get("kind")),
6751 Some(&json!("stale_identity_runtime_binding"))
6752 );
6753 assert_eq!(
6754 stale_response
6755 .error
6756 .as_ref()
6757 .and_then(|error| error.data.as_ref())
6758 .and_then(|data| data.get("registered_runtime_member_id")),
6759 Some(&json!("rt:review:singleton:1"))
6760 );
6761 assert_eq!(
6762 stale_response
6763 .error
6764 .as_ref()
6765 .and_then(|error| error.data.as_ref())
6766 .and_then(|data| data.get("live_runtime_member_id")),
6767 Some(&json!("rt:review:singleton:0"))
6768 );
6769
6770 Ok(())
6771 }
6772
6773 #[tokio::test]
6774 async fn durable_resolution_rejects_hidden_registered_live_binding()
6775 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6776 let temp_dir = tempfile::tempdir()?;
6777 let runtime = Box::pin(
6778 UnifiedRuntime::builder()
6779 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6780 .module_config(MobKitConfig {
6781 modules: Vec::new(),
6782 discovery: DiscoverySpec {
6783 namespace: "rpc-hidden-bound-test".to_string(),
6784 modules: Vec::new(),
6785 },
6786 pre_spawn: Vec::new(),
6787 })
6788 .timeout(Duration::from_secs(1))
6789 .build(),
6790 )
6791 .await?;
6792 runtime
6793 .spawn(
6794 SpawnMemberSpec::from_wire(
6795 "worker".to_string(),
6796 "rt:review:singleton:0".to_string(),
6797 Some("You are a hidden Review Agent.".into()),
6798 None,
6799 None,
6800 )
6801 .with_labels(BTreeMap::from([
6802 ("agent_identity".to_string(), "review:singleton".to_string()),
6803 ("role".to_string(), "delegate".to_string()),
6804 ("source_mob_id".to_string(), "upstream".to_string()),
6805 ])),
6806 )
6807 .await?;
6808 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6809 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6810 lease_provider: Arc::new(LocalLeaseProvider::new()),
6811 runtime_instance_id: "rpc-hidden-bound-test".to_string(),
6812 has_runtime_store: true,
6813 durability_policy: DurabilityPolicy::SyncWriteThrough,
6814 bridge: None,
6815 default_timeout: None,
6816 });
6817 let identity = AgentIdentity::parse("review:singleton")?;
6818 identity_rt
6819 .register(
6820 DurableAgentSpec {
6821 identity: identity.clone(),
6822 profile: meerkat_mob::ProfileName::from("worker"),
6823 addressability: AgentAddressability::Addressable,
6824 display_name: None,
6825 labels: BTreeMap::new(),
6826 context: None,
6827 additional_instructions: Vec::new(),
6828 initial_message: None,
6829 runtime_mode_override: None,
6830 backend: None,
6831 binding: None,
6832 },
6833 IdentityLifecycleState::Active,
6834 Some(ContinuityRecord {
6835 identity,
6836 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
6837 session_id: meerkat_core::types::SessionId::new(),
6838 generation: ContinuityGeneration::new(0),
6839 checkpoint_version: CheckpointVersion::new(0),
6840 }),
6841 None,
6842 )
6843 .await;
6844
6845 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
6846 let err =
6847 resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
6848 .await
6849 .expect_err("hidden registered live binding must not resolve");
6850 assert!(
6851 err.contains("identity hidden by policy"),
6852 "unexpected error for {requested_identity}: {err}"
6853 );
6854 }
6855
6856 Ok(())
6857 }
6858
6859 #[tokio::test]
6860 async fn live_only_hidden_alias_reports_policy_error()
6861 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6862 let temp_dir = tempfile::tempdir()?;
6863 let runtime = Box::pin(
6864 UnifiedRuntime::builder()
6865 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6866 .module_config(MobKitConfig {
6867 modules: Vec::new(),
6868 discovery: DiscoverySpec {
6869 namespace: "rpc-hidden-live-only-test".to_string(),
6870 modules: Vec::new(),
6871 },
6872 pre_spawn: Vec::new(),
6873 })
6874 .timeout(Duration::from_secs(1))
6875 .build(),
6876 )
6877 .await?;
6878 runtime
6879 .spawn(
6880 SpawnMemberSpec::from_wire(
6881 "worker".to_string(),
6882 "rt:review:singleton:0".to_string(),
6883 Some("You are a hidden Review Agent.".into()),
6884 None,
6885 None,
6886 )
6887 .with_labels(BTreeMap::from([
6888 ("agent_identity".to_string(), "review:singleton".to_string()),
6889 ("role".to_string(), "delegate".to_string()),
6890 ("source_mob_id".to_string(), "upstream".to_string()),
6891 ])),
6892 )
6893 .await?;
6894 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6895 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6896 lease_provider: Arc::new(LocalLeaseProvider::new()),
6897 runtime_instance_id: "rpc-hidden-live-only-test".to_string(),
6898 has_runtime_store: true,
6899 durability_policy: DurabilityPolicy::SyncWriteThrough,
6900 bridge: None,
6901 default_timeout: None,
6902 });
6903
6904 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
6905 let err =
6906 resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
6907 .await
6908 .expect_err("hidden live-only alias must not collapse into unknown identity");
6909 assert!(
6910 err.contains("identity hidden by policy"),
6911 "unexpected error for {requested_identity}: {err}"
6912 );
6913 }
6914
6915 let identity_ctx = IdentityFirstContext {
6916 runtime: Arc::new(identity_rt),
6917 roster_provider: Arc::new(EmptyRosterProvider),
6918 topology_provider: None,
6919 customizer: None,
6920 agent_memory_provider: None,
6921 mob_definition: None,
6922 };
6923 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
6924 let response: Value = serde_json::from_str(
6925 &handle_unified_rpc_json(
6926 &runtime,
6927 &json!({
6928 "jsonrpc": "2.0",
6929 "id": 1,
6930 "method": "mobkit/status_identity",
6931 "params": { "identity": requested_identity },
6932 })
6933 .to_string(),
6934 Duration::from_secs(1),
6935 None,
6936 Some(&identity_ctx),
6937 )
6938 .await,
6939 )?;
6940 assert_eq!(
6941 response["error"]["data"]["kind"],
6942 json!("identity_hidden_by_policy"),
6943 "unexpected hidden response for {requested_identity}: {response:#?}"
6944 );
6945 }
6946
6947 Ok(())
6948 }
6949
6950 #[tokio::test]
6951 async fn live_only_resolution_rejects_runtime_member_bound_to_other_durable_identity()
6952 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
6953 let temp_dir = tempfile::tempdir()?;
6954 let runtime = Box::pin(
6955 UnifiedRuntime::builder()
6956 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
6957 .module_config(MobKitConfig {
6958 modules: Vec::new(),
6959 discovery: DiscoverySpec {
6960 namespace: "rpc-identity-alias-test".to_string(),
6961 modules: Vec::new(),
6962 },
6963 pre_spawn: Vec::new(),
6964 })
6965 .timeout(Duration::from_secs(1))
6966 .build(),
6967 )
6968 .await?;
6969 let mut labels = BTreeMap::new();
6970 labels.insert("agent_identity".to_string(), "other:singleton".to_string());
6971 runtime
6972 .spawn(
6973 SpawnMemberSpec::from_wire(
6974 "worker".to_string(),
6975 "rt:review:singleton:0".to_string(),
6976 Some("You are a wrong-projected Review Agent.".into()),
6977 None,
6978 None,
6979 )
6980 .with_labels(labels),
6981 )
6982 .await?;
6983
6984 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
6985 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
6986 lease_provider: Arc::new(LocalLeaseProvider::new()),
6987 runtime_instance_id: "rpc-identity-alias-test".to_string(),
6988 has_runtime_store: true,
6989 durability_policy: DurabilityPolicy::SyncWriteThrough,
6990 bridge: None,
6991 default_timeout: None,
6992 });
6993 let identity = AgentIdentity::parse("review:singleton")?;
6994 let record = ContinuityRecord {
6995 identity: identity.clone(),
6996 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
6997 session_id: meerkat_core::types::SessionId::new(),
6998 generation: ContinuityGeneration::new(0),
6999 checkpoint_version: CheckpointVersion::new(0),
7000 };
7001 identity_rt
7002 .register(
7003 DurableAgentSpec {
7004 identity: identity.clone(),
7005 profile: meerkat_mob::ProfileName::from("worker"),
7006 addressability: AgentAddressability::Addressable,
7007 display_name: None,
7008 labels: BTreeMap::new(),
7009 context: None,
7010 additional_instructions: Vec::new(),
7011 initial_message: None,
7012 runtime_mode_override: None,
7013 backend: None,
7014 binding: None,
7015 },
7016 IdentityLifecycleState::Active,
7017 Some(record),
7018 Some(LeaseGrant {
7019 identity,
7020 fencing_token: FencingToken::new(1),
7021 ttl: Duration::from_mins(1),
7022 }),
7023 )
7024 .await;
7025
7026 let err = resolve_rpc_identity_control_target(&runtime, &identity_rt, "other:singleton")
7027 .await
7028 .expect_err("wrong-projected live alias must not resolve as live-only");
7029 assert!(
7030 err.contains("identity runtime binding belongs to review:singleton"),
7031 "unexpected error: {err}"
7032 );
7033
7034 Ok(())
7035 }
7036
7037 #[tokio::test]
7046 async fn send_message_resolves_bare_durable_identity_through_identity_bridge()
7047 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
7048 let temp_dir = tempfile::tempdir()?;
7049 let runtime = Box::pin(
7050 UnifiedRuntime::builder()
7051 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
7052 .module_config(MobKitConfig {
7053 modules: Vec::new(),
7054 discovery: DiscoverySpec {
7055 namespace: "rpc-send-message-identity-bridge-test".to_string(),
7056 modules: Vec::new(),
7057 },
7058 pre_spawn: Vec::new(),
7059 })
7060 .timeout(Duration::from_secs(5))
7061 .build(),
7062 )
7063 .await?;
7064
7065 for (runtime_id, durable) in [
7068 ("rt:atlas-base-001:0", "atlas-base-001"),
7069 ("rt:draco-base-001:0", "draco-base-001"),
7070 ] {
7071 runtime
7072 .spawn(
7073 SpawnMemberSpec::from_wire(
7074 "worker".to_string(),
7075 runtime_id.to_string(),
7076 Some("You are a swarm base agent.".into()),
7077 None,
7078 None,
7079 )
7080 .with_labels(BTreeMap::from([(
7081 "agent_identity".to_string(),
7082 durable.to_string(),
7083 )])),
7084 )
7085 .await?;
7086 }
7087 runtime
7090 .spawn(SpawnMemberSpec::from_wire(
7091 "worker".to_string(),
7092 "draco-base-001".to_string(),
7093 Some("You are the raw roster member.".into()),
7094 None,
7095 None,
7096 ))
7097 .await?;
7098
7099 let session_service = runtime
7100 .mob_runtime()
7101 .session_service()
7102 .cloned()
7103 .expect("test mob spec has a session service");
7104 let bridge: Arc<dyn crate::identity_first::SessionBridge> = Arc::new(
7105 crate::identity_first::MobSessionBridge::with_session_service(
7106 runtime.mob_handle(),
7107 session_service,
7108 ),
7109 );
7110 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
7111 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
7112 lease_provider: Arc::new(LocalLeaseProvider::new()),
7113 runtime_instance_id: "rpc-send-message-identity-bridge-test".to_string(),
7114 has_runtime_store: true,
7115 durability_policy: DurabilityPolicy::SyncWriteThrough,
7116 bridge: Some(bridge),
7117 default_timeout: None,
7118 })
7119 .with_runtime_services(crate::identity_first::AgentRuntimeServices::new(
7120 runtime.mob_handle(),
7121 ));
7122 for (durable, runtime_id) in [
7123 ("atlas-base-001", "rt:atlas-base-001:0"),
7124 ("draco-base-001", "rt:draco-base-001:0"),
7125 ] {
7126 let identity = AgentIdentity::parse(durable)?;
7127 let session_id = runtime
7128 .mob_handle()
7129 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(runtime_id))
7130 .await
7131 .unwrap_or_else(meerkat_core::types::SessionId::new);
7132 identity_rt
7133 .register(
7134 DurableAgentSpec {
7135 identity: identity.clone(),
7136 profile: meerkat_mob::ProfileName::from("worker"),
7137 addressability: AgentAddressability::Addressable,
7138 display_name: None,
7139 labels: BTreeMap::new(),
7140 context: None,
7141 additional_instructions: Vec::new(),
7142 initial_message: None,
7143 runtime_mode_override: None,
7144 backend: None,
7145 binding: None,
7146 },
7147 IdentityLifecycleState::Active,
7148 Some(ContinuityRecord {
7149 identity: identity.clone(),
7150 agent_runtime_id: AgentRuntimeId::parse(runtime_id)?,
7151 session_id,
7152 generation: ContinuityGeneration::new(0),
7153 checkpoint_version: CheckpointVersion::new(0),
7154 }),
7155 Some(LeaseGrant {
7156 identity,
7157 fencing_token: FencingToken::new(1),
7158 ttl: Duration::from_mins(1),
7159 }),
7160 )
7161 .await;
7162 }
7163 let identity_ctx = IdentityFirstContext {
7164 runtime: Arc::new(identity_rt),
7165 roster_provider: Arc::new(EmptyRosterProvider),
7166 topology_provider: None,
7167 customizer: None,
7168 agent_memory_provider: None,
7169 mob_definition: None,
7170 };
7171
7172 let send = |id: u64, params: Value| {
7173 let runtime = &runtime;
7174 let identity_ctx = &identity_ctx;
7175 async move {
7176 let raw = handle_unified_rpc_json(
7177 runtime,
7178 &json!({
7179 "jsonrpc": "2.0",
7180 "id": id,
7181 "method": "mobkit/send_message",
7182 "params": params,
7183 })
7184 .to_string(),
7185 Duration::from_secs(10),
7186 None,
7187 Some(identity_ctx),
7188 )
7189 .await;
7190 serde_json::from_str::<Value>(&raw)
7191 }
7192 };
7193
7194 let response = send(
7197 1,
7198 json!({ "member_id": "atlas-base-001", "message": "status check" }),
7199 )
7200 .await?;
7201 assert!(
7202 response["error"].is_null(),
7203 "bare identity send must bridge-resolve: {response:#?}"
7204 );
7205 assert_eq!(response["result"]["accepted"], json!(true));
7206 assert_eq!(response["result"]["member_id"], json!("atlas-base-001"));
7207 let atlas_session = runtime
7208 .mob_handle()
7209 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
7210 "rt:atlas-base-001:0",
7211 ))
7212 .await
7213 .expect("atlas runtime member has a bridge session after send")
7214 .to_string();
7215 assert_eq!(response["result"]["session_id"], json!(atlas_session));
7216
7217 let response = send(
7219 2,
7220 json!({
7221 "member_id": "atlas-base-001",
7222 "message": "steer: stand down",
7223 "handling_mode": "steer",
7224 }),
7225 )
7226 .await?;
7227 assert!(
7228 response["error"].is_null(),
7229 "bare identity steer must bridge-resolve: {response:#?}"
7230 );
7231 assert_eq!(response["result"]["accepted"], json!(true));
7232
7233 let response = send(
7237 3,
7238 json!({ "member_id": "draco-base-001", "message": "raw roster delivery" }),
7239 )
7240 .await?;
7241 assert!(
7242 response["error"].is_null(),
7243 "exact roster member send must keep raw semantics: {response:#?}"
7244 );
7245 assert_eq!(response["result"]["accepted"], json!(true));
7246 let draco_raw_session = runtime
7247 .mob_handle()
7248 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id("draco-base-001"))
7249 .await
7250 .expect("bare draco member has a bridge session after send")
7251 .to_string();
7252 assert_eq!(response["result"]["session_id"], json!(draco_raw_session));
7253 if let Some(draco_rt_session) = runtime
7254 .mob_handle()
7255 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
7256 "rt:draco-base-001:0",
7257 ))
7258 .await
7259 {
7260 assert_ne!(
7261 response["result"]["session_id"],
7262 json!(draco_rt_session.to_string()),
7263 "exact member id match must not be shadowed by identity resolution"
7264 );
7265 }
7266
7267 let response = send(
7269 4,
7270 json!({ "member_id": "phantom-base-999", "message": "nobody home" }),
7271 )
7272 .await?;
7273 assert_eq!(response["error"]["code"], json!(-32000), "{response:#?}");
7274 let message = response["error"]["message"]
7275 .as_str()
7276 .expect("error message");
7277 assert!(
7278 message.starts_with("send_message failed:"),
7279 "unexpected error message: {message}"
7280 );
7281
7282 Ok(())
7283 }
7284
7285 #[tokio::test]
7292 async fn member_state_wire_vocabulary_is_lowercase_on_both_surfaces()
7293 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
7294 let temp_dir = tempfile::tempdir()?;
7295 let runtime = Box::pin(
7296 UnifiedRuntime::builder()
7297 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
7298 .module_config(MobKitConfig {
7299 modules: Vec::new(),
7300 discovery: DiscoverySpec {
7301 namespace: "rpc-state-vocabulary-test".to_string(),
7302 modules: Vec::new(),
7303 },
7304 pre_spawn: Vec::new(),
7305 })
7306 .timeout(Duration::from_secs(5))
7307 .build(),
7308 )
7309 .await?;
7310 runtime
7311 .spawn(SpawnMemberSpec::from_wire(
7312 "worker".to_string(),
7313 "worker-one".to_string(),
7314 None,
7315 None,
7316 None,
7317 ))
7318 .await?;
7319
7320 let raw = handle_unified_rpc_json(
7322 &runtime,
7323 &json!({
7324 "jsonrpc": "2.0",
7325 "id": 1,
7326 "method": "mobkit/get_member",
7327 "params": { "member_id": "worker-one" },
7328 })
7329 .to_string(),
7330 Duration::from_secs(5),
7331 None,
7332 None,
7333 )
7334 .await;
7335 let response: Value = serde_json::from_str(&raw)?;
7336 assert_eq!(
7337 response["result"]["state"],
7338 json!("active"),
7339 "member rows must speak the lowercase SDK vocabulary: {response:#?}"
7340 );
7341
7342 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
7344 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
7345 lease_provider: Arc::new(LocalLeaseProvider::new()),
7346 runtime_instance_id: "rpc-state-vocabulary-test".to_string(),
7347 has_runtime_store: true,
7348 durability_policy: DurabilityPolicy::SyncWriteThrough,
7349 bridge: None,
7350 default_timeout: None,
7351 });
7352 let identity = AgentIdentity::parse("review:singleton")?;
7353 identity_rt
7354 .register(
7355 DurableAgentSpec {
7356 identity: identity.clone(),
7357 profile: meerkat_mob::ProfileName::from("worker"),
7358 addressability: AgentAddressability::Addressable,
7359 display_name: None,
7360 labels: BTreeMap::new(),
7361 context: None,
7362 additional_instructions: Vec::new(),
7363 initial_message: None,
7364 runtime_mode_override: None,
7365 backend: None,
7366 binding: None,
7367 },
7368 IdentityLifecycleState::Active,
7369 Some(ContinuityRecord {
7370 identity: identity.clone(),
7371 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
7372 session_id: meerkat_core::types::SessionId::new(),
7373 generation: ContinuityGeneration::new(0),
7374 checkpoint_version: CheckpointVersion::new(0),
7375 }),
7376 Some(LeaseGrant {
7377 identity,
7378 fencing_token: FencingToken::new(1),
7379 ttl: Duration::from_mins(1),
7380 }),
7381 )
7382 .await;
7383 let identity_ctx = IdentityFirstContext {
7384 runtime: Arc::new(identity_rt),
7385 roster_provider: Arc::new(EmptyRosterProvider),
7386 topology_provider: None,
7387 customizer: None,
7388 agent_memory_provider: None,
7389 mob_definition: None,
7390 };
7391 let raw = handle_unified_rpc_json(
7392 &runtime,
7393 &json!({
7394 "jsonrpc": "2.0",
7395 "id": 2,
7396 "method": "mobkit/status_identity",
7397 "params": { "identity": "review:singleton" },
7398 })
7399 .to_string(),
7400 Duration::from_secs(5),
7401 None,
7402 Some(&identity_ctx),
7403 )
7404 .await;
7405 let response: Value = serde_json::from_str(&raw)?;
7406 assert_eq!(
7407 response["result"]["state"],
7408 json!("active"),
7409 "identity status must speak the same lowercase vocabulary: {response:#?}"
7410 );
7411
7412 Ok(())
7413 }
7414
7415 #[tokio::test]
7424 async fn send_message_pins_baseline_member_over_identity_fallback()
7425 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
7426 let temp_dir = tempfile::tempdir()?;
7427 let runtime = Box::pin(
7428 UnifiedRuntime::builder()
7429 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
7430 .module_config(MobKitConfig {
7431 modules: Vec::new(),
7432 discovery: DiscoverySpec {
7433 namespace: "rpc-send-message-baseline-pin-test".to_string(),
7434 modules: Vec::new(),
7435 },
7436 pre_spawn: Vec::new(),
7437 })
7438 .timeout(Duration::from_secs(5))
7439 .build(),
7440 )
7441 .await?;
7442
7443 for member in ["rt:draco-base-001:0", "draco-base-001"] {
7446 runtime
7447 .spawn(SpawnMemberSpec::from_wire(
7448 "worker".to_string(),
7449 member.to_string(),
7450 None,
7451 None,
7452 None,
7453 ))
7454 .await?;
7455 }
7456
7457 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
7458 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
7459 lease_provider: Arc::new(LocalLeaseProvider::new()),
7460 runtime_instance_id: "rpc-send-message-baseline-pin-test".to_string(),
7461 has_runtime_store: true,
7462 durability_policy: DurabilityPolicy::SyncWriteThrough,
7463 bridge: None,
7464 default_timeout: None,
7465 });
7466 let identity = AgentIdentity::parse("draco-base-001")?;
7467 identity_rt
7468 .register(
7469 DurableAgentSpec {
7470 identity: identity.clone(),
7471 profile: meerkat_mob::ProfileName::from("worker"),
7472 addressability: AgentAddressability::Addressable,
7473 display_name: None,
7474 labels: BTreeMap::new(),
7475 context: None,
7476 additional_instructions: Vec::new(),
7477 initial_message: None,
7478 runtime_mode_override: None,
7479 backend: None,
7480 binding: None,
7481 },
7482 IdentityLifecycleState::Active,
7483 Some(ContinuityRecord {
7484 identity: identity.clone(),
7485 agent_runtime_id: AgentRuntimeId::parse("rt:draco-base-001:0")?,
7486 session_id: meerkat_core::types::SessionId::new(),
7487 generation: ContinuityGeneration::new(0),
7488 checkpoint_version: CheckpointVersion::new(0),
7489 }),
7490 Some(LeaseGrant {
7491 identity,
7492 fencing_token: FencingToken::new(1),
7493 ttl: Duration::from_mins(1),
7494 }),
7495 )
7496 .await;
7497 let identity_ctx = IdentityFirstContext {
7498 runtime: Arc::new(identity_rt),
7499 roster_provider: Arc::new(EmptyRosterProvider),
7500 topology_provider: None,
7501 customizer: None,
7502 agent_memory_provider: None,
7503 mob_definition: None,
7504 };
7505
7506 runtime
7508 .mob_runtime()
7509 .set_baseline_member_specs(vec![SpawnMemberSpec::new(
7510 meerkat_mob::ProfileName::from("worker"),
7511 meerkat_mob::AgentIdentity::from("draco-base-001"),
7512 )])
7513 .await;
7514 runtime
7517 .mob_handle()
7518 .retire(meerkat_mob::AgentIdentity::from("draco-base-001"))
7519 .await?;
7520
7521 let raw = handle_unified_rpc_json(
7522 &runtime,
7523 &json!({
7524 "jsonrpc": "2.0",
7525 "id": 1,
7526 "method": "mobkit/send_message",
7527 "params": { "member_id": "draco-base-001", "message": "mid-reconcile send" },
7528 })
7529 .to_string(),
7530 Duration::from_secs(10),
7531 None,
7532 Some(&identity_ctx),
7533 )
7534 .await;
7535 let response: Value = serde_json::from_str(&raw)?;
7536 assert_eq!(
7537 response["error"]["code"],
7538 json!(-32000),
7539 "transiently-absent baseline member must keep raw member-id semantics \
7540 instead of silently delivering through the identity bridge: {response:#?}"
7541 );
7542 let message = response["error"]["message"]
7543 .as_str()
7544 .expect("error message");
7545 assert!(
7546 message.starts_with("send_message failed:"),
7547 "unexpected error message: {message}"
7548 );
7549
7550 Ok(())
7551 }
7552
7553 #[tokio::test]
7562 async fn retire_and_respawn_rpcs_succeed_for_idle_member()
7563 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
7564 let temp_dir = tempfile::tempdir()?;
7565 let runtime = Box::pin(
7566 UnifiedRuntime::builder()
7567 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
7568 .module_config(MobKitConfig {
7569 modules: Vec::new(),
7570 discovery: DiscoverySpec {
7571 namespace: "rpc-idle-lifecycle-test".to_string(),
7572 modules: Vec::new(),
7573 },
7574 pre_spawn: Vec::new(),
7575 })
7576 .timeout(Duration::from_secs(5))
7577 .build(),
7578 )
7579 .await?;
7580
7581 for member in ["worker-one", "worker-two"] {
7582 runtime
7583 .spawn(SpawnMemberSpec::from_wire(
7584 "worker".to_string(),
7585 member.to_string(),
7586 None,
7587 None,
7588 None,
7589 ))
7590 .await?;
7591 }
7592
7593 let send = |id: u64, method: &'static str, params: Value| {
7594 let runtime = &runtime;
7595 async move {
7596 let raw = handle_unified_rpc_json(
7597 runtime,
7598 &json!({
7599 "jsonrpc": "2.0",
7600 "id": id,
7601 "method": method,
7602 "params": params,
7603 })
7604 .to_string(),
7605 Duration::from_secs(10),
7606 None,
7607 None,
7608 )
7609 .await;
7610 serde_json::from_str::<Value>(&raw)
7611 }
7612 };
7613
7614 let response = send(
7616 1,
7617 "mobkit/retire_member",
7618 json!({"member_id": "worker-one"}),
7619 )
7620 .await?;
7621 assert!(
7622 response["error"].is_null(),
7623 "retire_member must succeed for an idle member: {response:#?}"
7624 );
7625 assert_eq!(response["result"]["accepted"], json!(true));
7626 assert!(
7627 !runtime
7628 .mob_handle()
7629 .list_members_including_retiring()
7630 .await
7631 .iter()
7632 .any(|entry| entry.agent_identity.as_str() == "worker-one"),
7633 "retired member must leave the roster"
7634 );
7635
7636 let response = send(
7639 2,
7640 "mobkit/respawn_member",
7641 json!({"member_id": "worker-two"}),
7642 )
7643 .await?;
7644 assert!(
7645 response["error"].is_null(),
7646 "respawn_member must succeed for an idle member: {response:#?}"
7647 );
7648 assert_eq!(response["result"]["accepted"], json!(true));
7649 let members = runtime.mob_handle().list_members_including_retiring().await;
7650 let worker_two = members
7651 .iter()
7652 .find(|entry| entry.agent_identity.as_str() == "worker-two")
7653 .expect("respawned member must remain in the roster");
7654 assert_eq!(
7655 worker_two.status,
7656 meerkat_mob::MobMemberStatus::Active,
7657 "respawned member must be active, not wedged in retiring"
7658 );
7659
7660 Ok(())
7661 }
7662}