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 ElephantMemoryStoreError, GatingDecideError, GatingDecideRequest, GatingDecision,
17 GatingEvaluateRequest, GatingRiskTier, MemoryIndexError, MemoryIndexRequest,
18 MemoryQueryRequest, MobkitRuntimeHandle, ModuleRouteError, ModuleRouteRequest,
19 ROUTING_RETRY_MAX_CAP, RoutingResolveError, RoutingResolveRequest, RuntimeDecisionState,
20 RuntimeRoute, RuntimeRouteMutationError, ScheduleDefinition, ScheduleValidationError,
21 SessionPersistenceRow, SubscribeError, SubscribeRequest, SubscribeScope,
22 handle_console_rest_json_route, route_module_call, validate_schedules,
23};
24use crate::unified_runtime::{EventQuery, UnifiedRuntime};
25
26mod console_ingress;
27mod gating_methods;
28mod 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_memory_index_params, parse_memory_query_params,
44 parse_memory_stores_params,
45};
46use routing_delivery_methods::{
47 RoutingDeliveryParamsError, parse_delivery_history_params, parse_delivery_send_params,
48 parse_routing_resolve_params, parse_routing_route_add_params,
49 parse_routing_route_delete_params, parse_routing_routes_list_params,
50};
51use scheduling_methods::{format_schedule_validation_error, parse_scheduling_params};
52use session_store_methods::{
53 BigQuerySessionStoreRpcError, format_bigquery_store_error, parse_bigquery_session_store_params,
54 run_bigquery_session_store_request,
55};
56use subscribe_methods::{SubscribeParamsError, parse_subscribe_request};
57
58pub const JSONRPC_VERSION: &str = "2.0";
59pub const MOBKIT_CONTRACT_VERSION: &str = "0.4.0";
60pub const MAX_SCHEDULES_PER_REQUEST: usize = 256;
61pub(crate) const MOBPACK_AUTHORING_METHODS: &[&str] = &[
62 "mobkit/mobpacks/schema",
63 "mobkit/mobpacks/catalogs",
64 "mobkit/tools/catalog",
65 "mobkit/skills/catalog",
66 "mobkit/agent_definitions/list",
67 "mobkit/mobpacks/templates",
68 "mobkit/mobpacks/validate",
69 "mobkit/mobpacks/source",
70 "mobkit/mobpacks/export",
71 "mobkit/mobpacks/import",
72 "mobkit/mobpacks/list",
73 "mobkit/mobpacks/get",
74 "mobkit/mobpacks/create",
75 "mobkit/mobpacks/save",
76 "mobkit/mobpacks/delete",
77 "mobkit/mobpacks/undo",
78 "mobkit/mobpacks/redo",
79 "mobkit/mobpacks/apply_operation",
80 "mobkit/mobpacks/graph_projection",
81 "mobkit/mobpacks/graph_to_flow",
82 "mobkit/mobpacks/deploy_command",
83 "mobkit/mobpacks/deploy",
84];
85
86pub(crate) fn mobpack_authoring_capabilities() -> Value {
87 serde_json::json!({
88 "domain": "mobpack_authoring",
89 "runtime_mutation": false,
90 "host_mutation_methods": {
91 "mobkit/mobpacks/deploy": "when execute=true, writes a mobpack archive and runs rkat mob run on the host",
92 "mobkit/mobpacks/validate": "when rkat_validate=true, writes a mobpack archive and runs rkat mob validate on the host"
93 },
94 "deploy_command": "rkat mob run",
95 "methods": MOBPACK_AUTHORING_METHODS,
96 "operations": crate::mobpack::mobpack_authoring_operations(),
97 })
98}
99
100async fn mobpack_runtime_catalog_state(
101 runtime: &UnifiedRuntime,
102) -> crate::mobpack::MobpackRuntimeCatalogState {
103 let loaded_modules = runtime.loaded_modules().await;
104 let runtime_flow_rows = crate::mobpack::runtime_flow_registry_rows_from_definition(
105 runtime.mob_handle().definition(),
106 );
107 let runtime_agent_definition_sources =
108 crate::mobpack::runtime_agent_definition_sources_from_definition(
109 runtime.mob_handle().definition(),
110 );
111 let runtime_skill_realms =
112 crate::mobpack::runtime_skill_realms_from_definition(runtime.mob_handle().definition());
113 let mut runtime_methods = vec![
114 "mobkit/capabilities".to_string(),
115 "mobkit/models/catalog".to_string(),
116 "mobkit/spawn_member".to_string(),
117 "mobkit/list_members".to_string(),
118 "mobkit/get_member".to_string(),
119 "mobkit/run_flow".to_string(),
120 "mobkit/list_flows".to_string(),
121 "mobkit/list_runs".to_string(),
122 ];
123 runtime_methods.extend(
124 MOBPACK_AUTHORING_METHODS
125 .iter()
126 .map(std::string::ToString::to_string),
127 );
128 if runtime.has_contact_directory() {
129 runtime_methods.push("mobkit/cross_mob/directory".to_string());
130 }
131 if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
132 runtime_methods.extend([
133 "mobkit/cross_mob/wire".to_string(),
134 "mobkit/cross_mob/unwire".to_string(),
135 "mobkit/cross_mob/send".to_string(),
136 ]);
137 }
138 crate::mobpack::MobpackRuntimeCatalogState {
139 loaded_modules,
140 runtime_methods,
141 has_contact_directory: runtime.has_contact_directory(),
142 has_peer_mob_handles: runtime.has_peer_mob_handles().await,
143 has_inproc_contacts: runtime.has_inproc_contacts(),
144 runtime_flow_rows,
145 runtime_agent_definition_sources,
146 runtime_skill_realms,
147 }
148}
149
150async fn handle_unified_mobpack_authoring_rpc(
151 runtime: &UnifiedRuntime,
152 method: &str,
153 params: &Value,
154 response_id: Value,
155) -> JsonRpcResponse {
156 let runtime_catalog_state = match method {
157 "mobkit/mobpacks/schema"
158 | "mobkit/mobpacks/catalogs"
159 | "mobkit/tools/catalog"
160 | "mobkit/skills/catalog"
161 | "mobkit/agent_definitions/list"
162 | "mobkit/mobpacks/templates"
163 | "mobkit/mobpacks/list"
164 | "mobkit/mobpacks/get"
165 | "mobkit/mobpacks/apply_operation" => Some(mobpack_runtime_catalog_state(runtime).await),
166 _ => None,
167 };
168 let result = match method {
169 "mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
170 runtime_catalog_state.as_ref(),
171 )),
172 "mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
173 runtime_catalog_state.as_ref(),
174 )),
175 "mobkit/skills/catalog" => Ok(
176 crate::mobpack::mobpack_skills_catalog_response_with_runtime(
177 runtime_catalog_state.as_ref(),
178 ),
179 ),
180 "mobkit/agent_definitions/list" => Ok(
181 crate::mobpack::mobpack_agent_definitions_response_with_runtime(
182 runtime_catalog_state.as_ref(),
183 ),
184 ),
185 "mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
186 runtime_catalog_state.as_ref(),
187 )),
188 _ => {
189 return handle_mobpack_authoring_rpc_with_runtime(
190 method,
191 params,
192 response_id.clone(),
193 runtime_catalog_state.as_ref(),
194 )
195 .unwrap_or_else(|| JsonRpcResponse {
196 jsonrpc: JSONRPC_VERSION.to_string(),
197 id: response_id,
198 result: None,
199 error: Some(JsonRpcError {
200 code: -32601,
201 message: "Method not found".to_string(),
202 data: None,
203 }),
204 });
205 }
206 };
207 match result {
208 Ok(result) => JsonRpcResponse {
209 jsonrpc: JSONRPC_VERSION.to_string(),
210 id: response_id,
211 result: Some(result),
212 error: None,
213 },
214 Err(message) => JsonRpcResponse {
215 jsonrpc: JSONRPC_VERSION.to_string(),
216 id: response_id,
217 result: None,
218 error: Some(JsonRpcError {
219 code: -32602,
220 message,
221 data: None,
222 }),
223 },
224 }
225}
226
227pub(crate) fn handle_mobpack_authoring_rpc(
228 method: &str,
229 params: &Value,
230 response_id: Value,
231) -> Option<JsonRpcResponse> {
232 handle_mobpack_authoring_rpc_with_runtime(method, params, response_id, None)
233}
234
235pub(crate) fn handle_mobpack_authoring_rpc_with_runtime(
236 method: &str,
237 params: &Value,
238 response_id: Value,
239 runtime: Option<&crate::mobpack::MobpackRuntimeCatalogState>,
240) -> Option<JsonRpcResponse> {
241 let result = match method {
242 "mobkit/mobpacks/schema" => Ok(crate::mobpack::mobpack_schema_response_with_runtime(
243 runtime,
244 )),
245 "mobkit/mobpacks/catalogs" => Ok(crate::mobpack::mobpack_catalogs_response_with_runtime(
246 runtime,
247 )),
248 "mobkit/tools/catalog" => Ok(crate::mobpack::mobpack_tools_catalog_response_with_runtime(
249 runtime,
250 )),
251 "mobkit/skills/catalog" => {
252 Ok(crate::mobpack::mobpack_skills_catalog_response_with_runtime(runtime))
253 }
254 "mobkit/agent_definitions/list" => {
255 Ok(crate::mobpack::mobpack_agent_definitions_response_with_runtime(runtime))
256 }
257 "mobkit/mobpacks/templates" => Ok(crate::mobpack::mobpack_templates_response_with_runtime(
258 runtime,
259 )),
260 "mobkit/mobpacks/validate" => crate::mobpack::validate_mobpack(params)
261 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
262 "mobkit/mobpacks/source" => crate::mobpack::source_mobpack(params)
263 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
264 "mobkit/mobpacks/export" => crate::mobpack::export_mobpack(params)
265 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
266 "mobkit/mobpacks/import" => crate::mobpack::import_mobpack(params),
267 "mobkit/mobpacks/list" => crate::mobpack::list_mobpack_drafts_with_runtime(params, runtime),
268 "mobkit/mobpacks/get" => crate::mobpack::get_mobpack_draft_with_runtime(params, runtime),
269 "mobkit/mobpacks/create" => crate::mobpack::create_mobpack_draft(params),
270 "mobkit/mobpacks/save" => crate::mobpack::save_mobpack_draft(params),
271 "mobkit/mobpacks/delete" => crate::mobpack::delete_mobpack_draft(params),
272 "mobkit/mobpacks/undo" => crate::mobpack::undo_mobpack_draft(params),
273 "mobkit/mobpacks/redo" => crate::mobpack::redo_mobpack_draft(params),
274 "mobkit/mobpacks/apply_operation" => {
275 crate::mobpack::apply_mobpack_authoring_operation_with_runtime(params, runtime)
276 }
277 "mobkit/mobpacks/graph_projection" => crate::mobpack::graph_projection_mobpack(params),
278 "mobkit/mobpacks/graph_to_flow" => crate::mobpack::graph_to_flow_mobpack(params),
279 "mobkit/mobpacks/deploy_command" => crate::mobpack::deploy_command_preview(params)
280 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
281 "mobkit/mobpacks/deploy" => crate::mobpack::deploy_mobpack(params)
282 .and_then(|result| serde_json::to_value(result).map_err(|err| err.to_string())),
283 _ => return None,
284 };
285 Some(match result {
286 Ok(result) => JsonRpcResponse {
287 jsonrpc: JSONRPC_VERSION.to_string(),
288 id: response_id,
289 result: Some(result),
290 error: None,
291 },
292 Err(message) => JsonRpcResponse {
293 jsonrpc: JSONRPC_VERSION.to_string(),
294 id: response_id,
295 result: None,
296 error: Some(JsonRpcError {
297 code: -32602,
298 message,
299 data: None,
300 }),
301 },
302 })
303}
304
305#[derive(Debug, Clone, PartialEq, Eq)]
306pub enum RpcCapabilitiesError {
307 InvalidJson,
308 InvalidSchema,
309 MissingContractVersion,
310 InvalidContractVersion,
311}
312
313impl std::fmt::Display for RpcCapabilitiesError {
314 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
315 match self {
316 Self::InvalidJson => write!(f, "invalid JSON"),
317 Self::InvalidSchema => write!(f, "invalid schema"),
318 Self::MissingContractVersion => write!(f, "missing contract version"),
319 Self::InvalidContractVersion => write!(f, "invalid contract version"),
320 }
321 }
322}
323
324impl std::error::Error for RpcCapabilitiesError {}
325
326#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
327pub struct RpcCapabilities {
328 pub contract_version: String,
329 #[serde(flatten)]
330 pub extra: BTreeMap<String, Value>,
331}
332
333pub fn parse_rpc_capabilities(line: &str) -> Result<RpcCapabilities, RpcCapabilitiesError> {
334 let raw: Value = serde_json::from_str(line).map_err(|_| RpcCapabilitiesError::InvalidJson)?;
335 let object = raw.as_object().ok_or(RpcCapabilitiesError::InvalidSchema)?;
336 let contract = object
337 .get("contract_version")
338 .ok_or(RpcCapabilitiesError::MissingContractVersion)?;
339 let contract_str = contract
340 .as_str()
341 .ok_or(RpcCapabilitiesError::InvalidContractVersion)?;
342 if contract_str.trim().is_empty() {
343 return Err(RpcCapabilitiesError::InvalidContractVersion);
344 }
345 serde_json::from_value(raw).map_err(|_| RpcCapabilitiesError::InvalidSchema)
346}
347
348pub const MOB_EVENTS_STALE_CURSOR_CODE: i64 = -32010;
355
356pub const MEMORY_BACKEND_UNAVAILABLE_CODE: i64 = -32012;
363
364pub const CONSOLE_TIMELINE_REPLAY_UNAVAILABLE_CODE: i64 = -32013;
369
370#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
371pub struct JsonRpcRequest {
372 pub jsonrpc: String,
373 #[serde(default)]
374 pub id: Option<Value>,
375 pub method: String,
376 #[serde(default)]
377 pub params: Value,
378}
379
380#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
381pub struct JsonRpcError {
382 pub code: i64,
383 pub message: String,
384 #[serde(default, skip_serializing_if = "Option::is_none")]
389 pub data: Option<Value>,
390}
391
392impl JsonRpcError {
393 pub fn new(code: i64, message: impl Into<String>) -> Self {
394 Self {
395 code,
396 message: message.into(),
397 data: None,
398 }
399 }
400
401 #[must_use]
402 pub fn with_data(mut self, data: Value) -> Self {
403 self.data = Some(data);
404 self
405 }
406}
407
408#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
409pub struct JsonRpcResponse {
410 pub jsonrpc: String,
411 pub id: Value,
412 #[serde(skip_serializing_if = "Option::is_none")]
413 pub result: Option<Value>,
414 #[serde(skip_serializing_if = "Option::is_none")]
415 pub error: Option<JsonRpcError>,
416}
417
418pub fn handle_mobkit_rpc_json(
419 runtime: &mut MobkitRuntimeHandle,
420 request_json: &str,
421 timeout: Duration,
422) -> String {
423 let raw_request: Value = match serde_json::from_str(request_json) {
424 Ok(raw_request) => raw_request,
425 Err(_) => {
426 return serialize_response(&JsonRpcResponse {
427 jsonrpc: JSONRPC_VERSION.to_string(),
428 id: Value::Null,
429 result: None,
430 error: Some(JsonRpcError {
431 code: -32700,
432 message: "Parse error".to_string(),
433 data: None,
434 }),
435 });
436 }
437 };
438 let response_id = raw_request
439 .as_object()
440 .and_then(|object| object.get("id"))
441 .cloned()
442 .unwrap_or(Value::Null);
443 let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
444 Ok(request) => request,
445 Err(_) => {
446 return serialize_response(&JsonRpcResponse {
447 jsonrpc: JSONRPC_VERSION.to_string(),
448 id: response_id,
449 result: None,
450 error: Some(JsonRpcError {
451 code: -32600,
452 message: "Invalid Request".to_string(),
453 data: None,
454 }),
455 });
456 }
457 };
458 let is_notification = request.id.is_none();
459 let response_id = request.id.clone().unwrap_or(Value::Null);
460
461 if request.jsonrpc != "2.0" {
462 let response = JsonRpcResponse {
463 jsonrpc: JSONRPC_VERSION.to_string(),
464 id: response_id,
465 result: None,
466 error: Some(JsonRpcError {
467 code: -32600,
468 message: "Invalid Request".to_string(),
469 data: None,
470 }),
471 };
472 return if is_notification {
473 String::new()
474 } else {
475 serialize_response(&response)
476 };
477 }
478
479 let response = match request.method.as_str() {
480 "mobkit/status" => JsonRpcResponse {
481 jsonrpc: JSONRPC_VERSION.to_string(),
482 id: response_id,
483 result: Some(serde_json::json!({
484 "contract_version": MOBKIT_CONTRACT_VERSION,
485 "running": runtime.is_running(),
486 "loaded_modules": runtime.loaded_modules(),
487 })),
488 error: None,
489 },
490 "mobkit/capabilities" => {
491 let mut methods = vec![
492 "mobkit/status",
493 "mobkit/capabilities",
494 "mobkit/reconcile",
495 "mobkit/spawn_member",
496 "mobkit/scheduling/evaluate",
497 "mobkit/scheduling/dispatch",
498 "mobkit/routing/resolve",
499 "mobkit/routing/routes/list",
500 "mobkit/routing/routes/add",
501 "mobkit/routing/routes/delete",
502 "mobkit/delivery/send",
503 "mobkit/delivery/history",
504 "mobkit/events/subscribe",
505 "mobkit/memory/stores",
506 "mobkit/memory/index",
507 "mobkit/memory/query",
508 "mobkit/session_store/bigquery",
509 "mobkit/gating/evaluate",
510 "mobkit/gating/pending",
511 "mobkit/gating/decide",
512 "mobkit/gating/audit",
513 "mobkit/call_tool",
514 "mobkit/models/catalog",
515 ];
516 methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
517 JsonRpcResponse {
518 jsonrpc: JSONRPC_VERSION.to_string(),
519 id: response_id,
520 result: Some(serde_json::json!({
521 "contract_version": MOBKIT_CONTRACT_VERSION,
522 "methods": methods,
523 "loaded_modules": runtime.loaded_modules(),
524 "runtime_capabilities": {
525 "can_spawn_members": false,
526 "can_send_messages": false,
527 "can_wire_members": false,
528 "can_retire_members": false,
529 "available_spawn_modes": ["module"],
530 },
531 "authoring_capabilities": mobpack_authoring_capabilities(),
532 })),
533 error: None,
534 }
535 }
536 "mobkit/models/catalog" => JsonRpcResponse {
537 jsonrpc: JSONRPC_VERSION.to_string(),
538 id: response_id,
539 result: Some(build_models_catalog_result()),
540 error: None,
541 },
542 method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
543 handle_mobpack_authoring_rpc(method, &request.params, response_id.clone())
544 .unwrap_or_else(|| JsonRpcResponse {
545 jsonrpc: JSONRPC_VERSION.to_string(),
546 id: response_id,
547 result: None,
548 error: Some(JsonRpcError {
549 code: -32601,
550 message: "Method not found".to_string(),
551 data: None,
552 }),
553 })
554 }
555 "mobkit/reconcile" => {
556 let modules = match params::required_string_array(&request.params, "modules") {
557 Ok(m) => m,
558 Err(reason) => {
559 return serialize_response(&JsonRpcResponse {
560 jsonrpc: JSONRPC_VERSION.to_string(),
561 id: response_id,
562 result: None,
563 error: Some(JsonRpcError {
564 code: -32602,
565 message: format!("Invalid params: {reason}"),
566 data: None,
567 }),
568 });
569 }
570 };
571
572 match runtime.reconcile_modules(modules.clone(), timeout) {
573 Ok(added) => JsonRpcResponse {
574 jsonrpc: JSONRPC_VERSION.to_string(),
575 id: response_id,
576 result: Some(serde_json::json!({
577 "accepted": true,
578 "reconciled_modules": modules,
579 "added": added
580 })),
581 error: None,
582 },
583 Err(err) => JsonRpcResponse {
584 jsonrpc: JSONRPC_VERSION.to_string(),
585 id: response_id,
586 result: None,
587 error: Some(JsonRpcError {
588 code: -32602,
589 message: format!("Invalid params: {err:?}"),
590 data: None,
591 }),
592 },
593 }
594 }
595 "mobkit/spawn_member" => {
596 let module_id = request
597 .params
598 .get("module_id")
599 .and_then(Value::as_str)
600 .unwrap_or_default()
601 .to_string();
602 if module_id.is_empty() {
603 JsonRpcResponse {
604 jsonrpc: JSONRPC_VERSION.to_string(),
605 id: response_id,
606 result: None,
607 error: Some(JsonRpcError {
608 code: -32602,
609 message: "Invalid params: module_id required".to_string(),
610 data: None,
611 }),
612 }
613 } else {
614 match runtime.spawn_member(&module_id, timeout) {
615 Ok(()) => JsonRpcResponse {
616 jsonrpc: JSONRPC_VERSION.to_string(),
617 id: response_id,
618 result: Some(serde_json::json!({
619 "accepted": true,
620 "module_id": module_id
621 })),
622 error: None,
623 },
624 Err(err) => JsonRpcResponse {
625 jsonrpc: JSONRPC_VERSION.to_string(),
626 id: response_id,
627 result: None,
628 error: Some(JsonRpcError {
629 code: -32602,
630 message: format!("Invalid params: {err:?}"),
631 data: None,
632 }),
633 },
634 }
635 }
636 }
637 "mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
638 Ok((schedules, tick_ms)) => match runtime.evaluate_schedule_tick(&schedules, tick_ms) {
639 Ok(evaluation) => JsonRpcResponse {
640 jsonrpc: JSONRPC_VERSION.to_string(),
641 id: response_id,
642 result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
643 error: None,
644 },
645 Err(err) => JsonRpcResponse {
646 jsonrpc: JSONRPC_VERSION.to_string(),
647 id: response_id,
648 result: None,
649 error: Some(JsonRpcError {
650 code: -32602,
651 message: format!(
652 "Invalid params: {}",
653 format_schedule_validation_error(err)
654 ),
655 data: None,
656 }),
657 },
658 },
659 Err(message) => JsonRpcResponse {
660 jsonrpc: JSONRPC_VERSION.to_string(),
661 id: response_id,
662 result: None,
663 error: Some(JsonRpcError {
664 code: -32602,
665 message: format!("Invalid params: {message}"),
666 data: None,
667 }),
668 },
669 },
670 "mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
671 Ok((schedules, tick_ms)) => match runtime.dispatch_schedule_tick(&schedules, tick_ms) {
672 Ok(dispatch) => JsonRpcResponse {
673 jsonrpc: JSONRPC_VERSION.to_string(),
674 id: response_id,
675 result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
676 error: None,
677 },
678 Err(err) => JsonRpcResponse {
679 jsonrpc: JSONRPC_VERSION.to_string(),
680 id: response_id,
681 result: None,
682 error: Some(JsonRpcError {
683 code: -32602,
684 message: format!(
685 "Invalid params: {}",
686 format_schedule_validation_error(err)
687 ),
688 data: None,
689 }),
690 },
691 },
692 Err(message) => JsonRpcResponse {
693 jsonrpc: JSONRPC_VERSION.to_string(),
694 id: response_id,
695 result: None,
696 error: Some(JsonRpcError {
697 code: -32602,
698 message: format!("Invalid params: {message}"),
699 data: None,
700 }),
701 },
702 },
703 "mobkit/routing/resolve" => {
704 match parse_routing_resolve_params(&request.params).and_then(|resolve_request| {
705 runtime
706 .resolve_routing(resolve_request)
707 .map_err(RoutingDeliveryParamsError::Routing)
708 }) {
709 Ok(resolution) => JsonRpcResponse {
710 jsonrpc: JSONRPC_VERSION.to_string(),
711 id: response_id,
712 result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
713 error: None,
714 },
715 Err(err) => JsonRpcResponse {
716 jsonrpc: JSONRPC_VERSION.to_string(),
717 id: response_id,
718 result: None,
719 error: Some(JsonRpcError {
720 code: -32602,
721 message: format!("Invalid params: {}", err.message()),
722 data: None,
723 }),
724 },
725 }
726 }
727 "mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
728 Ok(()) => JsonRpcResponse {
729 jsonrpc: JSONRPC_VERSION.to_string(),
730 id: response_id,
731 result: Some(serde_json::json!({
732 "routes": runtime.list_runtime_routes()
733 })),
734 error: None,
735 },
736 Err(err) => JsonRpcResponse {
737 jsonrpc: JSONRPC_VERSION.to_string(),
738 id: response_id,
739 result: None,
740 error: Some(JsonRpcError {
741 code: -32602,
742 message: format!("Invalid params: {}", err.message()),
743 data: None,
744 }),
745 },
746 },
747 "mobkit/routing/routes/add" => match parse_routing_route_add_params(&request.params)
748 .and_then(|route| {
749 runtime
750 .add_runtime_route(route)
751 .map_err(RoutingDeliveryParamsError::RouteMutation)
752 }) {
753 Ok(route) => JsonRpcResponse {
754 jsonrpc: JSONRPC_VERSION.to_string(),
755 id: response_id,
756 result: Some(serde_json::json!({ "route": route })),
757 error: None,
758 },
759 Err(err) => JsonRpcResponse {
760 jsonrpc: JSONRPC_VERSION.to_string(),
761 id: response_id,
762 result: None,
763 error: Some(JsonRpcError {
764 code: -32602,
765 message: format!("Invalid params: {}", err.message()),
766 data: None,
767 }),
768 },
769 },
770 "mobkit/routing/routes/delete" => match parse_routing_route_delete_params(&request.params)
771 .and_then(|route_key| {
772 runtime
773 .delete_runtime_route(&route_key)
774 .map_err(RoutingDeliveryParamsError::RouteMutation)
775 }) {
776 Ok(route) => JsonRpcResponse {
777 jsonrpc: JSONRPC_VERSION.to_string(),
778 id: response_id,
779 result: Some(serde_json::json!({ "deleted": route })),
780 error: None,
781 },
782 Err(err) => JsonRpcResponse {
783 jsonrpc: JSONRPC_VERSION.to_string(),
784 id: response_id,
785 result: None,
786 error: Some(JsonRpcError {
787 code: -32602,
788 message: format!("Invalid params: {}", err.message()),
789 data: None,
790 }),
791 },
792 },
793 "mobkit/delivery/send" => {
794 match parse_delivery_send_params(&request.params).and_then(|send_request| {
795 runtime
796 .send_delivery(send_request)
797 .map_err(RoutingDeliveryParamsError::Delivery)
798 }) {
799 Ok(record) => JsonRpcResponse {
800 jsonrpc: JSONRPC_VERSION.to_string(),
801 id: response_id,
802 result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
803 error: None,
804 },
805 Err(err) => JsonRpcResponse {
806 jsonrpc: JSONRPC_VERSION.to_string(),
807 id: response_id,
808 result: None,
809 error: Some(JsonRpcError {
810 code: -32602,
811 message: format!("Invalid params: {}", err.message()),
812 data: None,
813 }),
814 },
815 }
816 }
817 "mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
818 Ok(history_request) => JsonRpcResponse {
819 jsonrpc: JSONRPC_VERSION.to_string(),
820 id: response_id,
821 result: Some(
822 serde_json::to_value(runtime.delivery_history(history_request))
823 .unwrap_or(Value::Null),
824 ),
825 error: None,
826 },
827 Err(err) => JsonRpcResponse {
828 jsonrpc: JSONRPC_VERSION.to_string(),
829 id: response_id,
830 result: None,
831 error: Some(JsonRpcError {
832 code: -32602,
833 message: format!("Invalid params: {}", err.message()),
834 data: None,
835 }),
836 },
837 },
838 "mobkit/events/subscribe" => {
839 match parse_subscribe_request(&request.params).and_then(|subscribe_request| {
840 runtime
841 .subscribe_events(subscribe_request)
842 .map_err(SubscribeParamsError::Runtime)
843 }) {
844 Ok(subscribe_result) => JsonRpcResponse {
845 jsonrpc: JSONRPC_VERSION.to_string(),
846 id: response_id,
847 result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
848 error: None,
849 },
850 Err(err) => JsonRpcResponse {
851 jsonrpc: JSONRPC_VERSION.to_string(),
852 id: response_id,
853 result: None,
854 error: Some(JsonRpcError {
855 code: -32602,
856 message: format!("Invalid params: {}", err.message()),
857 data: None,
858 }),
859 },
860 }
861 }
862 "mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
863 Ok(()) => JsonRpcResponse {
864 jsonrpc: JSONRPC_VERSION.to_string(),
865 id: response_id,
866 result: Some(serde_json::json!({
867 "stores": runtime.memory_stores(),
868 })),
869 error: None,
870 },
871 Err(err) => JsonRpcResponse {
872 jsonrpc: JSONRPC_VERSION.to_string(),
873 id: response_id,
874 result: None,
875 error: Some(JsonRpcError {
876 code: -32602,
877 message: format!("Invalid params: {}", err.message()),
878 data: None,
879 }),
880 },
881 },
882 "mobkit/memory/index" => match parse_memory_index_params(&request.params) {
883 Ok(index_request) => match runtime.memory_index(index_request) {
884 Ok(indexed) => JsonRpcResponse {
885 jsonrpc: JSONRPC_VERSION.to_string(),
886 id: response_id,
887 result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
888 error: None,
889 },
890 Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
891 jsonrpc: JSONRPC_VERSION.to_string(),
892 id: response_id,
893 result: None,
894 error: Some(JsonRpcError {
895 code: MEMORY_BACKEND_UNAVAILABLE_CODE,
896 message: format!(
897 "Memory backend unavailable: {}",
898 MemoryParamsError::backend_message(&error)
899 ),
900 data: None,
901 }),
902 },
903 Err(err) => JsonRpcResponse {
904 jsonrpc: JSONRPC_VERSION.to_string(),
905 id: response_id,
906 result: None,
907 error: Some(JsonRpcError {
908 code: -32602,
909 message: format!(
910 "Invalid params: {}",
911 MemoryParamsError::Index(err).message()
912 ),
913 data: None,
914 }),
915 },
916 },
917 Err(err) => JsonRpcResponse {
918 jsonrpc: JSONRPC_VERSION.to_string(),
919 id: response_id,
920 result: None,
921 error: Some(JsonRpcError {
922 code: -32602,
923 message: format!("Invalid params: {}", err.message()),
924 data: None,
925 }),
926 },
927 },
928 "mobkit/memory/query" => match parse_memory_query_params(&request.params) {
929 Ok(query_request) => JsonRpcResponse {
930 jsonrpc: JSONRPC_VERSION.to_string(),
931 id: response_id,
932 result: Some(
933 serde_json::to_value(runtime.memory_query(query_request))
934 .unwrap_or(Value::Null),
935 ),
936 error: None,
937 },
938 Err(err) => JsonRpcResponse {
939 jsonrpc: JSONRPC_VERSION.to_string(),
940 id: response_id,
941 result: None,
942 error: Some(JsonRpcError {
943 code: -32602,
944 message: format!("Invalid params: {}", err.message()),
945 data: None,
946 }),
947 },
948 },
949 "mobkit/session_store/bigquery" => {
950 match parse_bigquery_session_store_params(&request.params)
951 .and_then(run_bigquery_session_store_request)
952 {
953 Ok(result) => JsonRpcResponse {
954 jsonrpc: JSONRPC_VERSION.to_string(),
955 id: response_id,
956 result: Some(result),
957 error: None,
958 },
959 Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
960 jsonrpc: JSONRPC_VERSION.to_string(),
961 id: response_id,
962 result: None,
963 error: Some(JsonRpcError {
964 code: -32602,
965 message: format!("Invalid params: {message}"),
966 data: None,
967 }),
968 },
969 Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
970 jsonrpc: JSONRPC_VERSION.to_string(),
971 id: response_id,
972 result: None,
973 error: Some(JsonRpcError {
974 code: -32011,
975 message: format!(
976 "BigQuery session store request failed: {}",
977 format_bigquery_store_error(&error)
978 ),
979 data: None,
980 }),
981 },
982 }
983 }
984 "mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
985 Ok(gating_request) => JsonRpcResponse {
986 jsonrpc: JSONRPC_VERSION.to_string(),
987 id: response_id,
988 result: Some(
989 serde_json::to_value(runtime.evaluate_gating_action(gating_request))
990 .unwrap_or(Value::Null),
991 ),
992 error: None,
993 },
994 Err(err) => JsonRpcResponse {
995 jsonrpc: JSONRPC_VERSION.to_string(),
996 id: response_id,
997 result: None,
998 error: Some(JsonRpcError {
999 code: -32602,
1000 message: format!("Invalid params: {}", err.message()),
1001 data: None,
1002 }),
1003 },
1004 },
1005 "mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
1006 Ok(()) => JsonRpcResponse {
1007 jsonrpc: JSONRPC_VERSION.to_string(),
1008 id: response_id,
1009 result: Some(serde_json::json!({
1010 "pending": runtime.list_gating_pending(),
1011 })),
1012 error: None,
1013 },
1014 Err(err) => JsonRpcResponse {
1015 jsonrpc: JSONRPC_VERSION.to_string(),
1016 id: response_id,
1017 result: None,
1018 error: Some(JsonRpcError {
1019 code: -32602,
1020 message: format!("Invalid params: {}", err.message()),
1021 data: None,
1022 }),
1023 },
1024 },
1025 "mobkit/gating/decide" => {
1026 match parse_gating_decide_params(&request.params).and_then(|decide_request| {
1027 runtime
1028 .decide_gating_action(decide_request)
1029 .map_err(GatingParamsError::Decision)
1030 }) {
1031 Ok(result) => JsonRpcResponse {
1032 jsonrpc: JSONRPC_VERSION.to_string(),
1033 id: response_id,
1034 result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
1035 error: None,
1036 },
1037 Err(err) => JsonRpcResponse {
1038 jsonrpc: JSONRPC_VERSION.to_string(),
1039 id: response_id,
1040 result: None,
1041 error: Some(JsonRpcError {
1042 code: -32602,
1043 message: format!("Invalid params: {}", err.message()),
1044 data: None,
1045 }),
1046 },
1047 }
1048 }
1049 "mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
1050 Ok(limit) => JsonRpcResponse {
1051 jsonrpc: JSONRPC_VERSION.to_string(),
1052 id: response_id,
1053 result: Some(serde_json::json!({
1054 "entries": runtime.gating_audit_entries(limit),
1055 })),
1056 error: None,
1057 },
1058 Err(err) => JsonRpcResponse {
1059 jsonrpc: JSONRPC_VERSION.to_string(),
1060 id: response_id,
1061 result: None,
1062 error: Some(JsonRpcError {
1063 code: -32602,
1064 message: format!("Invalid params: {}", err.message()),
1065 data: None,
1066 }),
1067 },
1068 },
1069 "mobkit/call_tool" => {
1070 let module_id = request.params.get("module_id").and_then(Value::as_str);
1071 let tool = request.params.get("tool").and_then(Value::as_str);
1072 let arguments = request
1073 .params
1074 .get("arguments")
1075 .cloned()
1076 .unwrap_or(serde_json::json!({}));
1077
1078 match (module_id, tool) {
1079 (Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
1080 let route = route_module_call(
1081 runtime,
1082 &ModuleRouteRequest {
1083 module_id: module_id.to_string(),
1084 method: tool.to_string(),
1085 params: arguments,
1086 },
1087 timeout,
1088 );
1089 match route {
1090 Ok(response) => JsonRpcResponse {
1091 jsonrpc: JSONRPC_VERSION.to_string(),
1092 id: response_id,
1093 result: Some(serde_json::json!({
1094 "module_id": response.module_id,
1095 "tool": response.method,
1096 "result": response.payload
1097 })),
1098 error: None,
1099 },
1100 Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
1101 jsonrpc: JSONRPC_VERSION.to_string(),
1102 id: response_id,
1103 result: None,
1104 error: Some(JsonRpcError {
1105 code: -32601,
1106 message: format!("Module '{mid}' not loaded"),
1107 data: None,
1108 }),
1109 },
1110 Err(err) => JsonRpcResponse {
1111 jsonrpc: JSONRPC_VERSION.to_string(),
1112 id: response_id,
1113 result: None,
1114 error: Some(JsonRpcError {
1115 code: -32000,
1116 message: format!("Tool call failed: {err:?}"),
1117 data: None,
1118 }),
1119 },
1120 }
1121 }
1122 _ => JsonRpcResponse {
1123 jsonrpc: JSONRPC_VERSION.to_string(),
1124 id: response_id,
1125 result: None,
1126 error: Some(JsonRpcError {
1127 code: -32602,
1128 message: "Invalid params: module_id and tool required".to_string(),
1129 data: None,
1130 }),
1131 },
1132 }
1133 }
1134 method if method.contains('/') && !method.starts_with("mobkit/") => {
1135 let module_id = method
1136 .split('/')
1137 .next()
1138 .map(ToString::to_string)
1139 .unwrap_or_default();
1140 let route = route_module_call(
1141 runtime,
1142 &ModuleRouteRequest {
1143 module_id,
1144 method: method.to_string(),
1145 params: request.params,
1146 },
1147 timeout,
1148 );
1149 match route {
1150 Ok(response) => JsonRpcResponse {
1151 jsonrpc: JSONRPC_VERSION.to_string(),
1152 id: response_id,
1153 result: Some(serde_json::json!({
1154 "module_id": response.module_id,
1155 "method": response.method,
1156 "payload": response.payload
1157 })),
1158 error: None,
1159 },
1160 Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
1161 jsonrpc: JSONRPC_VERSION.to_string(),
1162 id: response_id,
1163 result: None,
1164 error: Some(JsonRpcError {
1165 code: -32601,
1166 message: format!("Module '{module_id}' not loaded"),
1167 data: None,
1168 }),
1169 },
1170 Err(err) => JsonRpcResponse {
1171 jsonrpc: JSONRPC_VERSION.to_string(),
1172 id: response_id,
1173 result: None,
1174 error: Some(JsonRpcError {
1175 code: -32000,
1176 message: format!("Module route failed: {err:?}"),
1177 data: None,
1178 }),
1179 },
1180 }
1181 }
1182 _ => JsonRpcResponse {
1183 jsonrpc: JSONRPC_VERSION.to_string(),
1184 id: response_id,
1185 result: None,
1186 error: Some(JsonRpcError {
1187 code: -32601,
1188 message: "Method not found".to_string(),
1189 data: None,
1190 }),
1191 },
1192 };
1193 if is_notification {
1194 String::new()
1195 } else {
1196 serialize_response(&response)
1197 }
1198}
1199
1200pub struct IdentityFirstContext {
1202 pub runtime: std::sync::Arc<crate::identity_first::IdentityRuntime>,
1203 pub roster_provider: std::sync::Arc<dyn crate::identity_first::contracts::RosterProvider>,
1204 pub topology_provider:
1205 Option<std::sync::Arc<dyn crate::identity_first::contracts::TopologyProvider>>,
1206 pub customizer: Option<std::sync::Arc<dyn crate::identity_first::contracts::AgentCustomizer>>,
1207}
1208
1209pub fn handle_unified_rpc_json<'a>(
1210 runtime: &'a UnifiedRuntime,
1211 request_json: &'a str,
1212 timeout: Duration,
1213 http_base_url: Option<&'a str>,
1214 identity_ctx: Option<&'a IdentityFirstContext>,
1215) -> Pin<Box<dyn Future<Output = String> + Send + 'a>> {
1216 Box::pin(handle_unified_rpc_json_inner(
1217 runtime,
1218 request_json,
1219 timeout,
1220 http_base_url,
1221 identity_ctx,
1222 ))
1223}
1224
1225async fn handle_unified_rpc_json_inner(
1226 runtime: &UnifiedRuntime,
1227 request_json: &str,
1228 timeout: Duration,
1229 http_base_url: Option<&str>,
1230 identity_ctx: Option<&IdentityFirstContext>,
1231) -> String {
1232 let raw_request: Value = match serde_json::from_str(request_json) {
1233 Ok(raw_request) => raw_request,
1234 Err(_) => {
1235 return serialize_response(&JsonRpcResponse {
1236 jsonrpc: JSONRPC_VERSION.to_string(),
1237 id: Value::Null,
1238 result: None,
1239 error: Some(JsonRpcError {
1240 code: -32700,
1241 message: "Parse error".to_string(),
1242 data: None,
1243 }),
1244 });
1245 }
1246 };
1247 let response_id = raw_request
1248 .as_object()
1249 .and_then(|object| object.get("id"))
1250 .cloned()
1251 .unwrap_or(Value::Null);
1252 let request: JsonRpcRequest = match serde_json::from_value(raw_request) {
1253 Ok(request) => request,
1254 Err(_) => {
1255 return serialize_response(&JsonRpcResponse {
1256 jsonrpc: JSONRPC_VERSION.to_string(),
1257 id: response_id,
1258 result: None,
1259 error: Some(JsonRpcError {
1260 code: -32600,
1261 message: "Invalid Request".to_string(),
1262 data: None,
1263 }),
1264 });
1265 }
1266 };
1267 let is_notification = request.id.is_none();
1268 let response_id = request.id.clone().unwrap_or(Value::Null);
1269
1270 if request.jsonrpc != "2.0" {
1271 let response = JsonRpcResponse {
1272 jsonrpc: JSONRPC_VERSION.to_string(),
1273 id: response_id,
1274 result: None,
1275 error: Some(JsonRpcError {
1276 code: -32600,
1277 message: "Invalid Request".to_string(),
1278 data: None,
1279 }),
1280 };
1281 return if is_notification {
1282 String::new()
1283 } else {
1284 serialize_response(&response)
1285 };
1286 }
1287
1288 let response = match request.method.as_str() {
1289 "mobkit/status" => {
1290 let mob_state = Some(runtime.mob_handle().status_observation_snapshot());
1291 let is_running = runtime.module_is_running().await;
1292 let loaded = runtime.loaded_modules().await;
1293 let mut result = serde_json::json!({
1294 "contract_version": MOBKIT_CONTRACT_VERSION,
1295 "running": is_running,
1296 "loaded_modules": loaded,
1297 "mob_state": format!("{mob_state:?}"),
1298 });
1299 if let Some(url) = http_base_url {
1300 result["http_base_url"] = Value::String(url.to_string());
1301 }
1302 JsonRpcResponse {
1303 jsonrpc: JSONRPC_VERSION.to_string(),
1304 id: response_id,
1305 result: Some(result),
1306 error: None,
1307 }
1308 }
1309 "mobkit/capabilities" => {
1310 let loaded = runtime.loaded_modules().await;
1311 let mut methods = vec![
1312 "mobkit/init",
1313 "mobkit/status",
1314 "mobkit/capabilities",
1315 "mobkit/reconcile",
1316 "mobkit/spawn_member",
1317 "mobkit/scheduling/evaluate",
1318 "mobkit/scheduling/dispatch",
1319 "mobkit/routing/resolve",
1320 "mobkit/routing/routes/list",
1321 "mobkit/routing/routes/add",
1322 "mobkit/routing/routes/delete",
1323 "mobkit/delivery/send",
1324 "mobkit/delivery/history",
1325 "mobkit/events/subscribe",
1326 "mobkit/query_events",
1327 "mobkit/memory/stores",
1328 "mobkit/memory/index",
1329 "mobkit/memory/query",
1330 "mobkit/session_store/bigquery",
1331 "mobkit/gating/evaluate",
1332 "mobkit/gating/pending",
1333 "mobkit/gating/decide",
1334 "mobkit/gating/audit",
1335 "mobkit/call_tool",
1336 "mobkit/models/catalog",
1337 "mobkit/blob/get",
1338 "mobkit/send_message",
1339 "mobkit/find_members",
1340 "mobkit/ensure_member",
1341 "mobkit/list_members",
1342 "mobkit/get_member",
1343 "mobkit/retire_member",
1344 "mobkit/respawn_member",
1345 "mobkit/reconcile_edges",
1346 "mobkit/rediscover",
1347 "mobkit/mob_events/query",
1348 "mobkit/mob_events/subscribe",
1349 "mobkit/cross_mob/peer_info",
1351 "mobkit/cross_mob/wire_local",
1352 "mobkit/cross_mob/unwire_local",
1353 "mobkit/peer_pubkey",
1354 "mobkit/member_status",
1355 "mobkit/force_cancel_member",
1356 "mobkit/spawn_helper",
1357 "mobkit/fork_helper",
1358 "mobkit/attach_existing_session",
1359 "mobkit/cancel_flow",
1360 "mobkit/flow_status",
1361 "mobkit/list_flows",
1362 "mobkit/list_runs",
1363 "mobkit/run_flow",
1364 "mobkit/collect_completed",
1365 "mobkit/wait_ready",
1366 "mobkit/mob_labels/set",
1367 "mobkit/mob_labels/get",
1368 "mobkit/mob_labels/delete",
1369 "mobkit/run_labels/set",
1370 "mobkit/run_labels/get",
1371 "mobkit/run_labels/delete",
1372 ];
1373 methods.extend_from_slice(MOBPACK_AUTHORING_METHODS);
1374 if identity_ctx.is_some() {
1375 methods.extend_from_slice(&[
1376 "mobkit/send",
1377 "mobkit/interact",
1378 "mobkit/dispatch",
1379 "mobkit/subscribe",
1380 "mobkit/status_identity",
1381 "mobkit/respawn",
1382 "mobkit/retire",
1383 "mobkit/reset",
1384 "mobkit/delete_identity",
1385 "mobkit/inspect_identity",
1386 "mobkit/reconcile_identity",
1387 ]);
1388 }
1389 if runtime.has_contact_directory() {
1391 methods.push("mobkit/cross_mob/directory");
1392 }
1393 if runtime.has_peer_mob_handles().await && runtime.has_inproc_contacts() {
1397 methods.extend_from_slice(&[
1398 "mobkit/cross_mob/wire",
1399 "mobkit/cross_mob/unwire",
1400 "mobkit/cross_mob/send",
1401 ]);
1402 }
1403 JsonRpcResponse {
1404 jsonrpc: JSONRPC_VERSION.to_string(),
1405 id: response_id,
1406 result: Some(serde_json::json!({
1407 "contract_version": MOBKIT_CONTRACT_VERSION,
1408 "runtime_type": "unified",
1409 "methods": methods,
1410 "loaded_modules": loaded,
1411 "runtime_capabilities": {
1412 "can_spawn_members": true,
1413 "can_send_messages": true,
1414 "can_wire_members": true,
1415 "can_retire_members": true,
1416 "available_spawn_modes": ["module", "profile"],
1417 },
1418 "authoring_capabilities": mobpack_authoring_capabilities(),
1419 })),
1420 error: None,
1421 }
1422 }
1423 "mobkit/reconcile" => {
1424 let modules = match params::required_string_array(&request.params, "modules") {
1425 Ok(m) => m,
1426 Err(reason) => {
1427 return serialize_response(&JsonRpcResponse {
1428 jsonrpc: JSONRPC_VERSION.to_string(),
1429 id: response_id,
1430 result: None,
1431 error: Some(JsonRpcError {
1432 code: -32602,
1433 message: format!("Invalid params: {reason}"),
1434 data: None,
1435 }),
1436 });
1437 }
1438 };
1439
1440 match runtime.reconcile_modules(modules.clone(), timeout).await {
1441 Ok(added) => JsonRpcResponse {
1442 jsonrpc: JSONRPC_VERSION.to_string(),
1443 id: response_id,
1444 result: Some(serde_json::json!({
1445 "accepted": true,
1446 "reconciled_modules": modules,
1447 "added": added
1448 })),
1449 error: None,
1450 },
1451 Err(err) => JsonRpcResponse {
1452 jsonrpc: JSONRPC_VERSION.to_string(),
1453 id: response_id,
1454 result: None,
1455 error: Some(JsonRpcError {
1456 code: -32602,
1457 message: format!("Invalid params: {err:?}"),
1458 data: None,
1459 }),
1460 },
1461 }
1462 }
1463 "mobkit/spawn_member" => {
1464 let module_id = request.params.get("module_id").and_then(Value::as_str);
1466 let profile = request.params.get("profile").and_then(Value::as_str);
1467 let meerkat_id = request.params.get("meerkat_id").and_then(Value::as_str);
1468
1469 if let Some(module_id) = module_id {
1470 if module_id.is_empty() {
1472 JsonRpcResponse {
1473 jsonrpc: JSONRPC_VERSION.to_string(),
1474 id: response_id,
1475 result: None,
1476 error: Some(JsonRpcError {
1477 code: -32602,
1478 message: "Invalid params: module_id required".to_string(),
1479 data: None,
1480 }),
1481 }
1482 } else {
1483 match runtime.spawn_member(module_id, timeout).await {
1484 Ok(()) => JsonRpcResponse {
1485 jsonrpc: JSONRPC_VERSION.to_string(),
1486 id: response_id,
1487 result: Some(serde_json::json!({
1488 "accepted": true,
1489 "module_id": module_id
1490 })),
1491 error: None,
1492 },
1493 Err(err) => JsonRpcResponse {
1494 jsonrpc: JSONRPC_VERSION.to_string(),
1495 id: response_id,
1496 result: None,
1497 error: Some(JsonRpcError {
1498 code: -32602,
1499 message: format!("Invalid params: {err:?}"),
1500 data: None,
1501 }),
1502 },
1503 }
1504 }
1505 } else if let (Some(profile), Some(meerkat_id)) = (profile, meerkat_id) {
1506 let spec = meerkat_mob::SpawnMemberSpec::from_wire(
1508 profile.to_string(),
1509 meerkat_id.to_string(),
1510 request
1511 .params
1512 .get("initial_message")
1513 .and_then(Value::as_str)
1514 .map(|s| meerkat_core::ContentInput::from(s.to_string())),
1515 None,
1516 None,
1517 );
1518 match Box::pin(runtime.spawn(spec)).await {
1519 Ok(_member_ref) => JsonRpcResponse {
1520 jsonrpc: JSONRPC_VERSION.to_string(),
1521 id: response_id,
1522 result: Some(serde_json::json!({
1523 "accepted": true,
1524 "meerkat_id": meerkat_id
1525 })),
1526 error: None,
1527 },
1528 Err(err) => JsonRpcResponse {
1529 jsonrpc: JSONRPC_VERSION.to_string(),
1530 id: response_id,
1531 result: None,
1532 error: Some(JsonRpcError {
1533 code: -32602,
1534 message: format!("Invalid params: {err}"),
1535 data: None,
1536 }),
1537 },
1538 }
1539 } else {
1540 JsonRpcResponse {
1541 jsonrpc: JSONRPC_VERSION.to_string(),
1542 id: response_id,
1543 result: None,
1544 error: Some(JsonRpcError {
1545 code: -32602,
1546 message: "Invalid params: module_id or (profile + meerkat_id) required"
1547 .to_string(),
1548 data: None,
1549 }),
1550 }
1551 }
1552 }
1553 "mobkit/scheduling/evaluate" => match parse_scheduling_params(&request.params) {
1554 Ok((schedules, tick_ms)) => {
1555 match runtime.evaluate_schedule_tick(&schedules, tick_ms).await {
1556 Ok(evaluation) => JsonRpcResponse {
1557 jsonrpc: JSONRPC_VERSION.to_string(),
1558 id: response_id,
1559 result: Some(serde_json::to_value(evaluation).unwrap_or(Value::Null)),
1560 error: None,
1561 },
1562 Err(err) => JsonRpcResponse {
1563 jsonrpc: JSONRPC_VERSION.to_string(),
1564 id: response_id,
1565 result: None,
1566 error: Some(JsonRpcError {
1567 code: -32602,
1568 message: format!(
1569 "Invalid params: {}",
1570 format_schedule_validation_error(err)
1571 ),
1572 data: None,
1573 }),
1574 },
1575 }
1576 }
1577 Err(message) => JsonRpcResponse {
1578 jsonrpc: JSONRPC_VERSION.to_string(),
1579 id: response_id,
1580 result: None,
1581 error: Some(JsonRpcError {
1582 code: -32602,
1583 message: format!("Invalid params: {message}"),
1584 data: None,
1585 }),
1586 },
1587 },
1588 "mobkit/scheduling/dispatch" => match parse_scheduling_params(&request.params) {
1589 Ok((schedules, tick_ms)) => {
1590 match runtime.dispatch_schedule_tick(&schedules, tick_ms).await {
1591 Ok(dispatch) => JsonRpcResponse {
1592 jsonrpc: JSONRPC_VERSION.to_string(),
1593 id: response_id,
1594 result: Some(serde_json::to_value(dispatch).unwrap_or(Value::Null)),
1595 error: None,
1596 },
1597 Err(err) => JsonRpcResponse {
1598 jsonrpc: JSONRPC_VERSION.to_string(),
1599 id: response_id,
1600 result: None,
1601 error: Some(JsonRpcError {
1602 code: -32602,
1603 message: format!("Invalid params: {err}"),
1604 data: None,
1605 }),
1606 },
1607 }
1608 }
1609 Err(message) => JsonRpcResponse {
1610 jsonrpc: JSONRPC_VERSION.to_string(),
1611 id: response_id,
1612 result: None,
1613 error: Some(JsonRpcError {
1614 code: -32602,
1615 message: format!("Invalid params: {message}"),
1616 data: None,
1617 }),
1618 },
1619 },
1620 "mobkit/routing/resolve" => {
1621 let resolve_result = match parse_routing_resolve_params(&request.params) {
1622 Ok(resolve_request) => runtime
1623 .resolve_routing(resolve_request)
1624 .await
1625 .map_err(RoutingDeliveryParamsError::Routing),
1626 Err(e) => Err(e),
1627 };
1628 match resolve_result {
1629 Ok(resolution) => JsonRpcResponse {
1630 jsonrpc: JSONRPC_VERSION.to_string(),
1631 id: response_id,
1632 result: Some(serde_json::to_value(resolution).unwrap_or(Value::Null)),
1633 error: None,
1634 },
1635 Err(err) => JsonRpcResponse {
1636 jsonrpc: JSONRPC_VERSION.to_string(),
1637 id: response_id,
1638 result: None,
1639 error: Some(JsonRpcError {
1640 code: -32602,
1641 message: format!("Invalid params: {}", err.message()),
1642 data: None,
1643 }),
1644 },
1645 }
1646 }
1647 "mobkit/routing/routes/list" => match parse_routing_routes_list_params(&request.params) {
1648 Ok(()) => {
1649 let routes = runtime.list_runtime_routes().await;
1650 JsonRpcResponse {
1651 jsonrpc: JSONRPC_VERSION.to_string(),
1652 id: response_id,
1653 result: Some(serde_json::json!({
1654 "routes": routes
1655 })),
1656 error: None,
1657 }
1658 }
1659 Err(err) => JsonRpcResponse {
1660 jsonrpc: JSONRPC_VERSION.to_string(),
1661 id: response_id,
1662 result: None,
1663 error: Some(JsonRpcError {
1664 code: -32602,
1665 message: format!("Invalid params: {}", err.message()),
1666 data: None,
1667 }),
1668 },
1669 },
1670 "mobkit/routing/routes/add" => {
1671 let add_result = match parse_routing_route_add_params(&request.params) {
1672 Ok(route) => runtime
1673 .add_runtime_route(route)
1674 .await
1675 .map_err(RoutingDeliveryParamsError::RouteMutation),
1676 Err(e) => Err(e),
1677 };
1678 match add_result {
1679 Ok(route) => JsonRpcResponse {
1680 jsonrpc: JSONRPC_VERSION.to_string(),
1681 id: response_id,
1682 result: Some(serde_json::json!({ "route": route })),
1683 error: None,
1684 },
1685 Err(err) => JsonRpcResponse {
1686 jsonrpc: JSONRPC_VERSION.to_string(),
1687 id: response_id,
1688 result: None,
1689 error: Some(JsonRpcError {
1690 code: -32602,
1691 message: format!("Invalid params: {}", err.message()),
1692 data: None,
1693 }),
1694 },
1695 }
1696 }
1697 "mobkit/routing/routes/delete" => {
1698 let delete_result = match parse_routing_route_delete_params(&request.params) {
1699 Ok(route_key) => runtime
1700 .delete_runtime_route(&route_key)
1701 .await
1702 .map_err(RoutingDeliveryParamsError::RouteMutation),
1703 Err(e) => Err(e),
1704 };
1705 match delete_result {
1706 Ok(route) => JsonRpcResponse {
1707 jsonrpc: JSONRPC_VERSION.to_string(),
1708 id: response_id,
1709 result: Some(serde_json::json!({ "deleted": route })),
1710 error: None,
1711 },
1712 Err(err) => JsonRpcResponse {
1713 jsonrpc: JSONRPC_VERSION.to_string(),
1714 id: response_id,
1715 result: None,
1716 error: Some(JsonRpcError {
1717 code: -32602,
1718 message: format!("Invalid params: {}", err.message()),
1719 data: None,
1720 }),
1721 },
1722 }
1723 }
1724 "mobkit/delivery/send" => {
1725 let send_result = match parse_delivery_send_params(&request.params) {
1726 Ok(send_request) => runtime
1727 .send_delivery(send_request)
1728 .await
1729 .map_err(RoutingDeliveryParamsError::Delivery),
1730 Err(e) => Err(e),
1731 };
1732 match send_result {
1733 Ok(record) => JsonRpcResponse {
1734 jsonrpc: JSONRPC_VERSION.to_string(),
1735 id: response_id,
1736 result: Some(serde_json::to_value(record).unwrap_or(Value::Null)),
1737 error: None,
1738 },
1739 Err(err) => JsonRpcResponse {
1740 jsonrpc: JSONRPC_VERSION.to_string(),
1741 id: response_id,
1742 result: None,
1743 error: Some(JsonRpcError {
1744 code: -32602,
1745 message: format!("Invalid params: {}", err.message()),
1746 data: None,
1747 }),
1748 },
1749 }
1750 }
1751 "mobkit/delivery/history" => match parse_delivery_history_params(&request.params) {
1752 Ok(history_request) => {
1753 let history = runtime.delivery_history(history_request).await;
1754 JsonRpcResponse {
1755 jsonrpc: JSONRPC_VERSION.to_string(),
1756 id: response_id,
1757 result: Some(serde_json::to_value(history).unwrap_or(Value::Null)),
1758 error: None,
1759 }
1760 }
1761 Err(err) => JsonRpcResponse {
1762 jsonrpc: JSONRPC_VERSION.to_string(),
1763 id: response_id,
1764 result: None,
1765 error: Some(JsonRpcError {
1766 code: -32602,
1767 message: format!("Invalid params: {}", err.message()),
1768 data: None,
1769 }),
1770 },
1771 },
1772 "mobkit/events/subscribe" => match parse_subscribe_request(&request.params) {
1773 Ok(subscribe_request) => match runtime.subscribe_events(subscribe_request).await {
1774 Ok(subscribe_result) => JsonRpcResponse {
1775 jsonrpc: JSONRPC_VERSION.to_string(),
1776 id: response_id,
1777 result: Some(serde_json::to_value(subscribe_result).unwrap_or(Value::Null)),
1778 error: None,
1779 },
1780 Err(err) => JsonRpcResponse {
1781 jsonrpc: JSONRPC_VERSION.to_string(),
1782 id: response_id,
1783 result: None,
1784 error: Some(JsonRpcError {
1785 code: -32602,
1786 message: format!("Invalid params: {err}"),
1787 data: None,
1788 }),
1789 },
1790 },
1791 Err(err) => JsonRpcResponse {
1792 jsonrpc: JSONRPC_VERSION.to_string(),
1793 id: response_id,
1794 result: None,
1795 error: Some(JsonRpcError {
1796 code: -32602,
1797 message: format!("Invalid params: {}", err.message()),
1798 data: None,
1799 }),
1800 },
1801 },
1802 "mobkit/query_events" => {
1803 let query: EventQuery = if request.params.is_null() {
1804 EventQuery::default()
1805 } else {
1806 match serde_json::from_value(request.params.clone()) {
1807 Ok(query) => query,
1808 Err(err) => {
1809 return serde_json::to_string(&JsonRpcResponse {
1810 jsonrpc: JSONRPC_VERSION.to_string(),
1811 id: response_id,
1812 result: None,
1813 error: Some(JsonRpcError {
1814 code: -32602,
1815 message: format!("Invalid params: invalid query params: {err}"),
1816 data: None,
1817 }),
1818 })
1819 .unwrap_or_else(|_| r#"{"error":"serialization failed"}"#.to_string());
1820 }
1821 }
1822 };
1823 match runtime.event_log_store() {
1824 Some(store) => match store.query(query).await {
1825 Ok(events) => JsonRpcResponse {
1826 jsonrpc: JSONRPC_VERSION.to_string(),
1827 id: response_id,
1828 result: Some(serde_json::to_value(events).unwrap_or(Value::Null)),
1829 error: None,
1830 },
1831 Err(err) => JsonRpcResponse {
1832 jsonrpc: JSONRPC_VERSION.to_string(),
1833 id: response_id,
1834 result: None,
1835 error: Some(JsonRpcError {
1836 code: -32603,
1837 message: format!("query_events failed: {err}"),
1838 data: None,
1839 }),
1840 },
1841 },
1842 None => JsonRpcResponse {
1843 jsonrpc: JSONRPC_VERSION.to_string(),
1844 id: response_id,
1845 result: Some(serde_json::json!({
1846 "status": "no_event_log_configured",
1847 "events": [],
1848 })),
1849 error: None,
1850 },
1851 }
1852 }
1853 "mobkit/memory/stores" => match parse_memory_stores_params(&request.params) {
1854 Ok(()) => {
1855 let stores = runtime.memory_stores().await;
1856 JsonRpcResponse {
1857 jsonrpc: JSONRPC_VERSION.to_string(),
1858 id: response_id,
1859 result: Some(serde_json::json!({
1860 "stores": stores,
1861 })),
1862 error: None,
1863 }
1864 }
1865 Err(err) => JsonRpcResponse {
1866 jsonrpc: JSONRPC_VERSION.to_string(),
1867 id: response_id,
1868 result: None,
1869 error: Some(JsonRpcError {
1870 code: -32602,
1871 message: format!("Invalid params: {}", err.message()),
1872 data: None,
1873 }),
1874 },
1875 },
1876 "mobkit/memory/index" => match parse_memory_index_params(&request.params) {
1877 Ok(index_request) => match runtime.memory_index(index_request).await {
1878 Ok(indexed) => JsonRpcResponse {
1879 jsonrpc: JSONRPC_VERSION.to_string(),
1880 id: response_id,
1881 result: Some(serde_json::to_value(indexed).unwrap_or(Value::Null)),
1882 error: None,
1883 },
1884 Err(MemoryIndexError::BackendPersistFailed(error)) => JsonRpcResponse {
1885 jsonrpc: JSONRPC_VERSION.to_string(),
1886 id: response_id,
1887 result: None,
1888 error: Some(JsonRpcError {
1889 code: MEMORY_BACKEND_UNAVAILABLE_CODE,
1890 message: format!(
1891 "Memory backend unavailable: {}",
1892 MemoryParamsError::backend_message(&error)
1893 ),
1894 data: None,
1895 }),
1896 },
1897 Err(err) => JsonRpcResponse {
1898 jsonrpc: JSONRPC_VERSION.to_string(),
1899 id: response_id,
1900 result: None,
1901 error: Some(JsonRpcError {
1902 code: -32602,
1903 message: format!(
1904 "Invalid params: {}",
1905 MemoryParamsError::Index(err).message()
1906 ),
1907 data: None,
1908 }),
1909 },
1910 },
1911 Err(err) => JsonRpcResponse {
1912 jsonrpc: JSONRPC_VERSION.to_string(),
1913 id: response_id,
1914 result: None,
1915 error: Some(JsonRpcError {
1916 code: -32602,
1917 message: format!("Invalid params: {}", err.message()),
1918 data: None,
1919 }),
1920 },
1921 },
1922 "mobkit/memory/query" => match parse_memory_query_params(&request.params) {
1923 Ok(query_request) => {
1924 let query_result = runtime.memory_query(query_request).await;
1925 JsonRpcResponse {
1926 jsonrpc: JSONRPC_VERSION.to_string(),
1927 id: response_id,
1928 result: Some(serde_json::to_value(query_result).unwrap_or(Value::Null)),
1929 error: None,
1930 }
1931 }
1932 Err(err) => JsonRpcResponse {
1933 jsonrpc: JSONRPC_VERSION.to_string(),
1934 id: response_id,
1935 result: None,
1936 error: Some(JsonRpcError {
1937 code: -32602,
1938 message: format!("Invalid params: {}", err.message()),
1939 data: None,
1940 }),
1941 },
1942 },
1943 "mobkit/session_store/bigquery" => {
1944 match parse_bigquery_session_store_params(&request.params)
1945 .and_then(run_bigquery_session_store_request)
1946 {
1947 Ok(result) => JsonRpcResponse {
1948 jsonrpc: JSONRPC_VERSION.to_string(),
1949 id: response_id,
1950 result: Some(result),
1951 error: None,
1952 },
1953 Err(BigQuerySessionStoreRpcError::Params(message)) => JsonRpcResponse {
1954 jsonrpc: JSONRPC_VERSION.to_string(),
1955 id: response_id,
1956 result: None,
1957 error: Some(JsonRpcError {
1958 code: -32602,
1959 message: format!("Invalid params: {message}"),
1960 data: None,
1961 }),
1962 },
1963 Err(BigQuerySessionStoreRpcError::Store(error)) => JsonRpcResponse {
1964 jsonrpc: JSONRPC_VERSION.to_string(),
1965 id: response_id,
1966 result: None,
1967 error: Some(JsonRpcError {
1968 code: -32011,
1969 message: format!(
1970 "BigQuery session store request failed: {}",
1971 format_bigquery_store_error(&error)
1972 ),
1973 data: None,
1974 }),
1975 },
1976 }
1977 }
1978 "mobkit/gating/evaluate" => match parse_gating_evaluate_params(&request.params) {
1979 Ok(gating_request) => {
1980 let gating_result = runtime.evaluate_gating_action(gating_request).await;
1981 JsonRpcResponse {
1982 jsonrpc: JSONRPC_VERSION.to_string(),
1983 id: response_id,
1984 result: Some(serde_json::to_value(gating_result).unwrap_or(Value::Null)),
1985 error: None,
1986 }
1987 }
1988 Err(err) => JsonRpcResponse {
1989 jsonrpc: JSONRPC_VERSION.to_string(),
1990 id: response_id,
1991 result: None,
1992 error: Some(JsonRpcError {
1993 code: -32602,
1994 message: format!("Invalid params: {}", err.message()),
1995 data: None,
1996 }),
1997 },
1998 },
1999 "mobkit/gating/pending" => match parse_gating_pending_params(&request.params) {
2000 Ok(()) => {
2001 let pending = runtime.list_gating_pending().await;
2002 JsonRpcResponse {
2003 jsonrpc: JSONRPC_VERSION.to_string(),
2004 id: response_id,
2005 result: Some(serde_json::json!({
2006 "pending": pending,
2007 })),
2008 error: None,
2009 }
2010 }
2011 Err(err) => JsonRpcResponse {
2012 jsonrpc: JSONRPC_VERSION.to_string(),
2013 id: response_id,
2014 result: None,
2015 error: Some(JsonRpcError {
2016 code: -32602,
2017 message: format!("Invalid params: {}", err.message()),
2018 data: None,
2019 }),
2020 },
2021 },
2022 "mobkit/gating/decide" => {
2023 let decide_result = match parse_gating_decide_params(&request.params) {
2024 Ok(decide_request) => runtime
2025 .decide_gating_action(decide_request)
2026 .await
2027 .map_err(GatingParamsError::Decision),
2028 Err(e) => Err(e),
2029 };
2030 match decide_result {
2031 Ok(result) => JsonRpcResponse {
2032 jsonrpc: JSONRPC_VERSION.to_string(),
2033 id: response_id,
2034 result: Some(serde_json::to_value(result).unwrap_or(Value::Null)),
2035 error: None,
2036 },
2037 Err(err) => JsonRpcResponse {
2038 jsonrpc: JSONRPC_VERSION.to_string(),
2039 id: response_id,
2040 result: None,
2041 error: Some(JsonRpcError {
2042 code: -32602,
2043 message: format!("Invalid params: {}", err.message()),
2044 data: None,
2045 }),
2046 },
2047 }
2048 }
2049 "mobkit/gating/audit" => match parse_gating_audit_params(&request.params) {
2050 Ok(limit) => {
2051 let entries = runtime.gating_audit_entries(limit).await;
2052 JsonRpcResponse {
2053 jsonrpc: JSONRPC_VERSION.to_string(),
2054 id: response_id,
2055 result: Some(serde_json::json!({
2056 "entries": entries,
2057 })),
2058 error: None,
2059 }
2060 }
2061 Err(err) => JsonRpcResponse {
2062 jsonrpc: JSONRPC_VERSION.to_string(),
2063 id: response_id,
2064 result: None,
2065 error: Some(JsonRpcError {
2066 code: -32602,
2067 message: format!("Invalid params: {}", err.message()),
2068 data: None,
2069 }),
2070 },
2071 },
2072 "mobkit/call_tool" => {
2073 let module_id = request.params.get("module_id").and_then(Value::as_str);
2074 let tool = request.params.get("tool").and_then(Value::as_str);
2075 let arguments = request
2076 .params
2077 .get("arguments")
2078 .cloned()
2079 .unwrap_or(serde_json::json!({}));
2080
2081 match (module_id, tool) {
2082 (Some(module_id), Some(tool)) if !module_id.is_empty() && !tool.is_empty() => {
2083 let route = runtime
2084 .route_module_call(
2085 &ModuleRouteRequest {
2086 module_id: module_id.to_string(),
2087 method: tool.to_string(),
2088 params: arguments,
2089 },
2090 timeout,
2091 )
2092 .await;
2093 match route {
2094 Ok(response) => JsonRpcResponse {
2095 jsonrpc: JSONRPC_VERSION.to_string(),
2096 id: response_id,
2097 result: Some(serde_json::json!({
2098 "module_id": response.module_id,
2099 "tool": response.method,
2100 "result": response.payload
2101 })),
2102 error: None,
2103 },
2104 Err(ModuleRouteError::UnloadedModule(mid)) => JsonRpcResponse {
2105 jsonrpc: JSONRPC_VERSION.to_string(),
2106 id: response_id,
2107 result: None,
2108 error: Some(JsonRpcError {
2109 code: -32601,
2110 message: format!("Module '{mid}' not loaded"),
2111 data: None,
2112 }),
2113 },
2114 Err(err) => JsonRpcResponse {
2115 jsonrpc: JSONRPC_VERSION.to_string(),
2116 id: response_id,
2117 result: None,
2118 error: Some(JsonRpcError {
2119 code: -32000,
2120 message: format!("Tool call failed: {err:?}"),
2121 data: None,
2122 }),
2123 },
2124 }
2125 }
2126 _ => JsonRpcResponse {
2127 jsonrpc: JSONRPC_VERSION.to_string(),
2128 id: response_id,
2129 result: None,
2130 error: Some(JsonRpcError {
2131 code: -32602,
2132 message: "Invalid params: module_id and tool required".to_string(),
2133 data: None,
2134 }),
2135 },
2136 }
2137 }
2138 "mobkit/models/catalog" => JsonRpcResponse {
2139 jsonrpc: JSONRPC_VERSION.to_string(),
2140 id: response_id,
2141 result: Some(build_models_catalog_result()),
2142 error: None,
2143 },
2144 method if MOBPACK_AUTHORING_METHODS.contains(&method) => {
2145 handle_unified_mobpack_authoring_rpc(runtime, method, &request.params, response_id)
2146 .await
2147 }
2148 "mobkit/blob/get" => {
2149 mob_methods::handle_blob_get(runtime, response_id, &request.params).await
2150 }
2151 "mobkit/send_message" => {
2152 Box::pin(mob_methods::handle_send_message(
2156 runtime,
2157 identity_ctx.map(|ctx| &ctx.runtime),
2158 response_id,
2159 &request.params,
2160 ))
2161 .await
2162 }
2163 "mobkit/find_members" => {
2164 mob_methods::handle_find_members(runtime, response_id, &request.params).await
2165 }
2166 "mobkit/ensure_member" => {
2167 Box::pin(mob_methods::handle_ensure_member(
2168 runtime,
2169 response_id,
2170 &request.params,
2171 ))
2172 .await
2173 }
2174 "mobkit/list_members" => mob_methods::handle_list_members(runtime, response_id).await,
2175 "mobkit/get_member" => {
2176 mob_methods::handle_get_member(runtime, response_id, &request.params).await
2177 }
2178 "mobkit/retire_member" => {
2179 mob_methods::handle_retire_member(runtime, response_id, &request.params).await
2180 }
2181 "mobkit/respawn_member" => {
2182 Box::pin(mob_methods::handle_respawn_member(
2183 runtime,
2184 response_id,
2185 &request.params,
2186 ))
2187 .await
2188 }
2189 "mobkit/reconcile_edges" => mob_methods::handle_reconcile_edges(runtime, response_id).await,
2190 "mobkit/rediscover" => mob_methods::handle_rediscover(runtime, response_id).await,
2191 "mobkit/mob_events/query" => {
2192 mob_methods::handle_mob_events_query(runtime, response_id, request.params).await
2193 }
2194 "mobkit/mob_events/subscribe" => {
2195 mob_methods::handle_mob_events_subscribe(runtime, response_id, request.params).await
2196 }
2197 "mobkit/cross_mob/wire" => {
2198 Box::pin(mob_methods::handle_cross_mob_wire(
2199 runtime,
2200 response_id,
2201 &request.params,
2202 ))
2203 .await
2204 }
2205 "mobkit/cross_mob/unwire" => {
2206 Box::pin(mob_methods::handle_cross_mob_unwire(
2207 runtime,
2208 response_id,
2209 &request.params,
2210 ))
2211 .await
2212 }
2213 "mobkit/cross_mob/send" => {
2214 mob_methods::handle_cross_mob_send(runtime, response_id, &request.params).await
2215 }
2216 "mobkit/cross_mob/directory" => {
2217 mob_methods::handle_cross_mob_directory(runtime, response_id).await
2218 }
2219 "mobkit/cross_mob/peer_info" => {
2220 mob_methods::handle_cross_mob_peer_info(runtime, response_id, &request.params).await
2221 }
2222 "mobkit/cross_mob/wire_local" => {
2223 mob_methods::handle_cross_mob_wire_local(runtime, response_id, &request.params).await
2224 }
2225 "mobkit/cross_mob/unwire_local" => {
2226 mob_methods::handle_cross_mob_unwire_local(runtime, response_id, &request.params).await
2227 }
2228 "mobkit/peer_pubkey" => mob_methods::handle_peer_pubkey(runtime, response_id).await,
2229 "mobkit/member_status" => {
2230 mob_methods::handle_member_status(runtime, response_id, &request.params).await
2231 }
2232 "mobkit/force_cancel_member" => {
2233 mob_methods::handle_force_cancel_member(runtime, response_id, &request.params).await
2234 }
2235 "mobkit/spawn_helper" => {
2236 Box::pin(mob_methods::handle_spawn_helper(
2237 runtime,
2238 response_id,
2239 &request.params,
2240 ))
2241 .await
2242 }
2243 "mobkit/fork_helper" => {
2244 Box::pin(mob_methods::handle_fork_helper(
2245 runtime,
2246 response_id,
2247 &request.params,
2248 ))
2249 .await
2250 }
2251 "mobkit/attach_existing_session" => {
2252 Box::pin(mob_methods::handle_attach_existing_session(
2253 runtime,
2254 response_id,
2255 &request.params,
2256 ))
2257 .await
2258 }
2259 "mobkit/cancel_flow" => {
2260 mob_methods::handle_cancel_flow(runtime, response_id, &request.params).await
2261 }
2262 "mobkit/flow_status" => {
2263 mob_methods::handle_flow_status(runtime, response_id, &request.params).await
2264 }
2265 "mobkit/list_flows" => mob_methods::handle_list_flows(runtime, response_id).await,
2266 "mobkit/list_runs" => {
2267 mob_methods::handle_list_runs(runtime, response_id, &request.params).await
2268 }
2269 "mobkit/run_flow" => {
2270 Box::pin(mob_methods::handle_run_flow(
2271 runtime,
2272 response_id,
2273 &request.params,
2274 ))
2275 .await
2276 }
2277 "mobkit/collect_completed" => {
2278 mob_methods::handle_collect_completed(runtime, response_id).await
2279 }
2280 "mobkit/wait_ready" => {
2281 mob_methods::handle_wait_ready(runtime, response_id, &request.params).await
2282 }
2283 "mobkit/mob_labels/set" => {
2284 mob_methods::handle_mob_labels_set(runtime, response_id, &request.params).await
2285 }
2286 "mobkit/mob_labels/get" => mob_methods::handle_mob_labels_get(runtime, response_id).await,
2287 "mobkit/mob_labels/delete" => {
2288 mob_methods::handle_mob_labels_delete(runtime, response_id).await
2289 }
2290 "mobkit/run_labels/set" => {
2291 mob_methods::handle_run_labels_set(runtime, response_id, &request.params).await
2292 }
2293 "mobkit/run_labels/get" => {
2294 mob_methods::handle_run_labels_get(runtime, response_id, &request.params).await
2295 }
2296 "mobkit/run_labels/delete" => {
2297 mob_methods::handle_run_labels_delete(runtime, response_id, &request.params).await
2298 }
2299 "mobkit/send" => {
2301 let identity_rt = match identity_ctx {
2302 Some(ctx) => &*ctx.runtime,
2303 None => return maybe_identity_not_configured(is_notification, response_id),
2304 };
2305 let identity_str = request
2306 .params
2307 .get("identity")
2308 .and_then(|v| v.as_str())
2309 .unwrap_or("");
2310 let target =
2311 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2312 {
2313 Ok(target) => target,
2314 Err(e) => {
2315 return maybe_error_response(
2316 is_notification,
2317 response_id,
2318 -32602,
2319 format!("invalid identity: {e}"),
2320 );
2321 }
2322 };
2323 let identity = target.identity.clone();
2324 if let Some(response) =
2325 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2326 {
2327 return if is_notification {
2328 String::new()
2329 } else {
2330 serialize_response(&response)
2331 };
2332 }
2333 let content_val = request
2334 .params
2335 .get("content")
2336 .cloned()
2337 .unwrap_or(Value::Null);
2338 let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
2339 Ok(content) => content,
2340 Err(err) => {
2341 return maybe_error_response(
2342 is_notification,
2343 response_id,
2344 -32602,
2345 format!("invalid content: {err}"),
2346 );
2347 }
2348 };
2349 match identity_rt.send(&identity, &content).await {
2350 Ok(token) => JsonRpcResponse {
2351 jsonrpc: JSONRPC_VERSION.to_string(),
2352 id: response_id,
2353 result: Some(serde_json::json!({ "fencing_token": token.get() })),
2354 error: None,
2355 },
2356 Err(e) => identity_error_response(response_id, &e),
2357 }
2358 }
2359 "mobkit/interact" => {
2360 let identity_rt = match identity_ctx {
2361 Some(ctx) => &*ctx.runtime,
2362 None => return maybe_identity_not_configured(is_notification, response_id),
2363 };
2364 let identity_str = request
2365 .params
2366 .get("identity")
2367 .and_then(|v| v.as_str())
2368 .unwrap_or("");
2369 let target =
2370 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2371 {
2372 Ok(target) => target,
2373 Err(e) => {
2374 return maybe_error_response(
2375 is_notification,
2376 response_id,
2377 -32602,
2378 format!("invalid identity: {e}"),
2379 );
2380 }
2381 };
2382 let identity = target.identity.clone();
2383 if let Some(response) =
2384 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2385 {
2386 return if is_notification {
2387 String::new()
2388 } else {
2389 serialize_response(&response)
2390 };
2391 }
2392 let content_val = request
2393 .params
2394 .get("content")
2395 .cloned()
2396 .unwrap_or(Value::Null);
2397 let content =
2398 match serde_json::from_value::<meerkat_core::ContentInput>(content_val.clone()) {
2399 Ok(content) => content,
2400 Err(err) => {
2401 return maybe_error_response(
2402 is_notification,
2403 response_id,
2404 -32602,
2405 format!("invalid content: {err}"),
2406 );
2407 }
2408 };
2409 let origin = request
2410 .params
2411 .get("origin")
2412 .and_then(|v| v.as_str())
2413 .unwrap_or("console");
2414 let interaction_id = request
2415 .params
2416 .get("interaction_id")
2417 .and_then(|v| v.as_str())
2418 .map(ToString::to_string)
2419 .unwrap_or_else(|| meerkat_core::types::SessionId::new().to_string());
2420 let runtime_member_id = identity_rt
2421 .status(&identity)
2422 .await
2423 .ok()
2424 .and_then(|status| status.agent_runtime_id.map(|id| id.as_str().to_string()));
2425
2426 if let Err(err) = runtime
2427 .reserve_identity_interaction(
2428 identity.as_str(),
2429 runtime_member_id.as_deref(),
2430 &interaction_id,
2431 origin,
2432 content_val,
2433 )
2434 .await
2435 {
2436 return maybe_error_response(
2437 is_notification,
2438 response_id,
2439 -32003,
2440 format!("failed to reserve interaction: {err}"),
2441 );
2442 }
2443
2444 match identity_rt.send(&identity, &content).await {
2445 Ok(token) => JsonRpcResponse {
2446 jsonrpc: JSONRPC_VERSION.to_string(),
2447 id: response_id,
2448 result: Some(serde_json::json!({
2449 "interaction_id": interaction_id,
2450 "fencing_token": token.get(),
2451 "stream": {
2452 "route": format!("/console/identity/{}/stream", identity.as_str()),
2453 "identity": identity.as_str(),
2454 }
2455 })),
2456 error: None,
2457 },
2458 Err(e) => {
2459 runtime
2460 .record_console_lifecycle(
2461 identity.as_str(),
2462 "interaction_failed",
2463 serde_json::json!({
2464 "interaction_id": interaction_id,
2465 "origin": origin,
2466 "error": e.to_string(),
2467 }),
2468 )
2469 .await;
2470 identity_error_response(response_id, &e)
2471 }
2472 }
2473 }
2474 "mobkit/dispatch" => {
2475 let identity_rt = match identity_ctx {
2476 Some(ctx) => &*ctx.runtime,
2477 None => return maybe_identity_not_configured(is_notification, response_id),
2478 };
2479 let identity_str = request
2480 .params
2481 .get("identity")
2482 .and_then(|v| v.as_str())
2483 .unwrap_or("");
2484 let target =
2485 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2486 {
2487 Ok(target) => target,
2488 Err(e) => {
2489 return maybe_error_response(
2490 is_notification,
2491 response_id,
2492 -32602,
2493 format!("invalid identity: {e}"),
2494 );
2495 }
2496 };
2497 let identity = target.identity.clone();
2498 if let Some(response) =
2499 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2500 {
2501 return if is_notification {
2502 String::new()
2503 } else {
2504 serialize_response(&response)
2505 };
2506 }
2507 let di_val = request
2508 .params
2509 .get("dispatch_input")
2510 .cloned()
2511 .unwrap_or(Value::Null);
2512 let content_val = di_val
2513 .get("content")
2514 .cloned()
2515 .unwrap_or_else(|| Value::String(String::new()));
2516 let content = match serde_json::from_value::<meerkat_core::ContentInput>(content_val) {
2517 Ok(content) => content,
2518 Err(err) => {
2519 return maybe_error_response(
2520 is_notification,
2521 response_id,
2522 -32602,
2523 format!("invalid dispatch_input.content: {err}"),
2524 );
2525 }
2526 };
2527 let origin_str = di_val
2528 .get("origin")
2529 .and_then(|v| v.as_str())
2530 .unwrap_or("system");
2531 let origin = match origin_str {
2532 "connector" => crate::identity_first::DispatchOrigin::Connector,
2533 "scheduler" => crate::identity_first::DispatchOrigin::Scheduler,
2534 "policy" => crate::identity_first::DispatchOrigin::Policy,
2535 "flow" => crate::identity_first::DispatchOrigin::Flow,
2536 _ => crate::identity_first::DispatchOrigin::System,
2537 };
2538 let correlation_id = di_val
2539 .get("correlation_id")
2540 .and_then(|v| v.as_str())
2541 .map(crate::identity_first::CorrelationId::new);
2542 let idempotency_key = di_val
2543 .get("idempotency_key")
2544 .and_then(|v| v.as_str())
2545 .map(crate::identity_first::DispatchIdempotencyKey::new);
2546 let dispatch_input = crate::identity_first::DispatchInput {
2547 content,
2548 origin,
2549 correlation_id,
2550 idempotency_key,
2551 };
2552 match identity_rt.dispatch(&identity, &dispatch_input).await {
2553 Ok((token, durable)) => JsonRpcResponse {
2554 jsonrpc: JSONRPC_VERSION.to_string(),
2555 id: response_id,
2556 result: Some(
2557 serde_json::json!({ "fencing_token": token.get(), "durable": durable }),
2558 ),
2559 error: None,
2560 },
2561 Err(e) => identity_error_response(response_id, &e),
2562 }
2563 }
2564 "mobkit/subscribe" => {
2565 let identity_rt = match identity_ctx {
2566 Some(ctx) => &*ctx.runtime,
2567 None => return maybe_identity_not_configured(is_notification, response_id),
2568 };
2569 let identity_str = request
2570 .params
2571 .get("identity")
2572 .and_then(|v| v.as_str())
2573 .unwrap_or("");
2574 let target =
2575 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2576 {
2577 Ok(target) => target,
2578 Err(e) => {
2579 return maybe_error_response(
2580 is_notification,
2581 response_id,
2582 -32602,
2583 format!("invalid identity: {e}"),
2584 );
2585 }
2586 };
2587 let identity = target.identity.clone();
2588 if let Some(response) =
2589 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2590 {
2591 return if is_notification {
2592 String::new()
2593 } else {
2594 serialize_response(&response)
2595 };
2596 }
2597 match identity_rt.subscribe(&identity).await {
2598 Ok(_receiver) => JsonRpcResponse {
2599 jsonrpc: JSONRPC_VERSION.to_string(),
2600 id: response_id,
2601 result: Some(serde_json::json!({
2602 "identity": identity.as_str(),
2603 "stream_id": identity.as_str(),
2604 "subscribed": true,
2605 })),
2606 error: None,
2607 },
2608 Err(e) => identity_error_response(response_id, &e),
2609 }
2610 }
2611 "mobkit/status_identity" => {
2612 let identity_rt = match identity_ctx {
2613 Some(ctx) => &*ctx.runtime,
2614 None => return maybe_identity_not_configured(is_notification, response_id),
2615 };
2616 let identity_str = request
2617 .params
2618 .get("identity")
2619 .and_then(|v| v.as_str())
2620 .unwrap_or("");
2621 let target =
2622 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2623 {
2624 Ok(target) => target,
2625 Err(e) => {
2626 return maybe_error_response(
2627 is_notification,
2628 response_id,
2629 -32602,
2630 format!("invalid identity: {e}"),
2631 );
2632 }
2633 };
2634 let identity = target.identity.clone();
2635 if let Some(response) =
2636 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2637 {
2638 return if is_notification {
2639 String::new()
2640 } else {
2641 serialize_response(&response)
2642 };
2643 }
2644 match identity_rt.status(&identity).await {
2645 Ok(status) => {
2646 let continuity_health =
2647 serde_json::to_value(&status.continuity_health).unwrap_or(Value::Null);
2648 let result = serde_json::json!({
2649 "state": identity_lifecycle_state_json(status.state),
2650 "identity": status.identity.as_str(),
2651 "agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
2652 "session_id": status.session_id.as_ref().map(ToString::to_string),
2653 "profile": status.profile.as_ref().map(meerkat_mob::ProfileName::as_str),
2654 "addressability": addressability_json(status.addressability),
2655 "display_name": status.display_name.as_ref().map(super::identity_first::DisplayName::as_str),
2656 "labels": status.labels,
2657 "generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
2658 "checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
2659 "continuity_health": continuity_health,
2660 "lease_healthy": status.lease.as_ref().map(|lease| lease.healthy),
2661 "lease": status.lease.as_ref().map(|lease| serde_json::json!({
2662 "fencing_token": lease.fencing_token.get(),
2663 "ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
2664 "healthy": lease.healthy,
2665 })),
2666 });
2667 JsonRpcResponse {
2668 jsonrpc: JSONRPC_VERSION.to_string(),
2669 id: response_id,
2670 result: Some(result),
2671 error: None,
2672 }
2673 }
2674 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
2675 if let Some(live) = target.live.as_ref() {
2676 JsonRpcResponse {
2677 jsonrpc: JSONRPC_VERSION.to_string(),
2678 id: response_id,
2679 result: Some(rpc_live_identity_status_json(live)),
2680 error: None,
2681 }
2682 } else {
2683 identity_error_response(response_id, &e)
2684 }
2685 }
2686 Err(e) => identity_error_response(response_id, &e),
2687 }
2688 }
2689 "mobkit/respawn" => {
2690 let identity_rt = match identity_ctx {
2691 Some(ctx) => &*ctx.runtime,
2692 None => return maybe_identity_not_configured(is_notification, response_id),
2693 };
2694 let identity_str = request
2695 .params
2696 .get("identity")
2697 .and_then(|v| v.as_str())
2698 .unwrap_or("");
2699 let target =
2700 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2701 {
2702 Ok(target) => target,
2703 Err(e) => {
2704 return maybe_error_response(
2705 is_notification,
2706 response_id,
2707 -32602,
2708 format!("invalid identity: {e}"),
2709 );
2710 }
2711 };
2712 let identity = target.identity.clone();
2713 if let Some(response) =
2714 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2715 {
2716 return if is_notification {
2717 String::new()
2718 } else {
2719 serialize_response(&response)
2720 };
2721 }
2722 let registered_status = match identity_rt.status(&identity).await {
2723 Ok(status) => Some(status),
2724 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => None,
2725 Err(e) => {
2726 let response = identity_error_response(response_id, &e);
2727 return if is_notification {
2728 String::new()
2729 } else {
2730 serialize_response(&response)
2731 };
2732 }
2733 };
2734 match identity_rt.respawn(&identity).await {
2735 Ok(mut record) => {
2736 let live_respawn_warning = match Box::pin(respawn_rpc_runtime_member_id(
2737 runtime,
2738 record.agent_runtime_id.as_str(),
2739 ))
2740 .await
2741 {
2742 Ok(live_result) => {
2743 let topology_restore_warning = live_result
2744 .get("topology_restore_warning")
2745 .filter(|warning| !warning.is_null())
2746 .cloned();
2747 let live_session_id =
2748 live_result.get("session_id").and_then(Value::as_str);
2749 if let Some(live_session_id) = live_session_id {
2750 match meerkat_core::types::SessionId::parse(live_session_id) {
2751 Ok(session_id) => {
2752 match identity_rt
2753 .rebind_session_after_live_respawn(
2754 &identity, session_id,
2755 )
2756 .await
2757 {
2758 Ok(updated_record) => {
2759 record = updated_record;
2760 topology_restore_warning
2761 }
2762 Err(err) => Some(serde_json::json!({
2763 "kind": "identity_rebind_failed_after_member_respawn",
2764 "message": err.to_string(),
2765 "identity": identity.as_str(),
2766 "agent_runtime_id": record.agent_runtime_id.as_str(),
2767 "live_session_id": live_session_id,
2768 })),
2769 }
2770 }
2771 Err(err) => Some(serde_json::json!({
2772 "kind": "member_respawn_session_id_invalid",
2773 "message": err.to_string(),
2774 "identity": identity.as_str(),
2775 "agent_runtime_id": record.agent_runtime_id.as_str(),
2776 "live_session_id": live_session_id,
2777 })),
2778 }
2779 } else {
2780 topology_restore_warning
2781 }
2782 }
2783 Err(err) => Some(serde_json::json!({
2784 "kind": "member_respawn_failed_after_identity_refresh",
2785 "message": err,
2786 "identity": identity.as_str(),
2787 "agent_runtime_id": record.agent_runtime_id.as_str(),
2788 })),
2789 };
2790 let cleanup_warning = if registered_status.is_some()
2791 && let Err(err) = retire_stale_rpc_members_for_identity(
2792 runtime,
2793 identity.as_str(),
2794 Some(record.agent_runtime_id.as_str()),
2795 )
2796 .await
2797 {
2798 Some(serde_json::json!({
2799 "kind": "stale_member_cleanup_failed_after_identity_respawn",
2800 "message": err,
2801 "identity": identity.as_str(),
2802 "agent_runtime_id": record.agent_runtime_id.as_str(),
2803 }))
2804 } else {
2805 None
2806 };
2807 runtime
2808 .record_console_lifecycle(
2809 identity.as_str(),
2810 "identity_respawned",
2811 serde_json::json!({
2812 "generation": record.generation.get(),
2813 "checkpoint_version": record.checkpoint_version.get(),
2814 "live_respawn_warning": live_respawn_warning.clone(),
2815 "cleanup_warning": cleanup_warning.clone(),
2816 }),
2817 )
2818 .await;
2819 JsonRpcResponse {
2820 jsonrpc: JSONRPC_VERSION.to_string(),
2821 id: response_id,
2822 result: Some(serde_json::json!({
2823 "identity": record.identity.as_str(),
2824 "agent_runtime_id": record.agent_runtime_id.as_str(),
2825 "session_id": record.session_id.to_string(),
2826 "generation": record.generation.get(),
2827 "checkpoint_version": record.checkpoint_version.get(),
2828 "live_respawn_warning": live_respawn_warning,
2829 "cleanup_warning": cleanup_warning,
2830 })),
2831 error: None,
2832 }
2833 }
2834 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
2835 if let Some(live) = target.live.as_ref() {
2836 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
2837 Ok(result) => {
2838 runtime
2839 .record_console_lifecycle(
2840 live.identity.as_str(),
2841 "identity_respawned",
2842 serde_json::json!({}),
2843 )
2844 .await;
2845 JsonRpcResponse {
2846 jsonrpc: JSONRPC_VERSION.to_string(),
2847 id: response_id,
2848 result: Some(result),
2849 error: None,
2850 }
2851 }
2852 Err(err) => JsonRpcResponse {
2853 jsonrpc: JSONRPC_VERSION.to_string(),
2854 id: response_id,
2855 result: None,
2856 error: Some(JsonRpcError {
2857 code: -32000,
2858 message: format!("respawn failed: {err}"),
2859 data: None,
2860 }),
2861 },
2862 }
2863 } else {
2864 identity_error_response(response_id, &e)
2865 }
2866 }
2867 Err(e) => identity_error_response(response_id, &e),
2868 }
2869 }
2870 "mobkit/retire" => {
2871 let identity_rt = match identity_ctx {
2872 Some(ctx) => &*ctx.runtime,
2873 None => return maybe_identity_not_configured(is_notification, response_id),
2874 };
2875 let identity_str = request
2876 .params
2877 .get("identity")
2878 .and_then(|v| v.as_str())
2879 .unwrap_or("");
2880 let target =
2881 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
2882 {
2883 Ok(target) => target,
2884 Err(e) => {
2885 return maybe_error_response(
2886 is_notification,
2887 response_id,
2888 -32602,
2889 format!("invalid identity: {e}"),
2890 );
2891 }
2892 };
2893 let identity = target.identity.clone();
2894 if let Some(response) =
2895 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
2896 {
2897 return if is_notification {
2898 String::new()
2899 } else {
2900 serialize_response(&response)
2901 };
2902 }
2903 let registered_status = match identity_rt.status(&identity).await {
2904 Ok(status) => Some(status),
2905 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => None,
2906 Err(e) => {
2907 let response = identity_error_response(response_id, &e);
2908 return if is_notification {
2909 String::new()
2910 } else {
2911 serialize_response(&response)
2912 };
2913 }
2914 };
2915 match identity_rt.retire(&identity).await {
2916 Ok(token) => {
2917 let keep_runtime_member_id = registered_status
2918 .as_ref()
2919 .and_then(|status| status.agent_runtime_id.as_ref())
2920 .filter(|_| identity_rt.has_session_bridge())
2921 .map(crate::identity_first::AgentRuntimeId::as_str);
2922 let cleanup_warning = if registered_status.is_some()
2923 && let Err(err) = retire_stale_rpc_members_for_identity(
2924 runtime,
2925 identity.as_str(),
2926 keep_runtime_member_id,
2927 )
2928 .await
2929 {
2930 Some(serde_json::json!({
2931 "kind": "stale_member_cleanup_failed_after_identity_retire",
2932 "message": err,
2933 "identity": identity.as_str(),
2934 }))
2935 } else {
2936 None
2937 };
2938 runtime
2939 .record_console_lifecycle(
2940 identity.as_str(),
2941 "identity_retired",
2942 serde_json::json!({
2943 "fencing_token": token.get(),
2944 "cleanup_warning": cleanup_warning.clone(),
2945 }),
2946 )
2947 .await;
2948 JsonRpcResponse {
2949 jsonrpc: JSONRPC_VERSION.to_string(),
2950 id: response_id,
2951 result: Some(serde_json::json!({
2952 "fencing_token": token.get(),
2953 "cleanup_warning": cleanup_warning,
2954 })),
2955 error: None,
2956 }
2957 }
2958 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
2959 if let Some(live) = target.live.as_ref() {
2960 match retire_rpc_live_identity(runtime, live).await {
2961 Ok(()) => {
2962 runtime
2963 .record_console_lifecycle(
2964 live.identity.as_str(),
2965 "identity_retired",
2966 serde_json::json!({}),
2967 )
2968 .await;
2969 JsonRpcResponse {
2970 jsonrpc: JSONRPC_VERSION.to_string(),
2971 id: response_id,
2972 result: Some(
2973 serde_json::json!({ "identity": live.identity.as_str() }),
2974 ),
2975 error: None,
2976 }
2977 }
2978 Err(err) => JsonRpcResponse {
2979 jsonrpc: JSONRPC_VERSION.to_string(),
2980 id: response_id,
2981 result: None,
2982 error: Some(JsonRpcError {
2983 code: -32000,
2984 message: format!("retire failed: {err}"),
2985 data: None,
2986 }),
2987 },
2988 }
2989 } else {
2990 identity_error_response(response_id, &e)
2991 }
2992 }
2993 Err(e) => identity_error_response(response_id, &e),
2994 }
2995 }
2996 "mobkit/reset" => {
2997 let identity_rt = match identity_ctx {
2998 Some(ctx) => &*ctx.runtime,
2999 None => return maybe_identity_not_configured(is_notification, response_id),
3000 };
3001 let identity_str = request
3002 .params
3003 .get("identity")
3004 .and_then(|v| v.as_str())
3005 .unwrap_or("");
3006 let target =
3007 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3008 {
3009 Ok(target) => target,
3010 Err(e) => {
3011 return maybe_error_response(
3012 is_notification,
3013 response_id,
3014 -32602,
3015 format!("invalid identity: {e}"),
3016 );
3017 }
3018 };
3019 let identity = target.identity.clone();
3020 if let Some(response) =
3021 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3022 {
3023 return if is_notification {
3024 String::new()
3025 } else {
3026 serialize_response(&response)
3027 };
3028 }
3029 let _registered_status = match identity_rt.status(&identity).await {
3030 Ok(status) => {
3031 if !identity_rt.has_session_bridge() {
3032 let response = rpc_reset_requires_session_bridge_response(response_id);
3033 return if is_notification {
3034 String::new()
3035 } else {
3036 serialize_response(&response)
3037 };
3038 }
3039 status
3040 }
3041 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3042 if let Some(live) = target.live.as_ref() {
3043 let response =
3044 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
3045 Ok(result) => {
3046 runtime
3047 .record_console_lifecycle(
3048 live.identity.as_str(),
3049 "identity_reset",
3050 serde_json::json!({}),
3051 )
3052 .await;
3053 JsonRpcResponse {
3054 jsonrpc: JSONRPC_VERSION.to_string(),
3055 id: response_id,
3056 result: Some(result),
3057 error: None,
3058 }
3059 }
3060 Err(err) => JsonRpcResponse {
3061 jsonrpc: JSONRPC_VERSION.to_string(),
3062 id: response_id,
3063 result: None,
3064 error: Some(JsonRpcError {
3065 code: -32000,
3066 message: format!("reset failed: {err}"),
3067 data: None,
3068 }),
3069 },
3070 };
3071 return if is_notification {
3072 String::new()
3073 } else {
3074 serialize_response(&response)
3075 };
3076 }
3077 let response = identity_error_response(response_id, &e);
3078 return if is_notification {
3079 String::new()
3080 } else {
3081 serialize_response(&response)
3082 };
3083 }
3084 Err(e) => {
3085 let response = identity_error_response(response_id, &e);
3086 return if is_notification {
3087 String::new()
3088 } else {
3089 serialize_response(&response)
3090 };
3091 }
3092 };
3093 match identity_rt.reset(&identity).await {
3094 Ok(record) => {
3095 let cleanup_warning = if let Err(err) = retire_stale_rpc_members_for_identity(
3096 runtime,
3097 identity.as_str(),
3098 Some(record.agent_runtime_id.as_str()),
3099 )
3100 .await
3101 {
3102 Some(serde_json::json!({
3103 "kind": "stale_member_cleanup_failed_after_identity_reset",
3104 "message": err,
3105 "identity": identity.as_str(),
3106 "agent_runtime_id": record.agent_runtime_id.as_str(),
3107 }))
3108 } else {
3109 None
3110 };
3111 runtime
3112 .record_console_lifecycle(
3113 identity.as_str(),
3114 "identity_reset",
3115 serde_json::json!({
3116 "generation": record.generation.get(),
3117 "checkpoint_version": record.checkpoint_version.get(),
3118 "cleanup_warning": cleanup_warning.clone(),
3119 }),
3120 )
3121 .await;
3122 JsonRpcResponse {
3123 jsonrpc: JSONRPC_VERSION.to_string(),
3124 id: response_id,
3125 result: Some(serde_json::json!({
3126 "identity": record.identity.as_str(),
3127 "agent_runtime_id": record.agent_runtime_id.as_str(),
3128 "session_id": record.session_id.to_string(),
3129 "generation": record.generation.get(),
3130 "checkpoint_version": record.checkpoint_version.get(),
3131 "cleanup_warning": cleanup_warning,
3132 })),
3133 error: None,
3134 }
3135 }
3136 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3137 if let Some(live) = target.live.as_ref() {
3138 match Box::pin(respawn_rpc_live_identity(runtime, live)).await {
3139 Ok(result) => {
3140 runtime
3141 .record_console_lifecycle(
3142 live.identity.as_str(),
3143 "identity_reset",
3144 serde_json::json!({}),
3145 )
3146 .await;
3147 JsonRpcResponse {
3148 jsonrpc: JSONRPC_VERSION.to_string(),
3149 id: response_id,
3150 result: Some(result),
3151 error: None,
3152 }
3153 }
3154 Err(err) => JsonRpcResponse {
3155 jsonrpc: JSONRPC_VERSION.to_string(),
3156 id: response_id,
3157 result: None,
3158 error: Some(JsonRpcError {
3159 code: -32000,
3160 message: format!("reset failed: {err}"),
3161 data: None,
3162 }),
3163 },
3164 }
3165 } else {
3166 identity_error_response(response_id, &e)
3167 }
3168 }
3169 Err(e) => identity_error_response(response_id, &e),
3170 }
3171 }
3172 "mobkit/delete_identity" => {
3173 let identity_rt = match identity_ctx {
3174 Some(ctx) => &*ctx.runtime,
3175 None => return maybe_identity_not_configured(is_notification, response_id),
3176 };
3177 let identity_str = request
3178 .params
3179 .get("identity")
3180 .and_then(|v| v.as_str())
3181 .unwrap_or("");
3182 let target =
3183 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3184 {
3185 Ok(target) => target,
3186 Err(e) => {
3187 return maybe_error_response(
3188 is_notification,
3189 response_id,
3190 -32602,
3191 format!("invalid identity: {e}"),
3192 );
3193 }
3194 };
3195 let identity = target.identity.clone();
3196 if let Some(response) =
3197 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3198 {
3199 return if is_notification {
3200 String::new()
3201 } else {
3202 serialize_response(&response)
3203 };
3204 }
3205 let registered_status = match identity_rt.status(&identity).await {
3206 Ok(status) => status,
3207 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3208 if target.live.is_some() {
3209 let response = JsonRpcResponse {
3210 jsonrpc: JSONRPC_VERSION.to_string(),
3211 id: response_id,
3212 result: None,
3213 error: Some(JsonRpcError {
3214 code: -32602,
3215 message: format!(
3216 "delete_identity requires durable identity: {} is live-only",
3217 identity.as_str()
3218 ),
3219 data: Some(serde_json::json!({
3220 "kind": "live_only_identity_delete_unsupported",
3221 "identity": identity.as_str(),
3222 })),
3223 }),
3224 };
3225 return if is_notification {
3226 String::new()
3227 } else {
3228 serialize_response(&response)
3229 };
3230 }
3231 let response = identity_error_response(response_id, &e);
3232 return if is_notification {
3233 String::new()
3234 } else {
3235 serialize_response(&response)
3236 };
3237 }
3238 Err(e) => {
3239 let response = identity_error_response(response_id, &e);
3240 return if is_notification {
3241 String::new()
3242 } else {
3243 serialize_response(&response)
3244 };
3245 }
3246 };
3247 let keep_runtime_member_id = registered_status
3248 .agent_runtime_id
3249 .as_ref()
3250 .filter(|_| identity_rt.has_session_bridge())
3251 .map(crate::identity_first::AgentRuntimeId::as_str);
3252 match identity_rt.delete_identity(&identity).await {
3253 Ok(()) => {
3254 let cleanup_warning = if 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_delete",
3263 "identity": identity.as_str(),
3264 "message": err,
3265 }))
3266 } else {
3267 None
3268 };
3269 runtime
3270 .record_console_lifecycle(
3271 identity.as_str(),
3272 "identity_deleted",
3273 serde_json::json!({
3274 "cleanup_warning": cleanup_warning,
3275 }),
3276 )
3277 .await;
3278 JsonRpcResponse {
3279 jsonrpc: JSONRPC_VERSION.to_string(),
3280 id: response_id,
3281 result: Some(serde_json::json!({
3282 "identity": identity.as_str(),
3283 "cleanup_warning": cleanup_warning,
3284 })),
3285 error: None,
3286 }
3287 }
3288 Err(e) => identity_error_response(response_id, &e),
3289 }
3290 }
3291 "mobkit/inspect_identity" => {
3292 let identity_rt = match identity_ctx {
3293 Some(ctx) => &*ctx.runtime,
3294 None => return maybe_identity_not_configured(is_notification, response_id),
3295 };
3296 let identity_str = request
3297 .params
3298 .get("identity")
3299 .and_then(|v| v.as_str())
3300 .unwrap_or("");
3301 let target =
3302 match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
3303 {
3304 Ok(target) => target,
3305 Err(e) => {
3306 return maybe_error_response(
3307 is_notification,
3308 response_id,
3309 -32602,
3310 format!("invalid identity: {e}"),
3311 );
3312 }
3313 };
3314 let identity = target.identity.clone();
3315 let status = identity_rt.status(&identity).await;
3316 if let Some(response) =
3317 rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
3318 {
3319 return if is_notification {
3320 String::new()
3321 } else {
3322 serialize_response(&response)
3323 };
3324 }
3325 match identity_rt.inspect(&identity).await {
3326 Ok(inspection) => {
3327 let status = status.ok();
3328 JsonRpcResponse {
3329 jsonrpc: JSONRPC_VERSION.to_string(),
3330 id: response_id,
3331 result: Some(serde_json::json!({
3332 "identity": identity.as_str(),
3333 "state": status.as_ref().map(|status| identity_lifecycle_state_json(status.state)),
3334 "profile": status.as_ref().and_then(|status| status.profile.as_ref().map(meerkat_mob::ProfileName::as_str)),
3335 "addressability": status.as_ref().map(|status| addressability_json(status.addressability)),
3336 "display_name": status.as_ref().and_then(|status| status.display_name.as_ref().map(super::identity_first::DisplayName::as_str)),
3337 "labels": status.as_ref().map(|status| status.labels.clone()).unwrap_or_default(),
3338 "generation": status.as_ref().and_then(|status| status.generation.map(super::identity_first::ContinuityGeneration::get)),
3339 "checkpoint_version": status.as_ref().and_then(|status| status.checkpoint_version.map(super::identity_first::CheckpointVersion::get)),
3340 "continuity_health": status.as_ref().and_then(|status| serde_json::to_value(&status.continuity_health).ok()).unwrap_or(Value::Null),
3341 "lease_healthy": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| lease.healthy)),
3342 "continuity": status.as_ref().map(|status| serde_json::json!({
3343 "generation": status.generation.map(super::identity_first::ContinuityGeneration::get),
3344 "checkpoint_version": status.checkpoint_version.map(super::identity_first::CheckpointVersion::get),
3345 "session_id": status.session_id.as_ref().map(ToString::to_string),
3346 "agent_runtime_id": status.agent_runtime_id.as_ref().map(super::identity_first::AgentRuntimeId::as_str),
3347 })).unwrap_or_else(|| serde_json::json!({})),
3348 "lease": status.as_ref().and_then(|status| status.lease.as_ref().map(|lease| serde_json::json!({
3349 "fencing_token": lease.fencing_token.get(),
3350 "ttl_remaining_ms": lease.ttl_remaining.as_millis() as u64,
3351 "healthy": lease.healthy,
3352 }))),
3353 "output_preview": inspection.output_preview,
3354 "is_final": inspection.is_final,
3355 "peer_reachable_count": inspection.peer_reachable_count,
3356 })),
3357 error: None,
3358 }
3359 }
3360 Err(e @ crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {
3361 if let Some(live) = target.live.as_ref() {
3362 JsonRpcResponse {
3363 jsonrpc: JSONRPC_VERSION.to_string(),
3364 id: response_id,
3365 result: Some(rpc_live_identity_inspect_json(runtime, live).await),
3366 error: None,
3367 }
3368 } else {
3369 identity_error_response(response_id, &e)
3370 }
3371 }
3372 Err(e) => identity_error_response(response_id, &e),
3373 }
3374 }
3375 "mobkit/reconcile_identity" => {
3376 let ctx = match identity_ctx {
3377 Some(ctx) => ctx,
3378 None => return maybe_identity_not_configured(is_notification, response_id),
3379 };
3380 let roster_specs = match ctx
3382 .roster_provider
3383 .roster(&crate::identity_first::RosterContext {
3384 mob_definition: None,
3385 previous_identities: Vec::new(),
3386 })
3387 .await
3388 {
3389 Ok(specs) => specs,
3390 Err(e) => {
3391 return maybe_error_response(
3392 is_notification,
3393 response_id,
3394 -32603,
3395 format!("roster provider failed: {e}"),
3396 );
3397 }
3398 };
3399 match crate::identity_first::restore_flow(
3400 &ctx.runtime,
3401 &roster_specs,
3402 ctx.topology_provider.as_deref(),
3403 ctx.customizer.as_deref(),
3404 )
3405 .await
3406 {
3407 Ok(result) => {
3408 let outcomes: serde_json::Map<String, Value> = result
3409 .outcomes
3410 .iter()
3411 .map(|(id, outcome)| {
3412 let val = match outcome {
3413 crate::identity_first::RestoreOutcome::Created {
3414 record, ..
3415 } => {
3416 serde_json::json!({
3417 "outcome": "created",
3418 "identity": record.identity.as_str(),
3419 "agent_runtime_id": record.agent_runtime_id.as_str(),
3420 "session_id": record.session_id.to_string(),
3421 "generation": record.generation.get(),
3422 })
3423 }
3424 crate::identity_first::RestoreOutcome::Dormant {
3425 record, ..
3426 } => {
3427 serde_json::json!({
3428 "outcome": "dormant",
3429 "identity": id.as_str(),
3430 "agent_runtime_id": record.as_ref().map(|record| record.agent_runtime_id.as_str()),
3431 "session_id": record.as_ref().map(|record| record.session_id.to_string()),
3432 "generation": record.as_ref().map(|record| record.generation.get()),
3433 })
3434 }
3435 crate::identity_first::RestoreOutcome::Resumed {
3436 record, ..
3437 } => {
3438 serde_json::json!({
3439 "outcome": "resumed",
3440 "identity": record.identity.as_str(),
3441 "agent_runtime_id": record.agent_runtime_id.as_str(),
3442 "session_id": record.session_id.to_string(),
3443 "generation": record.generation.get(),
3444 })
3445 }
3446 crate::identity_first::RestoreOutcome::Broken(failure) => {
3447 serde_json::json!({
3448 "outcome": "broken",
3449 "identity": failure.identity.as_str(),
3450 "detail": failure.detail,
3451 })
3452 }
3453 };
3454 (id.to_string(), val)
3455 })
3456 .collect();
3457 JsonRpcResponse {
3458 jsonrpc: JSONRPC_VERSION.to_string(),
3459 id: response_id,
3460 result: Some(serde_json::json!({
3461 "outcomes": outcomes,
3462 "managed_edges": result.managed_edges.len(),
3463 })),
3464 error: None,
3465 }
3466 }
3467 Err(e) => identity_error_response(response_id, &e),
3468 }
3469 }
3470 method if method.contains('/') && !method.starts_with("mobkit/") => {
3471 let module_id = method
3472 .split('/')
3473 .next()
3474 .map(ToString::to_string)
3475 .unwrap_or_default();
3476 let route = runtime
3477 .route_module_call(
3478 &ModuleRouteRequest {
3479 module_id: module_id.clone(),
3480 method: method.to_string(),
3481 params: request.params,
3482 },
3483 timeout,
3484 )
3485 .await;
3486 match route {
3487 Ok(response) => JsonRpcResponse {
3488 jsonrpc: JSONRPC_VERSION.to_string(),
3489 id: response_id,
3490 result: Some(serde_json::json!({
3491 "module_id": response.module_id,
3492 "method": response.method,
3493 "payload": response.payload
3494 })),
3495 error: None,
3496 },
3497 Err(ModuleRouteError::UnloadedModule(module_id)) => JsonRpcResponse {
3498 jsonrpc: JSONRPC_VERSION.to_string(),
3499 id: response_id,
3500 result: None,
3501 error: Some(JsonRpcError {
3502 code: -32601,
3503 message: format!("Module '{module_id}' not loaded"),
3504 data: None,
3505 }),
3506 },
3507 Err(err) => JsonRpcResponse {
3508 jsonrpc: JSONRPC_VERSION.to_string(),
3509 id: response_id,
3510 result: None,
3511 error: Some(JsonRpcError {
3512 code: -32000,
3513 message: format!("Module route failed: {err:?}"),
3514 data: None,
3515 }),
3516 },
3517 }
3518 }
3519 _ => JsonRpcResponse {
3520 jsonrpc: JSONRPC_VERSION.to_string(),
3521 id: response_id,
3522 result: None,
3523 error: Some(JsonRpcError {
3524 code: -32601,
3525 message: "Method not found".to_string(),
3526 data: None,
3527 }),
3528 },
3529 };
3530 if is_notification {
3531 String::new()
3532 } else {
3533 serialize_response(&response)
3534 }
3535}
3536
3537fn build_models_catalog_result() -> Value {
3538 let entries: Vec<Value> = meerkat_models::catalog()
3539 .iter()
3540 .filter_map(|e| {
3541 let mut val = serde_json::to_value(e).ok()?;
3542 if let Some(provider) = meerkat_core::Provider::parse_strict(e.provider)
3543 && let Some(profile) = meerkat_models::profile_for(provider, e.id)
3544 && let Ok(p) = serde_json::to_value(&profile)
3545 {
3546 val["profile"] = p;
3547 }
3548 Some(val)
3549 })
3550 .collect();
3551 let defaults: Vec<Value> = meerkat_models::provider_defaults()
3552 .iter()
3553 .filter_map(|d| serde_json::to_value(d).ok())
3554 .collect();
3555 serde_json::json!({
3556 "models": entries,
3557 "provider_defaults": defaults,
3558 })
3559}
3560
3561#[derive(Debug, Clone)]
3562struct RpcLiveIdentityAlias {
3563 identity: crate::identity_first::AgentIdentity,
3564 runtime_member_id: String,
3565 member: meerkat_mob::runtime::MobMemberListEntry,
3566 session_id: Option<String>,
3567}
3568
3569#[derive(Debug, Clone)]
3570struct RpcIdentityControlTarget {
3571 identity: crate::identity_first::AgentIdentity,
3572 live: Option<RpcLiveIdentityAlias>,
3573}
3574
3575fn rpc_member_durable_identity(member: &meerkat_mob::runtime::MobMemberListEntry) -> String {
3576 member
3577 .labels
3578 .get("agent_identity")
3579 .filter(|value| !value.trim().is_empty())
3580 .cloned()
3581 .unwrap_or_else(|| {
3584 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
3585 })
3586}
3587
3588async fn resolve_rpc_live_identity_alias(
3589 runtime: &UnifiedRuntime,
3590 requested_identity: &str,
3591) -> Result<Option<RpcLiveIdentityAlias>, String> {
3592 let matches = resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
3593 if matches.len() > 1 {
3594 return Err(format!(
3595 "ambiguous live identity alias {requested_identity}: candidates [{}]",
3596 matches
3597 .iter()
3598 .map(|entry| entry.runtime_member_id.clone())
3599 .collect::<Vec<_>>()
3600 .join(", ")
3601 ));
3602 }
3603 Ok(matches.into_iter().next())
3604}
3605
3606async fn resolve_rpc_live_runtime_member_alias(
3607 runtime: &UnifiedRuntime,
3608 runtime_member_id: &str,
3609) -> Result<Option<RpcLiveIdentityAlias>, String> {
3610 let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
3611 let handle = runtime.mob_handle();
3612 let Some(member) = handle
3613 .list_members_including_retiring()
3614 .await
3615 .into_iter()
3616 .find(|entry| entry.agent_identity == requested_member_id)
3617 else {
3618 return Ok(None);
3619 };
3620 if !rpc_live_identity_alias_member_visible(&member) {
3621 return Ok(None);
3622 }
3623 let durable_identity = rpc_member_durable_identity(&member);
3624 let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
3625 .map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
3626 let session_id = handle
3627 .resolve_bridge_session_id_observation(&member.agent_identity)
3628 .await
3629 .map(|session_id| session_id.to_string());
3630 Ok(Some(RpcLiveIdentityAlias {
3631 identity,
3632 runtime_member_id: crate::member_comms_id::runtime_alias_str(
3633 member.agent_identity.as_str(),
3634 )
3635 .into_owned(),
3636 member,
3637 session_id,
3638 }))
3639}
3640
3641async fn rpc_runtime_member_alias_exists_hidden(
3642 runtime: &UnifiedRuntime,
3643 runtime_member_id: &str,
3644) -> bool {
3645 let requested_member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
3646 runtime
3647 .mob_handle()
3648 .list_members_including_retiring()
3649 .await
3650 .into_iter()
3651 .find(|entry| entry.agent_identity == requested_member_id)
3652 .is_some_and(|member| !rpc_live_identity_alias_member_visible(&member))
3653}
3654
3655async fn rpc_live_identity_alias_exists_hidden(
3656 runtime: &UnifiedRuntime,
3657 requested_identity: &str,
3658) -> bool {
3659 let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
3660 runtime
3661 .mob_handle()
3662 .list_members_including_retiring()
3663 .await
3664 .into_iter()
3665 .any(|member| {
3666 (member.agent_identity == requested_member_id
3667 || member
3668 .labels
3669 .get("agent_identity")
3670 .is_some_and(|identity| identity == requested_identity))
3671 && !rpc_live_identity_alias_member_visible(&member)
3672 })
3673}
3674
3675async fn resolve_rpc_live_identity_alias_candidates(
3676 runtime: &UnifiedRuntime,
3677 requested_identity: &str,
3678) -> Result<Vec<RpcLiveIdentityAlias>, String> {
3679 let requested_member_id = crate::member_comms_id::mob_member_id(requested_identity);
3680 let handle = runtime.mob_handle();
3681 let members = handle.list_members_including_retiring().await;
3682 let exact_matches = members
3683 .iter()
3684 .filter(|entry| entry.agent_identity == requested_member_id)
3685 .cloned()
3686 .collect::<Vec<_>>();
3687 let label_matches = members
3688 .iter()
3689 .filter(|entry| {
3690 entry
3691 .labels
3692 .get("agent_identity")
3693 .is_some_and(|identity| identity == requested_identity)
3694 })
3695 .cloned()
3696 .collect::<Vec<_>>();
3697 let mut matches = exact_matches;
3698 matches.extend(label_matches);
3699 let mut seen_member_ids = BTreeSet::new();
3700 matches.retain(|entry| seen_member_ids.insert(entry.agent_identity.to_string()));
3701 let mut aliases = Vec::with_capacity(matches.len());
3702 for member in matches {
3703 if !rpc_live_identity_alias_member_visible(&member) {
3704 continue;
3705 }
3706 let durable_identity = rpc_member_durable_identity(&member);
3707 let identity = crate::identity_first::AgentIdentity::parse(&durable_identity)
3708 .map_err(|err| format!("invalid projected identity {durable_identity}: {err}"))?;
3709 let session_id = handle
3710 .resolve_bridge_session_id_observation(&member.agent_identity)
3711 .await
3712 .map(|session_id| session_id.to_string());
3713 aliases.push(RpcLiveIdentityAlias {
3714 identity,
3715 runtime_member_id: crate::member_comms_id::runtime_alias_str(
3716 member.agent_identity.as_str(),
3717 )
3718 .into_owned(),
3719 member,
3720 session_id,
3721 });
3722 }
3723 Ok(aliases)
3724}
3725
3726fn rpc_live_identity_alias_member_visible(
3727 member: &meerkat_mob::runtime::MobMemberListEntry,
3728) -> bool {
3729 rpc_live_identity_alias_visible(member.role.as_str(), &member.labels)
3730}
3731
3732fn rpc_live_identity_alias_visible(
3733 member_role: &str,
3734 labels: &std::collections::BTreeMap<String, String>,
3735) -> bool {
3736 let projected_role = labels
3737 .get("role")
3738 .map(String::as_str)
3739 .unwrap_or(member_role);
3740 !is_implicit_delegate_member(member_role, labels)
3741 && !is_implicit_delegate_member(projected_role, labels)
3742}
3743
3744async fn resolve_rpc_identity_control_target(
3745 runtime: &UnifiedRuntime,
3746 identity_rt: &crate::identity_first::IdentityRuntime,
3747 requested_identity: &str,
3748) -> Result<RpcIdentityControlTarget, String> {
3749 if requested_identity.starts_with("rt:") {
3750 for status in identity_rt.statuses().await {
3751 if status
3752 .agent_runtime_id
3753 .as_ref()
3754 .is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
3755 {
3756 let identity = status.identity;
3757 let registered_live =
3758 resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
3759 if let Some(registered) = registered_live {
3760 return Ok(RpcIdentityControlTarget {
3761 identity,
3762 live: Some(registered),
3763 });
3764 }
3765 if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
3766 return Err(format!("identity hidden by policy: {requested_identity}"));
3767 }
3768 let durable_live_candidates =
3769 resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
3770 let durable_live = if durable_live_candidates.len() > 1 {
3771 return Err(format!(
3772 "ambiguous live identity alias {}: candidates [{}]",
3773 identity.as_str(),
3774 durable_live_candidates
3775 .iter()
3776 .map(|alias| alias.runtime_member_id.clone())
3777 .collect::<Vec<_>>()
3778 .join(", ")
3779 ));
3780 } else {
3781 durable_live_candidates.into_iter().next()
3782 };
3783 return Ok(RpcIdentityControlTarget {
3784 identity,
3785 live: durable_live,
3786 });
3787 }
3788 }
3789 let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
3790 if let Some(live_alias) = live {
3791 let live_identity_candidates =
3792 resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
3793 .await?;
3794 if live_identity_candidates.len() > 1 {
3795 return Err(format!(
3796 "ambiguous live identity alias {}: candidates [{}]",
3797 live_alias.identity.as_str(),
3798 live_identity_candidates
3799 .iter()
3800 .map(|alias| alias.runtime_member_id.clone())
3801 .collect::<Vec<_>>()
3802 .join(", ")
3803 ));
3804 }
3805 return Ok(RpcIdentityControlTarget {
3806 identity: live_alias.identity.clone(),
3807 live: Some(live_alias),
3808 });
3809 }
3810 if rpc_runtime_member_alias_exists_hidden(runtime, requested_identity).await {
3811 return Err(format!("identity hidden by policy: {requested_identity}"));
3812 }
3813 return Err(format!("runtime identity not found: {requested_identity}"));
3814 }
3815 if let Ok(identity) = crate::identity_first::AgentIdentity::parse(requested_identity) {
3816 match identity_rt.status(&identity).await {
3817 Ok(status) => {
3818 let registered_live = match status.agent_runtime_id.as_ref() {
3819 Some(runtime_id) => {
3820 resolve_rpc_live_runtime_member_alias(runtime, runtime_id.as_str()).await?
3821 }
3822 None => None,
3823 };
3824 if let Some(registered) = registered_live {
3825 return Ok(RpcIdentityControlTarget {
3826 identity,
3827 live: Some(registered),
3828 });
3829 }
3830 if let Some(runtime_id) = status.agent_runtime_id.as_ref()
3831 && rpc_runtime_member_alias_exists_hidden(runtime, runtime_id.as_str()).await
3832 {
3833 return Err(format!("identity hidden by policy: {requested_identity}"));
3834 }
3835 let requested_live_candidates =
3836 resolve_rpc_live_identity_alias_candidates(runtime, requested_identity).await?;
3837 let requested_live = if requested_live_candidates.len() > 1 {
3838 return Err(format!(
3839 "ambiguous live identity alias {requested_identity}: candidates [{}]",
3840 requested_live_candidates
3841 .iter()
3842 .map(|alias| alias.runtime_member_id.clone())
3843 .collect::<Vec<_>>()
3844 .join(", ")
3845 ));
3846 } else {
3847 requested_live_candidates.into_iter().next()
3848 };
3849 return Ok(RpcIdentityControlTarget {
3850 identity,
3851 live: requested_live,
3852 });
3853 }
3854 Err(crate::identity_first::IdentityRuntimeError::UnknownIdentity(_)) => {}
3855 Err(err) => return Err(err.to_string()),
3856 }
3857 }
3858 for status in identity_rt.statuses().await {
3859 if status
3860 .agent_runtime_id
3861 .as_ref()
3862 .is_some_and(|runtime_id| runtime_id.as_str() == requested_identity)
3863 {
3864 let identity = status.identity;
3865 let registered_live =
3866 resolve_rpc_live_runtime_member_alias(runtime, requested_identity).await?;
3867 let durable_live_candidates =
3868 resolve_rpc_live_identity_alias_candidates(runtime, identity.as_str()).await?;
3869 let durable_live = if durable_live_candidates.len() > 1 {
3870 return Err(format!(
3871 "ambiguous live identity alias {}: candidates [{}]",
3872 identity.as_str(),
3873 durable_live_candidates
3874 .iter()
3875 .map(|alias| alias.runtime_member_id.clone())
3876 .collect::<Vec<_>>()
3877 .join(", ")
3878 ));
3879 } else {
3880 durable_live_candidates.into_iter().next()
3881 };
3882 let live = match (registered_live, durable_live) {
3883 (Some(registered), Some(durable))
3884 if registered.runtime_member_id == durable.runtime_member_id =>
3885 {
3886 Some(registered)
3887 }
3888 (Some(registered), None) => Some(registered),
3889 (Some(_registered), Some(durable)) => Some(durable),
3890 (None, durable) => durable,
3891 };
3892 return Ok(RpcIdentityControlTarget { identity, live });
3893 }
3894 }
3895 let live = resolve_rpc_live_identity_alias(runtime, requested_identity).await?;
3896 if let Some(live_alias) = live {
3897 if let Some(bound_status) = identity_rt.statuses().await.into_iter().find(|status| {
3898 status
3899 .agent_runtime_id
3900 .as_ref()
3901 .is_some_and(|runtime_id| runtime_id.as_str() == live_alias.runtime_member_id)
3902 }) && bound_status.identity != live_alias.identity
3903 {
3904 return Err(format!(
3905 "stale live identity alias: live console alias {} resolves to {}, but identity runtime binding belongs to {}",
3906 live_alias.identity.as_str(),
3907 live_alias.runtime_member_id,
3908 bound_status.identity.as_str(),
3909 ));
3910 }
3911 let live_identity_candidates =
3912 resolve_rpc_live_identity_alias_candidates(runtime, live_alias.identity.as_str())
3913 .await?;
3914 if live_identity_candidates.len() > 1 {
3915 return Err(format!(
3916 "ambiguous live identity alias {}: candidates [{}]",
3917 live_alias.identity.as_str(),
3918 live_identity_candidates
3919 .iter()
3920 .map(|alias| alias.runtime_member_id.clone())
3921 .collect::<Vec<_>>()
3922 .join(", ")
3923 ));
3924 }
3925 return Ok(RpcIdentityControlTarget {
3926 identity: live_alias.identity.clone(),
3927 live: Some(live_alias),
3928 });
3929 }
3930 if rpc_live_identity_alias_exists_hidden(runtime, requested_identity).await {
3931 return Err(format!("identity hidden by policy: {requested_identity}"));
3932 }
3933 let identity = crate::identity_first::AgentIdentity::parse(requested_identity)
3934 .map_err(|err| err.to_string())?;
3935 Ok(RpcIdentityControlTarget {
3936 identity,
3937 live: None,
3938 })
3939}
3940
3941fn rpc_reset_requires_session_bridge_response(response_id: Value) -> JsonRpcResponse {
3942 JsonRpcResponse {
3943 jsonrpc: JSONRPC_VERSION.to_string(),
3944 id: response_id,
3945 result: None,
3946 error: Some(JsonRpcError {
3947 code: -32602,
3948 message: "reset requires an identity runtime with a session bridge".to_string(),
3949 data: Some(serde_json::json!({
3950 "kind": "identity_reset_requires_session_bridge",
3951 })),
3952 }),
3953 }
3954}
3955
3956fn rpc_live_alias_matches_status_runtime(
3957 alias: Option<&RpcLiveIdentityAlias>,
3958 status: &crate::identity_first::IdentityStatus,
3959) -> bool {
3960 let Some(alias) = alias else {
3961 return true;
3962 };
3963 let session_matches = match (
3964 status.session_id.as_ref().map(ToString::to_string),
3965 alias.session_id.as_deref(),
3966 ) {
3967 (Some(status_session), Some(live_session)) => status_session == live_session,
3968 _ => true,
3969 };
3970 status
3971 .agent_runtime_id
3972 .as_ref()
3973 .is_some_and(|runtime_id| runtime_id.as_str() == alias.runtime_member_id)
3974 && alias.identity == status.identity
3975 && session_matches
3976}
3977
3978async fn rpc_stale_live_alias_error_response(
3979 identity_rt: &crate::identity_first::IdentityRuntime,
3980 target: &RpcIdentityControlTarget,
3981 response_id: Value,
3982) -> Option<JsonRpcResponse> {
3983 let live = target.live.as_ref()?;
3984 let Ok(status) = identity_rt.status(&target.identity).await else {
3985 return None;
3986 };
3987 if rpc_live_alias_matches_status_runtime(Some(live), &status) {
3988 return None;
3989 }
3990 Some(JsonRpcResponse {
3991 jsonrpc: JSONRPC_VERSION.to_string(),
3992 id: response_id,
3993 result: None,
3994 error: Some(JsonRpcError {
3995 code: -32000,
3996 message: format!(
3997 "identity runtime binding for {} points at {}, but requested live member is {}",
3998 target.identity.as_str(),
3999 status
4000 .agent_runtime_id
4001 .as_ref()
4002 .map(crate::identity_first::AgentRuntimeId::as_str)
4003 .unwrap_or("<none>"),
4004 live.runtime_member_id
4005 ),
4006 data: Some(serde_json::json!({
4007 "kind": "stale_identity_runtime_binding",
4008 "identity": target.identity.as_str(),
4009 "registered_runtime_member_id": status.agent_runtime_id.as_ref().map(crate::identity_first::AgentRuntimeId::as_str),
4010 "live_runtime_member_id": live.runtime_member_id,
4011 "registered_session_id": status.session_id.as_ref().map(ToString::to_string),
4012 "live_session_id": live.session_id,
4013 })),
4014 }),
4015 })
4016}
4017
4018fn rpc_member_is_addressable(member: &meerkat_mob::runtime::MobMemberListEntry) -> bool {
4019 member
4020 .labels
4021 .get("addressable")
4022 .map(|value| !value.eq_ignore_ascii_case("false"))
4023 .unwrap_or(true)
4024}
4025
4026fn rpc_live_identity_status_json(alias: &RpcLiveIdentityAlias) -> Value {
4027 serde_json::json!({
4028 "state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
4029 "identity": alias.identity.as_str(),
4030 "agent_runtime_id": alias.runtime_member_id,
4031 "session_id": alias.session_id,
4032 "profile": alias.member.role.to_string(),
4033 "addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
4034 "display_name": alias.member.labels.get("display_name"),
4035 "labels": alias.member.labels,
4036 "generation": Value::Null,
4037 "checkpoint_version": Value::Null,
4038 "continuity_health": Value::Null,
4039 "lease_healthy": Value::Null,
4040 "lease": Value::Null,
4041 })
4042}
4043
4044async fn rpc_live_identity_inspect_json(
4045 runtime: &UnifiedRuntime,
4046 alias: &RpcLiveIdentityAlias,
4047) -> Value {
4048 let snapshot = runtime
4049 .mob_handle()
4050 .member_status(&crate::member_comms_id::mob_member_id(
4051 alias.runtime_member_id.as_str(),
4052 ))
4053 .await
4054 .ok();
4055 serde_json::json!({
4056 "identity": alias.identity.as_str(),
4057 "state": crate::mob_handle_runtime::member_status_state_string(alias.member.status),
4058 "profile": alias.member.role.to_string(),
4059 "addressability": if rpc_member_is_addressable(&alias.member) { "addressable" } else { "internal_only" },
4060 "display_name": alias.member.labels.get("display_name"),
4061 "labels": alias.member.labels,
4062 "generation": Value::Null,
4063 "checkpoint_version": Value::Null,
4064 "continuity_health": Value::Null,
4065 "lease_healthy": Value::Null,
4066 "continuity": {
4067 "generation": Value::Null,
4068 "checkpoint_version": Value::Null,
4069 "session_id": alias.session_id,
4070 "agent_runtime_id": alias.runtime_member_id,
4071 },
4072 "lease": Value::Null,
4073 "output_preview": snapshot.as_ref().and_then(|snapshot| snapshot.output_preview.clone()),
4074 "is_final": snapshot.as_ref().map(|snapshot| snapshot.is_final).unwrap_or(false),
4075 "peer_reachable_count": alias.member.wired_to.len(),
4076 })
4077}
4078
4079async fn retire_rpc_live_identity(
4080 runtime: &UnifiedRuntime,
4081 alias: &RpcLiveIdentityAlias,
4082) -> Result<(), String> {
4083 retire_rpc_runtime_member_id(runtime, alias.runtime_member_id.as_str()).await
4084}
4085
4086async fn retire_rpc_runtime_member_id(
4087 runtime: &UnifiedRuntime,
4088 runtime_member_id: &str,
4089) -> Result<(), String> {
4090 match runtime
4091 .mob_handle()
4092 .retire(crate::member_comms_id::mob_member_id(runtime_member_id))
4093 .await
4094 {
4095 Ok(()) => Ok(()),
4096 Err(err) if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) => Ok(()),
4097 Err(err) => Err(err.to_string()),
4098 }
4099}
4100
4101fn rpc_member_id_matches_durable_identity(member_id: &str, durable_identity: &str) -> bool {
4102 crate::member_comms_id::runtime_alias_str(member_id) == durable_identity
4105}
4106
4107async fn retire_stale_rpc_members_for_identity(
4108 runtime: &UnifiedRuntime,
4109 durable_identity: &str,
4110 keep_runtime_member_id: Option<&str>,
4111) -> Result<(), String> {
4112 let stale_members = runtime
4113 .mob_handle()
4114 .list_members_including_retiring()
4115 .await
4116 .into_iter()
4117 .filter(|member| {
4118 if !rpc_live_identity_alias_member_visible(member) {
4119 return false;
4120 }
4121 (rpc_member_id_matches_durable_identity(
4122 member.agent_identity.as_str(),
4123 durable_identity,
4124 ) || member
4125 .labels
4126 .get("agent_identity")
4127 .is_some_and(|identity| identity == durable_identity))
4128 && keep_runtime_member_id
4129 .map(|keep| {
4130 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str())
4132 != keep
4133 })
4134 .unwrap_or(true)
4135 })
4136 .map(|member| {
4138 crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()).into_owned()
4139 })
4140 .collect::<Vec<_>>();
4141 for member_id in stale_members {
4142 retire_rpc_runtime_member_id(runtime, &member_id).await?;
4143 }
4144 Ok(())
4145}
4146
4147async fn respawn_rpc_live_identity(
4148 runtime: &UnifiedRuntime,
4149 alias: &RpcLiveIdentityAlias,
4150) -> Result<Value, String> {
4151 let mut result = Box::pin(respawn_rpc_runtime_member_id(
4152 runtime,
4153 alias.runtime_member_id.as_str(),
4154 ))
4155 .await?;
4156 result["identity"] = serde_json::json!(alias.identity.as_str());
4157 Ok(result)
4158}
4159
4160async fn respawn_rpc_runtime_member_id(
4161 runtime: &UnifiedRuntime,
4162 runtime_member_id: &str,
4163) -> Result<Value, String> {
4164 let handle = runtime.mob_handle();
4165 let member_id = crate::member_comms_id::mob_member_id(runtime_member_id);
4166 let entry_before_respawn = handle.get_member(&member_id).await.ok().flatten();
4169 let mut topology_restore_warning = None;
4170 match handle.respawn(member_id.clone(), None).await {
4171 Ok(_receipt) => {}
4172 Err(err) => {
4173 if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(&err) {
4174 tracing::warn!(
4175 member_id = %member_id,
4176 failed_peer_count = failed_peer_ids.len(),
4177 failed_peer_ids = ?failed_peer_ids,
4178 "rpc member respawn restored member with isolated peer edges; continuing degraded respawn"
4179 );
4180 topology_restore_warning = Some(topology_restore_warning_json(&failed_peer_ids));
4181 } else if mob_methods::lifecycle_archive_cleanup_completed(&err.to_string()) {
4182 if handle
4185 .get_member(&member_id)
4186 .await
4187 .map_err(|lookup_err| lookup_err.to_string())?
4188 .is_none()
4189 && let Some(entry) = entry_before_respawn
4190 {
4191 let mut spec =
4192 meerkat_mob::SpawnMemberSpec::new(entry.role.clone(), member_id.clone());
4193 if !entry.labels.is_empty() {
4194 spec = spec.with_labels(entry.labels.clone());
4195 }
4196 handle
4197 .ensure_member(spec)
4198 .await
4199 .map_err(|ensure_err| ensure_err.to_string())?;
4200 }
4201 } else {
4202 return Err(err.to_string());
4203 }
4204 }
4205 }
4206 let session_id = handle
4207 .resolve_bridge_session_id_observation(&member_id)
4208 .await
4209 .map(|session_id| session_id.to_string());
4210 Ok(serde_json::json!({
4211 "agent_runtime_id": runtime_member_id,
4212 "session_id": session_id,
4213 "generation": Value::Null,
4214 "checkpoint_version": Value::Null,
4215 "topology_restore_warning": topology_restore_warning,
4216 }))
4217}
4218
4219fn identity_not_configured(response_id: Value) -> String {
4220 error_response(response_id, -32601, "identity-first runtime not configured")
4221}
4222
4223fn maybe_identity_not_configured(is_notification: bool, response_id: Value) -> String {
4224 if is_notification {
4225 String::new()
4226 } else {
4227 identity_not_configured(response_id)
4228 }
4229}
4230
4231fn addressability_json(addressability: crate::identity_first::AgentAddressability) -> &'static str {
4232 match addressability {
4233 crate::identity_first::AgentAddressability::Addressable => "addressable",
4234 crate::identity_first::AgentAddressability::InternalOnly => "internal_only",
4235 }
4236}
4237
4238fn identity_lifecycle_state_json(
4241 state: crate::identity_first::IdentityLifecycleState,
4242) -> &'static str {
4243 state.wire_str()
4244}
4245
4246fn identity_error_response(
4247 response_id: Value,
4248 err: &crate::identity_first::IdentityRuntimeError,
4249) -> JsonRpcResponse {
4250 use crate::identity_first::IdentityRuntimeError;
4251 let (code, message) = match err {
4252 IdentityRuntimeError::UnknownIdentity(id) => (-32001, format!("unknown identity: {id}")),
4253 IdentityRuntimeError::NotAddressable(na) => {
4254 (-32002, format!("not addressable: {}", na.identity))
4255 }
4256 IdentityRuntimeError::NoActiveLease(id) => (-32003, format!("no active lease: {id}")),
4257 IdentityRuntimeError::LeaseLost(id) => (-32005, format!("lease lost: {id}")),
4263 _ => (-32603, format!("{err}")),
4264 };
4265 JsonRpcResponse {
4266 jsonrpc: JSONRPC_VERSION.to_string(),
4267 id: response_id,
4268 result: None,
4269 error: Some(JsonRpcError {
4270 code,
4271 message,
4272 data: None,
4273 }),
4274 }
4275}
4276
4277fn error_response(response_id: Value, code: i64, message: impl Into<String>) -> String {
4278 let message = message.into();
4279 let ambiguous_alias_rest = message
4280 .strip_prefix("ambiguous live identity alias ")
4281 .or_else(|| message.strip_prefix("invalid identity: ambiguous live identity alias "));
4282 let stale_live_alias_rest = message
4283 .strip_prefix("stale live identity alias: live console alias ")
4284 .or_else(|| {
4285 message.strip_prefix("invalid identity: stale live identity alias: live console alias ")
4286 });
4287 let hidden_policy_identity = message
4288 .strip_prefix("identity hidden by policy: ")
4289 .or_else(|| message.strip_prefix("invalid identity: identity hidden by policy: "));
4290 let data = if let Some(rest) = ambiguous_alias_rest {
4291 let (identity, candidates) = rest
4292 .split_once(": candidates [")
4293 .map(|(identity, candidates)| {
4294 (
4295 identity.to_string(),
4296 candidates
4297 .trim_end_matches(']')
4298 .split(',')
4299 .map(str::trim)
4300 .filter(|value| !value.is_empty())
4301 .map(str::to_string)
4302 .collect::<Vec<_>>(),
4303 )
4304 })
4305 .unwrap_or_else(|| (rest.to_string(), Vec::new()));
4306 Some(serde_json::json!({
4307 "kind": "ambiguous_live_identity_alias",
4308 "identity": identity,
4309 "candidates": candidates,
4310 }))
4311 } else if let Some(rest) = stale_live_alias_rest {
4312 let (identity, rest) = rest.split_once(" resolves to ").unwrap_or((rest, ""));
4313 let (runtime_member_id, bound_identity) = rest
4314 .split_once(", but identity runtime binding belongs to ")
4315 .unwrap_or((rest, ""));
4316 Some(serde_json::json!({
4317 "kind": "stale_live_identity_alias",
4318 "identity": identity,
4319 "live_runtime_member_id": runtime_member_id,
4320 "bound_identity": bound_identity,
4321 }))
4322 } else {
4323 hidden_policy_identity.map(|identity| {
4324 serde_json::json!({
4325 "kind": "identity_hidden_by_policy",
4326 "identity": identity,
4327 })
4328 })
4329 };
4330 serialize_response(&JsonRpcResponse {
4331 jsonrpc: JSONRPC_VERSION.to_string(),
4332 id: response_id,
4333 result: None,
4334 error: Some(JsonRpcError {
4335 code,
4336 message,
4337 data,
4338 }),
4339 })
4340}
4341
4342fn maybe_error_response(
4343 is_notification: bool,
4344 response_id: Value,
4345 code: i64,
4346 message: impl Into<String>,
4347) -> String {
4348 if is_notification {
4349 String::new()
4350 } else {
4351 error_response(response_id, code, message)
4352 }
4353}
4354
4355fn serialize_response(response: &JsonRpcResponse) -> String {
4356 serde_json::to_string(response).unwrap_or_else(|_| {
4357 r#"{"jsonrpc":"2.0","id":null,"error":{"code":-32603,"message":"Internal error"}}"#
4358 .to_string()
4359 })
4360}
4361
4362#[cfg(test)]
4363#[allow(clippy::expect_used)]
4364mod tests {
4365 use super::{
4366 error_response, handle_unified_rpc_json, identity_error_response,
4367 resolve_rpc_identity_control_target, rpc_live_identity_alias_visible,
4368 rpc_member_id_matches_durable_identity,
4369 };
4370 use crate::identity_first::contracts::RosterProvider;
4371 use crate::identity_first::{
4372 AgentAddressability, AgentIdentity, AgentRuntimeId, CheckpointVersion,
4373 ContinuityGeneration, ContinuityRecord, DurabilityPolicy, DurableAgentSpec, FencingToken,
4374 IdentityLifecycleState, IdentityRuntime, IdentityRuntimeConfig, LeaseGrant,
4375 LocalContinuityStore, LocalLeaseProvider, RosterContext, RosterError,
4376 };
4377 use crate::{
4378 DiscoverySpec, IdentityFirstContext, MobBootstrapOptions, MobBootstrapSpec, MobKitConfig,
4379 UnifiedRuntime,
4380 };
4381 use async_trait::async_trait;
4382 use meerkat::{AgentFactory, Config, build_ephemeral_service};
4383 use meerkat_client::TestClient;
4384 use meerkat_mob::{MobDefinition, MobStorage, SpawnMemberSpec};
4385 use serde_json::{Value, json};
4386 use std::collections::BTreeMap;
4387 use std::sync::Arc;
4388 use std::time::Duration;
4389
4390 #[derive(Debug, Default)]
4391 struct EmptyRosterProvider;
4392
4393 #[async_trait]
4394 impl RosterProvider for EmptyRosterProvider {
4395 async fn roster(
4396 &self,
4397 _context: &RosterContext,
4398 ) -> Result<Vec<DurableAgentSpec>, RosterError> {
4399 Ok(Vec::new())
4400 }
4401 }
4402
4403 fn rpc_test_mob_spec(
4404 temp_dir: &tempfile::TempDir,
4405 ) -> Result<MobBootstrapSpec, Box<dyn std::error::Error + Send + Sync>> {
4406 let session_path = temp_dir.path().join("sessions");
4407 std::fs::create_dir_all(&session_path)?;
4408 let factory = AgentFactory::new(&session_path).comms(true);
4409 let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
4410 let definition = MobDefinition::from_toml(
4411 r#"
4412[mob]
4413id = "rpc-identity-alias-test"
4414
4415[profiles.worker]
4416model = "gpt-5.5"
4417external_addressable = true
4418
4419[profiles.worker.tools]
4420comms = true
4421"#,
4422 )?;
4423 Ok(
4424 MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
4425 .with_options(MobBootstrapOptions {
4426 allow_ephemeral_sessions: true,
4427 notify_orchestrator_on_resume: true,
4428 default_llm_client: Some(Arc::new(TestClient::default())),
4429 }),
4430 )
4431 }
4432
4433 #[tokio::test]
4434 async fn unified_capabilities_separate_mobpack_authoring_from_runtime_controls()
4435 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
4436 let temp_dir = tempfile::tempdir()?;
4437 let runtime = Box::pin(
4438 UnifiedRuntime::builder()
4439 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
4440 .module_config(MobKitConfig {
4441 modules: Vec::new(),
4442 discovery: DiscoverySpec {
4443 namespace: "rpc-authoring-capabilities-test".to_string(),
4444 modules: Vec::new(),
4445 },
4446 pre_spawn: Vec::new(),
4447 })
4448 .timeout(Duration::from_secs(1))
4449 .build(),
4450 )
4451 .await?;
4452
4453 let response: Value = serde_json::from_str(
4454 &handle_unified_rpc_json(
4455 &runtime,
4456 &json!({
4457 "jsonrpc": "2.0",
4458 "id": 1,
4459 "method": "mobkit/capabilities",
4460 })
4461 .to_string(),
4462 Duration::from_secs(1),
4463 None,
4464 None,
4465 )
4466 .await,
4467 )?;
4468
4469 assert!(response["error"].is_null(), "{response:#?}");
4470 let methods = response["result"]["methods"]
4471 .as_array()
4472 .expect("methods array")
4473 .iter()
4474 .filter_map(Value::as_str)
4475 .collect::<Vec<_>>();
4476 for method in super::MOBPACK_AUTHORING_METHODS {
4477 assert!(
4478 methods.contains(method),
4479 "missing authoring method {method}"
4480 );
4481 }
4482 assert_eq!(
4483 response["result"]["authoring_capabilities"]["domain"],
4484 json!("mobpack_authoring")
4485 );
4486 assert_eq!(
4487 response["result"]["authoring_capabilities"]["runtime_mutation"],
4488 json!(false)
4489 );
4490 assert_eq!(
4491 response["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
4492 json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
4493 );
4494 assert_eq!(
4495 response["result"]["authoring_capabilities"]["methods"]
4496 .as_array()
4497 .expect("authoring methods")
4498 .iter()
4499 .filter_map(Value::as_str)
4500 .collect::<Vec<_>>(),
4501 super::MOBPACK_AUTHORING_METHODS
4502 );
4503 assert_eq!(
4504 response["result"]["authoring_capabilities"]["deploy_command"],
4505 json!("rkat mob run")
4506 );
4507
4508 Ok(())
4509 }
4510
4511 #[tokio::test]
4512 async fn unified_rpc_dispatches_mobpack_authoring_methods()
4513 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
4514 let temp_dir = tempfile::tempdir()?;
4515 let runtime = Box::pin(
4516 UnifiedRuntime::builder()
4517 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
4518 .module_config(MobKitConfig {
4519 modules: Vec::new(),
4520 discovery: DiscoverySpec {
4521 namespace: "rpc-authoring-dispatch-test".to_string(),
4522 modules: Vec::new(),
4523 },
4524 pre_spawn: Vec::new(),
4525 })
4526 .timeout(Duration::from_secs(1))
4527 .build(),
4528 )
4529 .await?;
4530
4531 let response: Value = serde_json::from_str(
4532 &handle_unified_rpc_json(
4533 &runtime,
4534 &json!({
4535 "jsonrpc": "2.0",
4536 "id": 1,
4537 "method": "mobkit/mobpacks/schema",
4538 })
4539 .to_string(),
4540 Duration::from_secs(1),
4541 None,
4542 None,
4543 )
4544 .await,
4545 )?;
4546
4547 assert!(response["error"].is_null(), "{response:#?}");
4548 assert_eq!(
4549 response["result"]["media_type"],
4550 json!("application/vnd.meerkat.mobpack")
4551 );
4552 assert_eq!(
4553 response["result"]["commands"]["deploy_rpc"],
4554 json!("mobkit/mobpacks/deploy")
4555 );
4556 assert_eq!(
4557 response["result"]["deploy_settings"]["runtime_backed"],
4558 json!(true)
4559 );
4560 assert_eq!(
4561 response["result"]["deploy_settings"]["authoring_provider"]["runtime_binding"],
4562 json!("bound")
4563 );
4564 assert_eq!(
4565 response["result"]["deploy_settings"]["provenance"]["source"],
4566 json!("UnifiedRuntime.authoring_provider.deploy_target")
4567 );
4568 assert!(response["result"]["sample_mobpacks"].is_null());
4569 assert!(response["result"]["agent_definitions"].is_null());
4570
4571 let catalogs: Value = serde_json::from_str(
4572 &handle_unified_rpc_json(
4573 &runtime,
4574 &json!({
4575 "jsonrpc": "2.0",
4576 "id": 2,
4577 "method": "mobkit/mobpacks/catalogs",
4578 })
4579 .to_string(),
4580 Duration::from_secs(1),
4581 None,
4582 None,
4583 )
4584 .await,
4585 )?;
4586 assert!(catalogs["error"].is_null(), "{catalogs:#?}");
4587 assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
4588 assert_eq!(
4589 catalogs["result"]["authoring_provider"]["id"],
4590 json!("unified_runtime")
4591 );
4592 assert_eq!(
4593 catalogs["result"]["authoring_provider"]["runtime_binding"],
4594 json!("bound")
4595 );
4596 assert_eq!(
4597 catalogs["result"]["sources"]["runtime"],
4598 json!("unified_runtime")
4599 );
4600 assert_eq!(
4601 catalogs["result"]["sources"]["runtime_binding"],
4602 json!("bound")
4603 );
4604 assert!(
4605 catalogs["result"]["runtime_unavailable_reason"].is_null(),
4606 "{catalogs:#?}"
4607 );
4608 assert_eq!(
4609 catalogs["result"]["catalog_snapshot"]["runtime_backed"],
4610 json!(true)
4611 );
4612 assert_eq!(
4613 catalogs["result"]["authoring_provider"]["deploy_target"]["command"],
4614 json!("rkat mob run")
4615 );
4616 assert!(
4617 catalogs["result"]["authoring_provider"]["runtime_methods"]
4618 .as_array()
4619 .is_some_and(|methods| methods.contains(&json!("mobkit/mobpacks/deploy"))),
4620 "{catalogs:#?}"
4621 );
4622
4623 Ok(())
4624 }
4625
4626 #[test]
4627 fn identity_lease_lost_maps_off_capability_unavailable_code() {
4628 let identity = AgentIdentity::parse("review:singleton").expect("valid identity");
4634 let err = crate::identity_first::IdentityRuntimeError::LeaseLost(identity);
4635 let response = identity_error_response(json!("req-1"), &err);
4636 let error = response.error.expect("lease-lost must surface an error");
4637 assert_ne!(
4638 error.code, -32004,
4639 "LeaseLost must not use the capability code"
4640 );
4641 assert_eq!(
4642 error.code, -32005,
4643 "LeaseLost has its own identity-plane code"
4644 );
4645 assert!(error.message.contains("lease lost"));
4646 }
4647
4648 #[test]
4649 fn mobpack_authoring_rpc_helper_preserves_runtime_catalog_binding()
4650 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
4651 let runtime = crate::mobpack::MobpackRuntimeCatalogState {
4652 loaded_modules: vec!["editor-host".to_string()],
4653 runtime_methods: vec![
4654 "mobkit/mobpacks/catalogs".to_string(),
4655 "mobkit/mobpacks/apply_operation".to_string(),
4656 "mobkit/mobpacks/deploy".to_string(),
4657 ],
4658 has_contact_directory: true,
4659 has_peer_mob_handles: false,
4660 has_inproc_contacts: false,
4661 runtime_flow_rows: vec![json!({
4662 "id": "runtime_rpc_main",
4663 "source": "mobkit/runtime/flow_projection",
4664 "document": {
4665 "mob_id": "runtime_rpc",
4666 "flow": { "name": "main", "steps": [] },
4667 "members": []
4668 },
4669 "validation": { "ok": true }
4670 })],
4671 runtime_agent_definition_sources: vec![json!({
4672 "id": "runtime_profiles_rpc",
4673 "name": "Runtime RPC profiles",
4674 "source": "mobkit/runtime/agent-definitions",
4675 "document": {
4676 "mob_id": "runtime_rpc",
4677 "members": [{
4678 "id": "m_runtime_reviewer",
4679 "name": "Runtime reviewer",
4680 "role": "runtime_reviewer",
4681 "profileBinding": "inline",
4682 "model": "gpt-5.5",
4683 "runtimeMode": "turn_driven",
4684 "tools": ["builtins"],
4685 "skills": ["mob.runtime.review"],
4686 "schema": ""
4687 }],
4688 "schemas": []
4689 }
4690 })],
4691 runtime_skill_realms: vec![json!({
4692 "id": "runtime_rpc",
4693 "label": "Runtime RPC",
4694 "source": "mobkit/runtime/skills",
4695 "skills": [{
4696 "id": "mob.runtime.review",
4697 "label": "Runtime review",
4698 "source": "inline",
4699 "content": "Review runtime work."
4700 }]
4701 })],
4702 };
4703
4704 let catalogs = super::handle_mobpack_authoring_rpc_with_runtime(
4705 "mobkit/mobpacks/catalogs",
4706 &json!({}),
4707 json!(1),
4708 Some(&runtime),
4709 )
4710 .expect("catalogs method");
4711 let catalogs: Value = serde_json::to_value(catalogs)?;
4712 assert_eq!(catalogs["result"]["runtime_backed"], json!(true));
4713 assert_eq!(
4714 catalogs["result"]["authoring_provider"]["runtime_binding"],
4715 json!("bound")
4716 );
4717 assert_eq!(
4718 catalogs["result"]["runtime_flows"][0]["id"],
4719 json!("runtime_rpc_main")
4720 );
4721 let listed = super::handle_mobpack_authoring_rpc_with_runtime(
4722 "mobkit/mobpacks/list",
4723 &json!({}),
4724 json!(2),
4725 Some(&runtime),
4726 )
4727 .expect("list method");
4728 let listed: Value = serde_json::to_value(listed)?;
4729 assert_eq!(listed["result"]["runtime_backed"], json!(true));
4730 assert_eq!(listed["result"]["rows"][0]["id"], json!("runtime_rpc_main"));
4731
4732 let fetched = super::handle_mobpack_authoring_rpc_with_runtime(
4733 "mobkit/mobpacks/get",
4734 &json!({ "id": "runtime_rpc_main" }),
4735 json!(3),
4736 Some(&runtime),
4737 )
4738 .expect("get method");
4739 let fetched: Value = serde_json::to_value(fetched)?;
4740 assert_eq!(fetched["result"]["runtime_backed"], json!(true));
4741 assert_eq!(fetched["result"]["row"]["id"], json!("runtime_rpc_main"));
4742
4743 let definitions = super::handle_mobpack_authoring_rpc_with_runtime(
4744 "mobkit/agent_definitions/list",
4745 &json!({}),
4746 json!(4),
4747 Some(&runtime),
4748 )
4749 .expect("agent definitions method");
4750 let definitions: Value = serde_json::to_value(definitions)?;
4751 let runtime_definition = definitions["result"]["agent_definitions"]
4752 .as_array()
4753 .and_then(|rows| {
4754 rows.iter()
4755 .find(|row| row["sourceOrigin"] == "mobkit/runtime/agent-definitions")
4756 })
4757 .expect("runtime profile definition");
4758 assert_eq!(runtime_definition["role"], json!("runtime_reviewer"));
4759 assert_eq!(
4760 runtime_definition["toolDefinitions"][0]["id"],
4761 json!("builtins")
4762 );
4763 assert_eq!(
4764 runtime_definition["skillDefinitions"][0]["id"],
4765 json!("mob.runtime.review")
4766 );
4767 assert_eq!(
4768 definitions["result"]["catalog_snapshot"]["runtime_backed"],
4769 json!(true)
4770 );
4771
4772 let tools = super::handle_mobpack_authoring_rpc_with_runtime(
4773 "mobkit/tools/catalog",
4774 &json!({}),
4775 json!(5),
4776 Some(&runtime),
4777 )
4778 .expect("tools catalog method");
4779 let tools: Value = serde_json::to_value(tools)?;
4780 let mob_tool = tools["result"]["tool_catalog"]
4781 .as_array()
4782 .expect("tools")
4783 .iter()
4784 .find(|tool| tool["id"] == "mob")
4785 .expect("mob tool");
4786 assert_eq!(mob_tool["runtime_availability"]["available"], json!(false));
4787
4788 let agents = super::handle_mobpack_authoring_rpc_with_runtime(
4789 "mobkit/agent_definitions/list",
4790 &json!({}),
4791 json!(3),
4792 Some(&runtime),
4793 )
4794 .expect("agent definitions method");
4795 let agents: Value = serde_json::to_value(agents)?;
4796 assert_eq!(agents["result"]["runtime_backed"], json!(true));
4797 let planner = agents["result"]["agent_definitions"]
4798 .as_array()
4799 .expect("agent definitions")
4800 .iter()
4801 .find(|definition| definition["role"] == "planner")
4802 .expect("planner definition");
4803 let planner_mob_tool = planner["toolDefinitions"]
4804 .as_array()
4805 .expect("planner tools")
4806 .iter()
4807 .find(|tool| tool["id"] == "mob")
4808 .expect("planner mob tool");
4809 assert_eq!(
4810 planner_mob_tool["runtimeAvailability"]["state"],
4811 json!("unavailable")
4812 );
4813
4814 Ok(())
4815 }
4816
4817 #[test]
4818 fn module_rpc_dispatches_mobpack_authoring_methods()
4819 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
4820 let config = MobKitConfig {
4821 modules: Vec::new(),
4822 discovery: DiscoverySpec {
4823 namespace: "module-rpc-authoring-dispatch-test".to_string(),
4824 modules: Vec::new(),
4825 },
4826 pre_spawn: Vec::new(),
4827 };
4828 let mut runtime = crate::start_mobkit_runtime(config, Vec::new(), Duration::from_secs(1))?;
4829
4830 let capabilities: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
4831 &mut runtime,
4832 &json!({
4833 "jsonrpc": "2.0",
4834 "id": 1,
4835 "method": "mobkit/capabilities",
4836 })
4837 .to_string(),
4838 Duration::from_secs(1),
4839 ))?;
4840 let methods = capabilities["result"]["methods"]
4841 .as_array()
4842 .expect("methods array")
4843 .iter()
4844 .filter_map(Value::as_str)
4845 .collect::<Vec<_>>();
4846 for method in super::MOBPACK_AUTHORING_METHODS {
4847 assert!(
4848 methods.contains(method),
4849 "missing authoring method {method}"
4850 );
4851 }
4852 assert_eq!(
4853 capabilities["result"]["authoring_capabilities"]["runtime_mutation"],
4854 json!(false)
4855 );
4856 assert_eq!(
4857 capabilities["result"]["authoring_capabilities"]["host_mutation_methods"]["mobkit/mobpacks/deploy"],
4858 json!("when execute=true, writes a mobpack archive and runs rkat mob run on the host")
4859 );
4860
4861 let schema: Value = serde_json::from_str(&super::handle_mobkit_rpc_json(
4862 &mut runtime,
4863 &json!({
4864 "jsonrpc": "2.0",
4865 "id": 2,
4866 "method": "mobkit/mobpacks/schema",
4867 })
4868 .to_string(),
4869 Duration::from_secs(1),
4870 ))?;
4871 assert!(schema["error"].is_null(), "{schema:#?}");
4872 assert_eq!(
4873 schema["result"]["commands"]["deploy_rpc"],
4874 json!("mobkit/mobpacks/deploy")
4875 );
4876 assert!(schema["result"]["agent_definitions"].is_null());
4877
4878 let _ = runtime.shutdown();
4879 Ok(())
4880 }
4881
4882 #[test]
4883 fn generated_runtime_ids_match_their_durable_identity_prefix() {
4884 assert!(!rpc_member_id_matches_durable_identity(
4885 "rt:review:singleton:0",
4886 "review:singleton",
4887 ));
4888 assert!(!rpc_member_id_matches_durable_identity(
4889 "review:singleton:gen1",
4890 "review:singleton",
4891 ));
4892 assert!(!rpc_member_id_matches_durable_identity(
4893 "review:singleton:1",
4894 "review:singleton",
4895 ));
4896 assert!(!rpc_member_id_matches_durable_identity(
4897 "rt:reviewer:singleton:0",
4898 "review:singleton",
4899 ));
4900 assert!(!rpc_member_id_matches_durable_identity(
4901 "rt:review:singleton:qa:0",
4902 "review:singleton",
4903 ));
4904 assert!(!rpc_member_id_matches_durable_identity(
4905 "review:singleton:qa",
4906 "review:singleton",
4907 ));
4908 }
4909
4910 #[test]
4911 fn rpc_live_identity_visibility_matches_delegate_projection_labels() {
4912 assert!(rpc_live_identity_alias_visible("worker", &BTreeMap::new()));
4913
4914 let mut labels = BTreeMap::new();
4915 labels.insert("role".to_string(), "delegate".to_string());
4916 labels.insert("source_mob_id".to_string(), "mob-a".to_string());
4917 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
4918 assert!(!rpc_live_identity_alias_visible("worker", &labels));
4919 assert!(!rpc_live_identity_alias_visible("delegate", &labels));
4920 }
4921
4922 #[test]
4923 fn ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error> {
4924 let response: Value = serde_json::from_str(&error_response(
4925 json!(1),
4926 -32602,
4927 "ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
4928 ))?;
4929
4930 assert_eq!(
4931 response["error"]["data"]["kind"],
4932 json!("ambiguous_live_identity_alias")
4933 );
4934 assert_eq!(
4935 response["error"]["data"]["identity"],
4936 json!("review:singleton")
4937 );
4938 assert_eq!(
4939 response["error"]["data"]["candidates"],
4940 json!(["rt:review:singleton:0", "rt:review:singleton:1"])
4941 );
4942 Ok(())
4943 }
4944
4945 #[test]
4946 fn wrapped_ambiguous_live_alias_errors_include_structured_data() -> Result<(), serde_json::Error>
4947 {
4948 let response: Value = serde_json::from_str(&error_response(
4949 json!(1),
4950 -32602,
4951 "invalid identity: ambiguous live identity alias review:singleton: candidates [rt:review:singleton:0, rt:review:singleton:1]",
4952 ))?;
4953
4954 assert_eq!(
4955 response["error"]["data"]["kind"],
4956 json!("ambiguous_live_identity_alias")
4957 );
4958 assert_eq!(
4959 response["error"]["data"]["identity"],
4960 json!("review:singleton")
4961 );
4962 assert_eq!(
4963 response["error"]["data"]["candidates"],
4964 json!(["rt:review:singleton:0", "rt:review:singleton:1"])
4965 );
4966 Ok(())
4967 }
4968
4969 #[tokio::test]
4970 async fn runtime_id_live_only_resolution_rejects_duplicate_projected_identity()
4971 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
4972 let temp_dir = tempfile::tempdir()?;
4973 let runtime = Box::pin(
4974 UnifiedRuntime::builder()
4975 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
4976 .module_config(MobKitConfig {
4977 modules: Vec::new(),
4978 discovery: DiscoverySpec {
4979 namespace: "rpc-identity-alias-test".to_string(),
4980 modules: Vec::new(),
4981 },
4982 pre_spawn: Vec::new(),
4983 })
4984 .timeout(Duration::from_secs(1))
4985 .build(),
4986 )
4987 .await?;
4988 for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
4989 let mut labels = BTreeMap::new();
4990 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
4991 runtime
4992 .spawn(
4993 SpawnMemberSpec::from_wire(
4994 "worker".to_string(),
4995 runtime_id.to_string(),
4996 Some("You are a duplicate Review Agent.".into()),
4997 None,
4998 None,
4999 )
5000 .with_labels(labels),
5001 )
5002 .await?;
5003 }
5004 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5005 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5006 lease_provider: Arc::new(LocalLeaseProvider::new()),
5007 runtime_instance_id: "rpc-identity-alias-test".to_string(),
5008 has_runtime_store: true,
5009 durability_policy: DurabilityPolicy::SyncWriteThrough,
5010 bridge: None,
5011 default_timeout: None,
5012 });
5013
5014 let err =
5015 resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:0")
5016 .await
5017 .expect_err("runtime-id live-only fallback should reject duplicate durable alias");
5018 assert!(
5019 err.contains("ambiguous live identity alias review:singleton"),
5020 "unexpected error: {err}"
5021 );
5022
5023 Ok(())
5024 }
5025
5026 #[tokio::test]
5027 async fn durable_resolution_prefers_registered_live_binding_over_stale_duplicates()
5028 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5029 let temp_dir = tempfile::tempdir()?;
5030 let runtime = Box::pin(
5031 UnifiedRuntime::builder()
5032 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5033 .module_config(MobKitConfig {
5034 modules: Vec::new(),
5035 discovery: DiscoverySpec {
5036 namespace: "rpc-identity-alias-test".to_string(),
5037 modules: Vec::new(),
5038 },
5039 pre_spawn: Vec::new(),
5040 })
5041 .timeout(Duration::from_secs(1))
5042 .build(),
5043 )
5044 .await?;
5045 for runtime_id in ["rt:review:singleton:0", "rt:review:singleton:1"] {
5046 let mut labels = BTreeMap::new();
5047 labels.insert("agent_identity".to_string(), "review:singleton".to_string());
5048 runtime
5049 .spawn(
5050 SpawnMemberSpec::from_wire(
5051 "worker".to_string(),
5052 runtime_id.to_string(),
5053 Some("You are a duplicate Review Agent.".into()),
5054 None,
5055 None,
5056 )
5057 .with_labels(labels),
5058 )
5059 .await?;
5060 }
5061 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5062 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5063 lease_provider: Arc::new(LocalLeaseProvider::new()),
5064 runtime_instance_id: "rpc-identity-alias-test".to_string(),
5065 has_runtime_store: true,
5066 durability_policy: DurabilityPolicy::SyncWriteThrough,
5067 bridge: None,
5068 default_timeout: None,
5069 });
5070 let identity = AgentIdentity::parse("review:singleton")?;
5071 let record = ContinuityRecord {
5072 identity: identity.clone(),
5073 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:1")?,
5074 session_id: meerkat_core::types::SessionId::new(),
5075 generation: ContinuityGeneration::new(1),
5076 checkpoint_version: CheckpointVersion::new(0),
5077 };
5078 identity_rt
5079 .register(
5080 DurableAgentSpec {
5081 identity,
5082 profile: meerkat_mob::ProfileName::from("worker"),
5083 addressability: AgentAddressability::Addressable,
5084 display_name: None,
5085 labels: BTreeMap::new(),
5086 context: None,
5087 additional_instructions: Vec::new(),
5088 initial_message: None,
5089 runtime_mode_override: None,
5090 backend: None,
5091 binding: None,
5092 },
5093 IdentityLifecycleState::Active,
5094 Some(record),
5095 None,
5096 )
5097 .await;
5098
5099 let target =
5100 resolve_rpc_identity_control_target(&runtime, &identity_rt, "review:singleton").await?;
5101 assert_eq!(target.identity.as_str(), "review:singleton");
5102 assert_eq!(
5103 target
5104 .live
5105 .as_ref()
5106 .map(|alias| alias.runtime_member_id.as_str()),
5107 Some("rt:review:singleton:1")
5108 );
5109
5110 let target =
5111 resolve_rpc_identity_control_target(&runtime, &identity_rt, "rt:review:singleton:1")
5112 .await?;
5113 assert_eq!(
5114 target
5115 .live
5116 .as_ref()
5117 .map(|alias| alias.runtime_member_id.as_str()),
5118 Some("rt:review:singleton:1")
5119 );
5120
5121 Ok(())
5122 }
5123
5124 #[tokio::test]
5125 async fn durable_resolution_rejects_hidden_registered_live_binding()
5126 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5127 let temp_dir = tempfile::tempdir()?;
5128 let runtime = Box::pin(
5129 UnifiedRuntime::builder()
5130 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5131 .module_config(MobKitConfig {
5132 modules: Vec::new(),
5133 discovery: DiscoverySpec {
5134 namespace: "rpc-hidden-bound-test".to_string(),
5135 modules: Vec::new(),
5136 },
5137 pre_spawn: Vec::new(),
5138 })
5139 .timeout(Duration::from_secs(1))
5140 .build(),
5141 )
5142 .await?;
5143 runtime
5144 .spawn(
5145 SpawnMemberSpec::from_wire(
5146 "worker".to_string(),
5147 "rt:review:singleton:0".to_string(),
5148 Some("You are a hidden Review Agent.".into()),
5149 None,
5150 None,
5151 )
5152 .with_labels(BTreeMap::from([
5153 ("agent_identity".to_string(), "review:singleton".to_string()),
5154 ("role".to_string(), "delegate".to_string()),
5155 ("source_mob_id".to_string(), "upstream".to_string()),
5156 ])),
5157 )
5158 .await?;
5159 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5160 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5161 lease_provider: Arc::new(LocalLeaseProvider::new()),
5162 runtime_instance_id: "rpc-hidden-bound-test".to_string(),
5163 has_runtime_store: true,
5164 durability_policy: DurabilityPolicy::SyncWriteThrough,
5165 bridge: None,
5166 default_timeout: None,
5167 });
5168 let identity = AgentIdentity::parse("review:singleton")?;
5169 identity_rt
5170 .register(
5171 DurableAgentSpec {
5172 identity: identity.clone(),
5173 profile: meerkat_mob::ProfileName::from("worker"),
5174 addressability: AgentAddressability::Addressable,
5175 display_name: None,
5176 labels: BTreeMap::new(),
5177 context: None,
5178 additional_instructions: Vec::new(),
5179 initial_message: None,
5180 runtime_mode_override: None,
5181 backend: None,
5182 binding: None,
5183 },
5184 IdentityLifecycleState::Active,
5185 Some(ContinuityRecord {
5186 identity,
5187 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
5188 session_id: meerkat_core::types::SessionId::new(),
5189 generation: ContinuityGeneration::new(0),
5190 checkpoint_version: CheckpointVersion::new(0),
5191 }),
5192 None,
5193 )
5194 .await;
5195
5196 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
5197 let err =
5198 resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
5199 .await
5200 .expect_err("hidden registered live binding must not resolve");
5201 assert!(
5202 err.contains("identity hidden by policy"),
5203 "unexpected error for {requested_identity}: {err}"
5204 );
5205 }
5206
5207 Ok(())
5208 }
5209
5210 #[tokio::test]
5211 async fn live_only_hidden_alias_reports_policy_error()
5212 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5213 let temp_dir = tempfile::tempdir()?;
5214 let runtime = Box::pin(
5215 UnifiedRuntime::builder()
5216 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5217 .module_config(MobKitConfig {
5218 modules: Vec::new(),
5219 discovery: DiscoverySpec {
5220 namespace: "rpc-hidden-live-only-test".to_string(),
5221 modules: Vec::new(),
5222 },
5223 pre_spawn: Vec::new(),
5224 })
5225 .timeout(Duration::from_secs(1))
5226 .build(),
5227 )
5228 .await?;
5229 runtime
5230 .spawn(
5231 SpawnMemberSpec::from_wire(
5232 "worker".to_string(),
5233 "rt:review:singleton:0".to_string(),
5234 Some("You are a hidden Review Agent.".into()),
5235 None,
5236 None,
5237 )
5238 .with_labels(BTreeMap::from([
5239 ("agent_identity".to_string(), "review:singleton".to_string()),
5240 ("role".to_string(), "delegate".to_string()),
5241 ("source_mob_id".to_string(), "upstream".to_string()),
5242 ])),
5243 )
5244 .await?;
5245 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5246 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5247 lease_provider: Arc::new(LocalLeaseProvider::new()),
5248 runtime_instance_id: "rpc-hidden-live-only-test".to_string(),
5249 has_runtime_store: true,
5250 durability_policy: DurabilityPolicy::SyncWriteThrough,
5251 bridge: None,
5252 default_timeout: None,
5253 });
5254
5255 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
5256 let err =
5257 resolve_rpc_identity_control_target(&runtime, &identity_rt, requested_identity)
5258 .await
5259 .expect_err("hidden live-only alias must not collapse into unknown identity");
5260 assert!(
5261 err.contains("identity hidden by policy"),
5262 "unexpected error for {requested_identity}: {err}"
5263 );
5264 }
5265
5266 let identity_ctx = IdentityFirstContext {
5267 runtime: Arc::new(identity_rt),
5268 roster_provider: Arc::new(EmptyRosterProvider),
5269 topology_provider: None,
5270 customizer: None,
5271 };
5272 for requested_identity in ["review:singleton", "rt:review:singleton:0"] {
5273 let response: Value = serde_json::from_str(
5274 &handle_unified_rpc_json(
5275 &runtime,
5276 &json!({
5277 "jsonrpc": "2.0",
5278 "id": 1,
5279 "method": "mobkit/status_identity",
5280 "params": { "identity": requested_identity },
5281 })
5282 .to_string(),
5283 Duration::from_secs(1),
5284 None,
5285 Some(&identity_ctx),
5286 )
5287 .await,
5288 )?;
5289 assert_eq!(
5290 response["error"]["data"]["kind"],
5291 json!("identity_hidden_by_policy"),
5292 "unexpected hidden response for {requested_identity}: {response:#?}"
5293 );
5294 }
5295
5296 Ok(())
5297 }
5298
5299 #[tokio::test]
5300 async fn live_only_resolution_rejects_runtime_member_bound_to_other_durable_identity()
5301 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5302 let temp_dir = tempfile::tempdir()?;
5303 let runtime = Box::pin(
5304 UnifiedRuntime::builder()
5305 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5306 .module_config(MobKitConfig {
5307 modules: Vec::new(),
5308 discovery: DiscoverySpec {
5309 namespace: "rpc-identity-alias-test".to_string(),
5310 modules: Vec::new(),
5311 },
5312 pre_spawn: Vec::new(),
5313 })
5314 .timeout(Duration::from_secs(1))
5315 .build(),
5316 )
5317 .await?;
5318 let mut labels = BTreeMap::new();
5319 labels.insert("agent_identity".to_string(), "other:singleton".to_string());
5320 runtime
5321 .spawn(
5322 SpawnMemberSpec::from_wire(
5323 "worker".to_string(),
5324 "rt:review:singleton:0".to_string(),
5325 Some("You are a wrong-projected Review Agent.".into()),
5326 None,
5327 None,
5328 )
5329 .with_labels(labels),
5330 )
5331 .await?;
5332
5333 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5334 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5335 lease_provider: Arc::new(LocalLeaseProvider::new()),
5336 runtime_instance_id: "rpc-identity-alias-test".to_string(),
5337 has_runtime_store: true,
5338 durability_policy: DurabilityPolicy::SyncWriteThrough,
5339 bridge: None,
5340 default_timeout: None,
5341 });
5342 let identity = AgentIdentity::parse("review:singleton")?;
5343 let record = ContinuityRecord {
5344 identity: identity.clone(),
5345 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
5346 session_id: meerkat_core::types::SessionId::new(),
5347 generation: ContinuityGeneration::new(0),
5348 checkpoint_version: CheckpointVersion::new(0),
5349 };
5350 identity_rt
5351 .register(
5352 DurableAgentSpec {
5353 identity: identity.clone(),
5354 profile: meerkat_mob::ProfileName::from("worker"),
5355 addressability: AgentAddressability::Addressable,
5356 display_name: None,
5357 labels: BTreeMap::new(),
5358 context: None,
5359 additional_instructions: Vec::new(),
5360 initial_message: None,
5361 runtime_mode_override: None,
5362 backend: None,
5363 binding: None,
5364 },
5365 IdentityLifecycleState::Active,
5366 Some(record),
5367 Some(LeaseGrant {
5368 identity,
5369 fencing_token: FencingToken::new(1),
5370 ttl: Duration::from_mins(1),
5371 }),
5372 )
5373 .await;
5374
5375 let err = resolve_rpc_identity_control_target(&runtime, &identity_rt, "other:singleton")
5376 .await
5377 .expect_err("wrong-projected live alias must not resolve as live-only");
5378 assert!(
5379 err.contains("identity runtime binding belongs to review:singleton"),
5380 "unexpected error: {err}"
5381 );
5382
5383 Ok(())
5384 }
5385
5386 #[tokio::test]
5395 async fn send_message_resolves_bare_durable_identity_through_identity_bridge()
5396 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5397 let temp_dir = tempfile::tempdir()?;
5398 let runtime = Box::pin(
5399 UnifiedRuntime::builder()
5400 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5401 .module_config(MobKitConfig {
5402 modules: Vec::new(),
5403 discovery: DiscoverySpec {
5404 namespace: "rpc-send-message-identity-bridge-test".to_string(),
5405 modules: Vec::new(),
5406 },
5407 pre_spawn: Vec::new(),
5408 })
5409 .timeout(Duration::from_secs(5))
5410 .build(),
5411 )
5412 .await?;
5413
5414 for (runtime_id, durable) in [
5417 ("rt:atlas-base-001:0", "atlas-base-001"),
5418 ("rt:draco-base-001:0", "draco-base-001"),
5419 ] {
5420 runtime
5421 .spawn(
5422 SpawnMemberSpec::from_wire(
5423 "worker".to_string(),
5424 runtime_id.to_string(),
5425 Some("You are a swarm base agent.".into()),
5426 None,
5427 None,
5428 )
5429 .with_labels(BTreeMap::from([(
5430 "agent_identity".to_string(),
5431 durable.to_string(),
5432 )])),
5433 )
5434 .await?;
5435 }
5436 runtime
5439 .spawn(SpawnMemberSpec::from_wire(
5440 "worker".to_string(),
5441 "draco-base-001".to_string(),
5442 Some("You are the raw roster member.".into()),
5443 None,
5444 None,
5445 ))
5446 .await?;
5447
5448 let session_service = runtime
5449 .mob_runtime()
5450 .session_service()
5451 .cloned()
5452 .expect("test mob spec has a session service");
5453 let bridge: Arc<dyn crate::identity_first::SessionBridge> = Arc::new(
5454 crate::identity_first::MobSessionBridge::with_session_service(
5455 runtime.mob_handle(),
5456 session_service,
5457 ),
5458 );
5459 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5460 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5461 lease_provider: Arc::new(LocalLeaseProvider::new()),
5462 runtime_instance_id: "rpc-send-message-identity-bridge-test".to_string(),
5463 has_runtime_store: true,
5464 durability_policy: DurabilityPolicy::SyncWriteThrough,
5465 bridge: Some(bridge),
5466 default_timeout: None,
5467 })
5468 .with_runtime_services(crate::identity_first::AgentRuntimeServices::new(
5469 runtime.mob_handle(),
5470 ));
5471 for (durable, runtime_id) in [
5472 ("atlas-base-001", "rt:atlas-base-001:0"),
5473 ("draco-base-001", "rt:draco-base-001:0"),
5474 ] {
5475 let identity = AgentIdentity::parse(durable)?;
5476 let session_id = runtime
5477 .mob_handle()
5478 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(runtime_id))
5479 .await
5480 .unwrap_or_else(meerkat_core::types::SessionId::new);
5481 identity_rt
5482 .register(
5483 DurableAgentSpec {
5484 identity: identity.clone(),
5485 profile: meerkat_mob::ProfileName::from("worker"),
5486 addressability: AgentAddressability::Addressable,
5487 display_name: None,
5488 labels: BTreeMap::new(),
5489 context: None,
5490 additional_instructions: Vec::new(),
5491 initial_message: None,
5492 runtime_mode_override: None,
5493 backend: None,
5494 binding: None,
5495 },
5496 IdentityLifecycleState::Active,
5497 Some(ContinuityRecord {
5498 identity: identity.clone(),
5499 agent_runtime_id: AgentRuntimeId::parse(runtime_id)?,
5500 session_id,
5501 generation: ContinuityGeneration::new(0),
5502 checkpoint_version: CheckpointVersion::new(0),
5503 }),
5504 Some(LeaseGrant {
5505 identity,
5506 fencing_token: FencingToken::new(1),
5507 ttl: Duration::from_mins(1),
5508 }),
5509 )
5510 .await;
5511 }
5512 let identity_ctx = IdentityFirstContext {
5513 runtime: Arc::new(identity_rt),
5514 roster_provider: Arc::new(EmptyRosterProvider),
5515 topology_provider: None,
5516 customizer: None,
5517 };
5518
5519 let send = |id: u64, params: Value| {
5520 let runtime = &runtime;
5521 let identity_ctx = &identity_ctx;
5522 async move {
5523 let raw = handle_unified_rpc_json(
5524 runtime,
5525 &json!({
5526 "jsonrpc": "2.0",
5527 "id": id,
5528 "method": "mobkit/send_message",
5529 "params": params,
5530 })
5531 .to_string(),
5532 Duration::from_secs(10),
5533 None,
5534 Some(identity_ctx),
5535 )
5536 .await;
5537 serde_json::from_str::<Value>(&raw)
5538 }
5539 };
5540
5541 let response = send(
5544 1,
5545 json!({ "member_id": "atlas-base-001", "message": "status check" }),
5546 )
5547 .await?;
5548 assert!(
5549 response["error"].is_null(),
5550 "bare identity send must bridge-resolve: {response:#?}"
5551 );
5552 assert_eq!(response["result"]["accepted"], json!(true));
5553 assert_eq!(response["result"]["member_id"], json!("atlas-base-001"));
5554 let atlas_session = runtime
5555 .mob_handle()
5556 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
5557 "rt:atlas-base-001:0",
5558 ))
5559 .await
5560 .expect("atlas runtime member has a bridge session after send")
5561 .to_string();
5562 assert_eq!(response["result"]["session_id"], json!(atlas_session));
5563
5564 let response = send(
5566 2,
5567 json!({
5568 "member_id": "atlas-base-001",
5569 "message": "steer: stand down",
5570 "handling_mode": "steer",
5571 }),
5572 )
5573 .await?;
5574 assert!(
5575 response["error"].is_null(),
5576 "bare identity steer must bridge-resolve: {response:#?}"
5577 );
5578 assert_eq!(response["result"]["accepted"], json!(true));
5579
5580 let response = send(
5584 3,
5585 json!({ "member_id": "draco-base-001", "message": "raw roster delivery" }),
5586 )
5587 .await?;
5588 assert!(
5589 response["error"].is_null(),
5590 "exact roster member send must keep raw semantics: {response:#?}"
5591 );
5592 assert_eq!(response["result"]["accepted"], json!(true));
5593 let draco_raw_session = runtime
5594 .mob_handle()
5595 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id("draco-base-001"))
5596 .await
5597 .expect("bare draco member has a bridge session after send")
5598 .to_string();
5599 assert_eq!(response["result"]["session_id"], json!(draco_raw_session));
5600 if let Some(draco_rt_session) = runtime
5601 .mob_handle()
5602 .resolve_bridge_session_id(&crate::member_comms_id::mob_member_id(
5603 "rt:draco-base-001:0",
5604 ))
5605 .await
5606 {
5607 assert_ne!(
5608 response["result"]["session_id"],
5609 json!(draco_rt_session.to_string()),
5610 "exact member id match must not be shadowed by identity resolution"
5611 );
5612 }
5613
5614 let response = send(
5616 4,
5617 json!({ "member_id": "phantom-base-999", "message": "nobody home" }),
5618 )
5619 .await?;
5620 assert_eq!(response["error"]["code"], json!(-32000), "{response:#?}");
5621 let message = response["error"]["message"]
5622 .as_str()
5623 .expect("error message");
5624 assert!(
5625 message.starts_with("send_message failed:"),
5626 "unexpected error message: {message}"
5627 );
5628
5629 Ok(())
5630 }
5631
5632 #[tokio::test]
5639 async fn member_state_wire_vocabulary_is_lowercase_on_both_surfaces()
5640 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5641 let temp_dir = tempfile::tempdir()?;
5642 let runtime = Box::pin(
5643 UnifiedRuntime::builder()
5644 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5645 .module_config(MobKitConfig {
5646 modules: Vec::new(),
5647 discovery: DiscoverySpec {
5648 namespace: "rpc-state-vocabulary-test".to_string(),
5649 modules: Vec::new(),
5650 },
5651 pre_spawn: Vec::new(),
5652 })
5653 .timeout(Duration::from_secs(5))
5654 .build(),
5655 )
5656 .await?;
5657 runtime
5658 .spawn(SpawnMemberSpec::from_wire(
5659 "worker".to_string(),
5660 "worker-one".to_string(),
5661 None,
5662 None,
5663 None,
5664 ))
5665 .await?;
5666
5667 let raw = handle_unified_rpc_json(
5669 &runtime,
5670 &json!({
5671 "jsonrpc": "2.0",
5672 "id": 1,
5673 "method": "mobkit/get_member",
5674 "params": { "member_id": "worker-one" },
5675 })
5676 .to_string(),
5677 Duration::from_secs(5),
5678 None,
5679 None,
5680 )
5681 .await;
5682 let response: Value = serde_json::from_str(&raw)?;
5683 assert_eq!(
5684 response["result"]["state"],
5685 json!("active"),
5686 "member rows must speak the lowercase SDK vocabulary: {response:#?}"
5687 );
5688
5689 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5691 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5692 lease_provider: Arc::new(LocalLeaseProvider::new()),
5693 runtime_instance_id: "rpc-state-vocabulary-test".to_string(),
5694 has_runtime_store: true,
5695 durability_policy: DurabilityPolicy::SyncWriteThrough,
5696 bridge: None,
5697 default_timeout: None,
5698 });
5699 let identity = AgentIdentity::parse("review:singleton")?;
5700 identity_rt
5701 .register(
5702 DurableAgentSpec {
5703 identity: identity.clone(),
5704 profile: meerkat_mob::ProfileName::from("worker"),
5705 addressability: AgentAddressability::Addressable,
5706 display_name: None,
5707 labels: BTreeMap::new(),
5708 context: None,
5709 additional_instructions: Vec::new(),
5710 initial_message: None,
5711 runtime_mode_override: None,
5712 backend: None,
5713 binding: None,
5714 },
5715 IdentityLifecycleState::Active,
5716 Some(ContinuityRecord {
5717 identity: identity.clone(),
5718 agent_runtime_id: AgentRuntimeId::parse("rt:review:singleton:0")?,
5719 session_id: meerkat_core::types::SessionId::new(),
5720 generation: ContinuityGeneration::new(0),
5721 checkpoint_version: CheckpointVersion::new(0),
5722 }),
5723 Some(LeaseGrant {
5724 identity,
5725 fencing_token: FencingToken::new(1),
5726 ttl: Duration::from_mins(1),
5727 }),
5728 )
5729 .await;
5730 let identity_ctx = IdentityFirstContext {
5731 runtime: Arc::new(identity_rt),
5732 roster_provider: Arc::new(EmptyRosterProvider),
5733 topology_provider: None,
5734 customizer: None,
5735 };
5736 let raw = handle_unified_rpc_json(
5737 &runtime,
5738 &json!({
5739 "jsonrpc": "2.0",
5740 "id": 2,
5741 "method": "mobkit/status_identity",
5742 "params": { "identity": "review:singleton" },
5743 })
5744 .to_string(),
5745 Duration::from_secs(5),
5746 None,
5747 Some(&identity_ctx),
5748 )
5749 .await;
5750 let response: Value = serde_json::from_str(&raw)?;
5751 assert_eq!(
5752 response["result"]["state"],
5753 json!("active"),
5754 "identity status must speak the same lowercase vocabulary: {response:#?}"
5755 );
5756
5757 Ok(())
5758 }
5759
5760 #[tokio::test]
5769 async fn send_message_pins_baseline_member_over_identity_fallback()
5770 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5771 let temp_dir = tempfile::tempdir()?;
5772 let runtime = Box::pin(
5773 UnifiedRuntime::builder()
5774 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5775 .module_config(MobKitConfig {
5776 modules: Vec::new(),
5777 discovery: DiscoverySpec {
5778 namespace: "rpc-send-message-baseline-pin-test".to_string(),
5779 modules: Vec::new(),
5780 },
5781 pre_spawn: Vec::new(),
5782 })
5783 .timeout(Duration::from_secs(5))
5784 .build(),
5785 )
5786 .await?;
5787
5788 for member in ["rt:draco-base-001:0", "draco-base-001"] {
5791 runtime
5792 .spawn(SpawnMemberSpec::from_wire(
5793 "worker".to_string(),
5794 member.to_string(),
5795 None,
5796 None,
5797 None,
5798 ))
5799 .await?;
5800 }
5801
5802 let identity_rt = IdentityRuntime::new(IdentityRuntimeConfig {
5803 continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5804 lease_provider: Arc::new(LocalLeaseProvider::new()),
5805 runtime_instance_id: "rpc-send-message-baseline-pin-test".to_string(),
5806 has_runtime_store: true,
5807 durability_policy: DurabilityPolicy::SyncWriteThrough,
5808 bridge: None,
5809 default_timeout: None,
5810 });
5811 let identity = AgentIdentity::parse("draco-base-001")?;
5812 identity_rt
5813 .register(
5814 DurableAgentSpec {
5815 identity: identity.clone(),
5816 profile: meerkat_mob::ProfileName::from("worker"),
5817 addressability: AgentAddressability::Addressable,
5818 display_name: None,
5819 labels: BTreeMap::new(),
5820 context: None,
5821 additional_instructions: Vec::new(),
5822 initial_message: None,
5823 runtime_mode_override: None,
5824 backend: None,
5825 binding: None,
5826 },
5827 IdentityLifecycleState::Active,
5828 Some(ContinuityRecord {
5829 identity: identity.clone(),
5830 agent_runtime_id: AgentRuntimeId::parse("rt:draco-base-001:0")?,
5831 session_id: meerkat_core::types::SessionId::new(),
5832 generation: ContinuityGeneration::new(0),
5833 checkpoint_version: CheckpointVersion::new(0),
5834 }),
5835 Some(LeaseGrant {
5836 identity,
5837 fencing_token: FencingToken::new(1),
5838 ttl: Duration::from_mins(1),
5839 }),
5840 )
5841 .await;
5842 let identity_ctx = IdentityFirstContext {
5843 runtime: Arc::new(identity_rt),
5844 roster_provider: Arc::new(EmptyRosterProvider),
5845 topology_provider: None,
5846 customizer: None,
5847 };
5848
5849 runtime
5851 .mob_runtime()
5852 .set_baseline_member_specs(vec![SpawnMemberSpec::new(
5853 meerkat_mob::ProfileName::from("worker"),
5854 meerkat_mob::AgentIdentity::from("draco-base-001"),
5855 )])
5856 .await;
5857 runtime
5860 .mob_handle()
5861 .retire(meerkat_mob::AgentIdentity::from("draco-base-001"))
5862 .await?;
5863
5864 let raw = handle_unified_rpc_json(
5865 &runtime,
5866 &json!({
5867 "jsonrpc": "2.0",
5868 "id": 1,
5869 "method": "mobkit/send_message",
5870 "params": { "member_id": "draco-base-001", "message": "mid-reconcile send" },
5871 })
5872 .to_string(),
5873 Duration::from_secs(10),
5874 None,
5875 Some(&identity_ctx),
5876 )
5877 .await;
5878 let response: Value = serde_json::from_str(&raw)?;
5879 assert_eq!(
5880 response["error"]["code"],
5881 json!(-32000),
5882 "transiently-absent baseline member must keep raw member-id semantics \
5883 instead of silently delivering through the identity bridge: {response:#?}"
5884 );
5885 let message = response["error"]["message"]
5886 .as_str()
5887 .expect("error message");
5888 assert!(
5889 message.starts_with("send_message failed:"),
5890 "unexpected error message: {message}"
5891 );
5892
5893 Ok(())
5894 }
5895
5896 #[tokio::test]
5905 async fn retire_and_respawn_rpcs_succeed_for_idle_member()
5906 -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
5907 let temp_dir = tempfile::tempdir()?;
5908 let runtime = Box::pin(
5909 UnifiedRuntime::builder()
5910 .mob_spec(rpc_test_mob_spec(&temp_dir)?)
5911 .module_config(MobKitConfig {
5912 modules: Vec::new(),
5913 discovery: DiscoverySpec {
5914 namespace: "rpc-idle-lifecycle-test".to_string(),
5915 modules: Vec::new(),
5916 },
5917 pre_spawn: Vec::new(),
5918 })
5919 .timeout(Duration::from_secs(5))
5920 .build(),
5921 )
5922 .await?;
5923
5924 for member in ["worker-one", "worker-two"] {
5925 runtime
5926 .spawn(SpawnMemberSpec::from_wire(
5927 "worker".to_string(),
5928 member.to_string(),
5929 None,
5930 None,
5931 None,
5932 ))
5933 .await?;
5934 }
5935
5936 let send = |id: u64, method: &'static str, params: Value| {
5937 let runtime = &runtime;
5938 async move {
5939 let raw = handle_unified_rpc_json(
5940 runtime,
5941 &json!({
5942 "jsonrpc": "2.0",
5943 "id": id,
5944 "method": method,
5945 "params": params,
5946 })
5947 .to_string(),
5948 Duration::from_secs(10),
5949 None,
5950 None,
5951 )
5952 .await;
5953 serde_json::from_str::<Value>(&raw)
5954 }
5955 };
5956
5957 let response = send(
5959 1,
5960 "mobkit/retire_member",
5961 json!({"member_id": "worker-one"}),
5962 )
5963 .await?;
5964 assert!(
5965 response["error"].is_null(),
5966 "retire_member must succeed for an idle member: {response:#?}"
5967 );
5968 assert_eq!(response["result"]["accepted"], json!(true));
5969 assert!(
5970 !runtime
5971 .mob_handle()
5972 .list_members_including_retiring()
5973 .await
5974 .iter()
5975 .any(|entry| entry.agent_identity.as_str() == "worker-one"),
5976 "retired member must leave the roster"
5977 );
5978
5979 let response = send(
5982 2,
5983 "mobkit/respawn_member",
5984 json!({"member_id": "worker-two"}),
5985 )
5986 .await?;
5987 assert!(
5988 response["error"].is_null(),
5989 "respawn_member must succeed for an idle member: {response:#?}"
5990 );
5991 assert_eq!(response["result"]["accepted"], json!(true));
5992 let members = runtime.mob_handle().list_members_including_retiring().await;
5993 let worker_two = members
5994 .iter()
5995 .find(|entry| entry.agent_identity.as_str() == "worker-two")
5996 .expect("respawned member must remain in the roster");
5997 assert_eq!(
5998 worker_two.status,
5999 meerkat_mob::MobMemberStatus::Active,
6000 "respawned member must be active, not wedged in retiring"
6001 );
6002
6003 Ok(())
6004 }
6005}