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