use super::connection::WireAuthBindingRef;
use super::runtime::WireTurnMetadataOverride;
use super::session::WireContentInput;
use super::supervisor_bridge::BridgeBootstrapToken;
use base64::{Engine, engine::general_purpose::STANDARD as BASE64};
use meerkat_core::OutputSchema;
use meerkat_core::{
HandlingMode,
types::{RenderClass, RenderMetadata, RenderSalience},
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeMap;
use meerkat_core::{SurfaceMetadata, SurfaceMetadataError};
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobBackendKind {
#[default]
Session,
External,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
pub enum WireRuntimeBinding {
Session,
External {
address: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
bootstrap_token: Option<BridgeBootstrapToken>,
identity: WireTrustedPeerIdentity,
},
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobRuntimeMode {
#[default]
AutonomousHost,
TurnDriven,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "mode", rename_all = "snake_case")]
pub enum WireMemberLaunchMode {
Fresh,
Resume {
bridge_session_id: String,
},
Fork {
source_member_id: String,
#[serde(default)]
fork_context: WireForkContext,
},
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum WireForkContext {
#[default]
FullHistory,
LastMessages {
count: u32,
},
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "type", content = "value", rename_all = "snake_case")]
pub enum WireToolAccessPolicy {
#[default]
Inherit,
AllowList(Vec<String>),
DenyList(Vec<String>),
}
impl WireToolAccessPolicy {
#[must_use]
pub fn into_core(self) -> meerkat_core::ops::ToolAccessPolicy {
match self {
Self::Inherit => meerkat_core::ops::ToolAccessPolicy::Inherit,
Self::AllowList(names) => {
meerkat_core::ops::ToolAccessPolicy::AllowList(names.into_iter().collect())
}
Self::DenyList(names) => {
meerkat_core::ops::ToolAccessPolicy::DenyList(names.into_iter().collect())
}
}
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub enum WireToolFilter {
#[default]
All,
Allow(Vec<String>),
Deny(Vec<String>),
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireMobToolConfig {
#[serde(default)]
pub builtins: bool,
#[serde(default)]
pub shell: bool,
#[serde(default)]
pub comms: bool,
#[serde(default)]
pub memory: bool,
#[serde(default)]
pub workgraph: bool,
#[serde(default)]
pub mob: bool,
#[serde(default)]
pub schedule: bool,
#[serde(default)]
pub image_generation: bool,
#[serde(default)]
pub mcp: Vec<String>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobResumeOverrideField {
Model,
Provider,
ProviderParams,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireMobProfile {
pub model: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider: Option<meerkat_core::Provider>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub self_hosted_server_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub image_generation_provider: Option<meerkat_core::Provider>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auto_compact_threshold: Option<std::num::NonZeroU64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub resume_overrides: Vec<WireMobResumeOverrideField>,
#[serde(default)]
pub skills: Vec<String>,
#[serde(default)]
pub tools: WireMobToolConfig,
#[serde(default)]
pub peer_description: String,
#[serde(default)]
pub external_addressable: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
#[serde(default)]
pub runtime_mode: WireMobRuntimeMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_inline_peer_notifications: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_schema: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider_params: Option<crate::wire::runtime::WireProviderParamsOverride>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobOrchestratorInput {
pub profile: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "source", rename_all = "snake_case")]
pub enum MobSkillSourceInput {
Inline { content: String },
Path { path: String },
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRoleWiringRuleInput {
pub a: String,
pub b: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobWiringRulesInput {
#[serde(default)]
pub auto_wire_orchestrator: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub role_wiring: Vec<MobRoleWiringRuleInput>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobToolConfigInput {
#[serde(default)]
pub builtins: bool,
#[serde(default)]
pub shell: bool,
#[serde(default)]
pub comms: bool,
#[serde(default)]
pub memory: bool,
#[serde(default)]
pub workgraph: bool,
#[serde(default)]
pub mob: bool,
#[serde(default)]
pub schedule: bool,
#[serde(default)]
pub image_generation: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub mcp: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[allow(clippy::large_enum_variant)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(untagged)]
pub enum MobProfileBindingInput {
RealmRef {
realm_profile: String,
},
Inline(MobProfileInput),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobProfileInput {
pub model: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider: Option<meerkat_core::Provider>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub self_hosted_server_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub image_generation_provider: Option<meerkat_core::Provider>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auto_compact_threshold: Option<std::num::NonZeroU64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub resume_overrides: Vec<WireMobResumeOverrideField>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub skills: Vec<String>,
#[serde(default)]
pub tools: MobToolConfigInput,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub peer_description: String,
#[serde(default)]
pub external_addressable: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
#[serde(default)]
pub runtime_mode: WireMobRuntimeMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_inline_peer_notifications: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_schema: Option<OutputSchema>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider_params: Option<crate::wire::runtime::WireProviderParamsOverride>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobExternalBackendConfigInput {
pub address_base: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub supervisor_bridge: Option<MobSupervisorBridgeEndpointConfigInput>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSupervisorBridgeEndpointConfigInput {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub bind_address: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub advertised_address: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobBackendConfigInput {
#[serde(default)]
pub default: WireMobBackendKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub external: Option<MobExternalBackendConfigInput>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobDispatchModeInput {
#[default]
FanOut,
OneToOne,
FanIn,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum MobCollectionPolicyInput {
#[default]
All,
Any,
Quorum {
n: u8,
},
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobDependencyModeInput {
#[default]
All,
Any,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobStepOutputFormatInput {
Json,
Text,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "op", rename_all = "snake_case")]
pub enum MobConditionExprInput {
Eq { path: String, value: Value },
In { path: String, values: Vec<Value> },
Gt { path: String, value: Value },
Lt { path: String, value: Value },
And { exprs: Vec<MobConditionExprInput> },
Or { exprs: Vec<MobConditionExprInput> },
Not { expr: Box<MobConditionExprInput> },
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFrameSpecInput {
pub nodes: BTreeMap<String, MobFlowNodeInput>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum MobFlowNodeInput {
Step(MobFrameStepInput),
RepeatUntil(MobRepeatUntilInput),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFrameStepInput {
pub step_id: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default)]
pub depends_on_mode: MobDependencyModeInput,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRepeatUntilInput {
pub loop_id: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default)]
pub depends_on_mode: MobDependencyModeInput,
pub body: MobFrameSpecInput,
pub until: MobConditionExprInput,
pub max_iterations: u32,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFlowStepInput {
pub role: String,
pub message: WireContentInput,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default)]
pub dispatch_mode: MobDispatchModeInput,
#[serde(default)]
pub collection_policy: MobCollectionPolicyInput,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub condition: Option<MobConditionExprInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expected_schema_ref: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
#[serde(default)]
pub depends_on_mode: MobDependencyModeInput,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub allowed_tools: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub blocked_tools: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_format: Option<MobStepOutputFormatInput>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFlowSpecInput {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub steps: BTreeMap<String, MobFlowStepInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub root: Option<MobFrameSpecInput>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobPolicyModeInput {
#[default]
Advisory,
Strict,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobTopologyRuleInput {
pub from_role: String,
pub to_role: String,
pub allowed: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobTopologySpecInput {
pub mode: MobPolicyModeInput,
pub rules: Vec<MobTopologyRuleInput>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSupervisorSpecInput {
pub role: String,
pub escalation_threshold: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub escalation_turn_timeout_ms: Option<u64>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobLimitsSpecInput {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_flow_duration_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_step_retries: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_orphaned_turns: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cancel_grace_timeout_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_active_nodes: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_active_frames: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_frame_depth: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "mode", rename_all = "snake_case")]
pub enum MobSpawnPolicyInput {
None,
Auto {
profile_map: BTreeMap<String, String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobEventRouterConfigInput {
#[serde(default = "default_event_router_buffer_size")]
pub buffer_size: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub include_patterns: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exclude_patterns: Option<Vec<String>>,
}
const fn default_event_router_buffer_size() -> usize {
256
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobDefinitionInput {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub orchestrator: Option<MobOrchestratorInput>,
pub profiles: BTreeMap<String, MobProfileBindingInput>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub models: BTreeMap<String, meerkat_core::config::CustomModelConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub image_generation_provider: Option<meerkat_core::Provider>,
#[serde(default)]
pub wiring: MobWiringRulesInput,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub skills: BTreeMap<String, MobSkillSourceInput>,
#[serde(default)]
pub backend: MobBackendConfigInput,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub flows: BTreeMap<String, MobFlowSpecInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub topology: Option<MobTopologySpecInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub supervisor: Option<MobSupervisorSpecInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub limits: Option<MobLimitsSpecInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub spawn_policy: Option<MobSpawnPolicyInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub event_router: Option<MobEventRouterConfigInput>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobCreateParams {
pub definition: MobDefinitionInput,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobCreateResult {
pub mob_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobIdParams {
pub mob_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberParams {
pub mob_id: String,
pub agent_identity: String,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub enum WireMobLifecycleStatus {
Creating,
Running,
Stopped,
Completed,
Destroyed,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobStatusResult {
pub mob_id: String,
pub status: WireMobLifecycleStatus,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobListResult {
pub mobs: Vec<MobStatusResult>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnParams {
pub mob_id: String,
pub profile: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub initial_message: Option<WireContentInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub labels: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub binding: Option<WireRuntimeBinding>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub shell_env: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auto_wire_parent: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub launch_mode: Option<WireMemberLaunchMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tool_access_policy: Option<WireToolAccessPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub inherited_tool_filter: Option<WireToolFilter>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub override_profile: Option<WireMobProfile>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_override: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub placement: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobSpawnResult {
pub mob_id: String,
pub agent_identity: String,
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnSpecParams {
pub profile: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub initial_message: Option<WireContentInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub labels: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub placement: Option<WireHostRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_override: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnManyParams {
pub mob_id: String,
pub specs: Vec<MobSpawnSpecParams>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobSpawnManyResultStatus {
Spawned,
Failed,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnManySpawnedResult {
pub agent_identity: String,
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobSpawnManyFailureCause {
ProfileNotFound,
MemberNotFound,
MemberAlreadyExists,
NotExternallyAddressable,
InvalidTransition,
WiringError,
BridgeCommandRejected,
MemberRestoreFailed,
KickoffWaitTimedOut,
ReadyWaitTimedOut,
DefinitionError,
FlowNotFound,
FlowFailed,
RunNotFound,
RunCanceled,
FlowTurnTimedOut,
FrameDepthLimitExceeded,
FrameAtomicPersistenceUnavailable,
SpecRevisionConflict,
SchemaValidation,
InsufficientTargets,
TopologyViolation,
BridgeDeliveryRejected,
SupervisorEscalation,
UnsupportedForMode,
MissingMemberCapability,
ResetBarrier,
StorageError,
SessionError,
CommsError,
CallbackPending,
StaleFenceToken,
StaleEventCursor,
WorkNotFound,
Internal,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnManyFailedResult {
pub cause: MobSpawnManyFailureCause,
pub message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(untagged)]
pub enum MobSpawnManyResultPayload {
Spawned(MobSpawnManySpawnedResult),
Failed(MobSpawnManyFailedResult),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(try_from = "MobSpawnManyResultEntryRaw")]
pub struct MobSpawnManyResultEntry {
pub status: MobSpawnManyResultStatus,
pub result: MobSpawnManyResultPayload,
}
#[derive(Debug, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
struct MobSpawnManyResultEntryRaw {
status: MobSpawnManyResultStatus,
result: MobSpawnManyResultPayload,
}
impl TryFrom<MobSpawnManyResultEntryRaw> for MobSpawnManyResultEntry {
type Error = String;
fn try_from(raw: MobSpawnManyResultEntryRaw) -> Result<Self, Self::Error> {
let entry = Self {
status: raw.status,
result: raw.result,
};
entry.validate().map_err(str::to_owned)?;
Ok(entry)
}
}
impl MobSpawnManyResultEntry {
pub fn spawned(agent_identity: impl Into<String>, member_ref: WireMemberRef) -> Self {
Self {
status: MobSpawnManyResultStatus::Spawned,
result: MobSpawnManyResultPayload::Spawned(MobSpawnManySpawnedResult {
agent_identity: agent_identity.into(),
member_ref,
}),
}
}
pub fn failed(cause: MobSpawnManyFailureCause, message: impl Into<String>) -> Self {
Self {
status: MobSpawnManyResultStatus::Failed,
result: MobSpawnManyResultPayload::Failed(MobSpawnManyFailedResult {
cause,
message: message.into(),
}),
}
}
pub fn validate(&self) -> Result<(), &'static str> {
match (&self.status, &self.result) {
(MobSpawnManyResultStatus::Spawned, MobSpawnManyResultPayload::Spawned(_))
| (MobSpawnManyResultStatus::Failed, MobSpawnManyResultPayload::Failed(_)) => Ok(()),
(MobSpawnManyResultStatus::Spawned, MobSpawnManyResultPayload::Failed(_)) => {
Err("mob spawn_many result status spawned requires spawned result")
}
(MobSpawnManyResultStatus::Failed, MobSpawnManyResultPayload::Spawned(_)) => {
Err("mob spawn_many result status failed requires failed result")
}
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobSpawnManyResult {
pub results: Vec<MobSpawnManyResultEntry>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRetireResult {
pub retired: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRespawnParams {
pub mob_id: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub initial_message: Option<WireContentInput>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRespawnReceipt {
pub identity: String,
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobRespawnOutcome {
Completed,
TopologyRestoreFailed,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRespawnResult {
pub status: WireMobRespawnOutcome,
pub receipt: MobRespawnReceipt,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub failed_peer_ids: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMembersResult {
pub mob_id: String,
pub members: Vec<MobMemberListEntryWire>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobEventsParams {
pub mob_id: String,
#[serde(default)]
pub after_cursor: u64,
#[serde(default = "default_mob_events_limit")]
pub limit: usize,
#[serde(default)]
pub strict: bool,
}
const fn default_mob_events_limit() -> usize {
100
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobEventsResult {
pub events: Vec<Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
pub enum WireTrustedPeerIdentity {
Ed25519PublicKey { public_key: String },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ResolvedWireTrustedPeerIdentity {
pub peer_id: meerkat_core::comms::PeerId,
pub pubkey: [u8; 32],
}
#[derive(Debug, Clone, thiserror::Error, PartialEq, Eq)]
pub enum WireTrustedPeerIdentityError {
#[error("external peer identity public_key must start with 'ed25519:'")]
MissingEd25519Prefix,
#[error("external peer identity public_key is not valid base64: {0}")]
InvalidBase64(String),
#[error("external peer identity public_key must decode to 32 bytes, got {actual}")]
InvalidLength { actual: usize },
#[error("external peer identity public_key must be non-zero")]
ZeroPublicKey,
}
impl WireTrustedPeerIdentity {
pub fn resolve(&self) -> Result<ResolvedWireTrustedPeerIdentity, WireTrustedPeerIdentityError> {
match self {
Self::Ed25519PublicKey { public_key } => {
let pubkey = parse_ed25519_public_key(public_key)?;
if pubkey == [0u8; 32] {
return Err(WireTrustedPeerIdentityError::ZeroPublicKey);
}
Ok(ResolvedWireTrustedPeerIdentity {
peer_id: meerkat_core::comms::PeerId::from_ed25519_pubkey(&pubkey),
pubkey,
})
}
}
}
}
fn parse_ed25519_public_key(raw: &str) -> Result<[u8; 32], WireTrustedPeerIdentityError> {
const PREFIX: &str = "ed25519:";
let encoded = raw
.strip_prefix(PREFIX)
.ok_or(WireTrustedPeerIdentityError::MissingEd25519Prefix)?;
let bytes = BASE64
.decode(encoded)
.map_err(|err| WireTrustedPeerIdentityError::InvalidBase64(err.to_string()))?;
let actual = bytes.len();
let pubkey: [u8; 32] = bytes
.try_into()
.map_err(|_| WireTrustedPeerIdentityError::InvalidLength { actual })?;
Ok(pubkey)
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireTrustedPeerSpec {
pub name: String,
pub address: String,
pub identity: WireTrustedPeerIdentity,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum MobPeerTarget {
Local(String),
External(WireTrustedPeerSpec),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobWireParams {
pub mob_id: String,
pub member: String,
pub peer: MobPeerTarget,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobWireResult {
pub wired: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobWireMembersBatchEdge {
pub a: String,
pub b: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobWireMembersBatchParams {
pub mob_id: String,
pub edges: Vec<MobWireMembersBatchEdge>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobWireMembersBatchResult {
pub requested: usize,
pub wired: Vec<MobWireMembersBatchEdge>,
pub already_wired: Vec<MobWireMembersBatchEdge>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobUnwireParams {
pub mob_id: String,
pub member: String,
pub peer: MobPeerTarget,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobUnwireResult {
pub unwired: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberSendParams {
pub mob_id: String,
pub agent_identity: String,
pub content: WireContentInput,
#[serde(default)]
pub handling_mode: WireHandlingMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub render_metadata: Option<WireRenderMetadata>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireAgentRuntimeId {
pub identity: String,
pub generation: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMemberSendResult {
pub mob_id: String,
pub agent_identity: String,
pub member_ref: WireMemberRef,
pub handling_mode: WireHandlingMode,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobIngressInteractionParams {
pub mob_id: String,
pub spec: MobMemberSpecWire,
pub content: WireContentInput,
#[serde(default)]
pub handling_mode: WireHandlingMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub render_metadata: Option<WireRenderMetadata>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobIngressInteractionResult {
pub mob_id: String,
pub agent_identity: String,
pub member_ref: WireMemberRef,
pub ensure_outcome: MobEnsureMemberOutcomeWire,
pub delivery: MobMemberSendResult,
pub events_after_cursor: u64,
pub latest_event_cursor: u64,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireHandlingMode {
#[default]
Queue,
Steer,
}
impl From<WireHandlingMode> for HandlingMode {
fn from(mode: WireHandlingMode) -> Self {
match mode {
WireHandlingMode::Queue => HandlingMode::Queue,
WireHandlingMode::Steer => HandlingMode::Steer,
}
}
}
impl From<HandlingMode> for WireHandlingMode {
fn from(mode: HandlingMode) -> Self {
match mode {
HandlingMode::Queue => WireHandlingMode::Queue,
HandlingMode::Steer => WireHandlingMode::Steer,
}
}
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireRenderClass {
UserPrompt,
PeerMessage,
PeerRequest,
PeerResponse,
ExternalEvent,
FlowStep,
Continuation,
SystemNotice,
ToolScopeNotice,
OpsProgress,
}
impl From<WireRenderClass> for RenderClass {
fn from(class: WireRenderClass) -> Self {
match class {
WireRenderClass::UserPrompt => RenderClass::UserPrompt,
WireRenderClass::PeerMessage => RenderClass::PeerMessage,
WireRenderClass::PeerRequest => RenderClass::PeerRequest,
WireRenderClass::PeerResponse => RenderClass::PeerResponse,
WireRenderClass::ExternalEvent => RenderClass::ExternalEvent,
WireRenderClass::FlowStep => RenderClass::FlowStep,
WireRenderClass::Continuation => RenderClass::Continuation,
WireRenderClass::SystemNotice => RenderClass::SystemNotice,
WireRenderClass::ToolScopeNotice => RenderClass::ToolScopeNotice,
WireRenderClass::OpsProgress => RenderClass::OpsProgress,
}
}
}
impl From<RenderClass> for WireRenderClass {
fn from(class: RenderClass) -> Self {
match class {
RenderClass::UserPrompt => WireRenderClass::UserPrompt,
RenderClass::PeerMessage => WireRenderClass::PeerMessage,
RenderClass::PeerRequest => WireRenderClass::PeerRequest,
RenderClass::PeerResponse => WireRenderClass::PeerResponse,
RenderClass::ExternalEvent => WireRenderClass::ExternalEvent,
RenderClass::FlowStep => WireRenderClass::FlowStep,
RenderClass::Continuation => WireRenderClass::Continuation,
RenderClass::SystemNotice => WireRenderClass::SystemNotice,
RenderClass::ToolScopeNotice => WireRenderClass::ToolScopeNotice,
RenderClass::OpsProgress => WireRenderClass::OpsProgress,
}
}
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireRenderSalience {
Background,
Normal,
Important,
Urgent,
}
impl From<WireRenderSalience> for RenderSalience {
fn from(salience: WireRenderSalience) -> Self {
match salience {
WireRenderSalience::Background => RenderSalience::Background,
WireRenderSalience::Normal => RenderSalience::Normal,
WireRenderSalience::Important => RenderSalience::Important,
WireRenderSalience::Urgent => RenderSalience::Urgent,
}
}
}
impl From<RenderSalience> for WireRenderSalience {
fn from(salience: RenderSalience) -> Self {
match salience {
RenderSalience::Background => WireRenderSalience::Background,
RenderSalience::Normal => WireRenderSalience::Normal,
RenderSalience::Important => WireRenderSalience::Important,
RenderSalience::Urgent => WireRenderSalience::Urgent,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireRenderMetadata {
pub class: WireRenderClass,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub salience: Option<WireRenderSalience>,
}
impl From<WireRenderMetadata> for RenderMetadata {
fn from(metadata: WireRenderMetadata) -> Self {
Self {
class: metadata.class.into(),
salience: metadata
.salience
.unwrap_or(WireRenderSalience::Normal)
.into(),
}
}
}
impl From<RenderMetadata> for WireRenderMetadata {
fn from(metadata: RenderMetadata) -> Self {
Self {
class: metadata.class.into(),
salience: Some(metadata.salience.into()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMemberSpecWire {
pub profile: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub initial_message: Option<WireContentInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub placement: Option<WireHostRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub binding: Option<WireRuntimeBinding>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub labels: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auto_wire_parent: Option<bool>,
}
impl MobMemberSpecWire {
#[must_use]
pub fn surface_metadata(&self) -> SurfaceMetadata {
SurfaceMetadata::from_optional_parts(self.labels.clone(), self.context.clone())
}
pub fn validate_public_surface_metadata(&self) -> Result<(), SurfaceMetadataError> {
self.surface_metadata().validate_public()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobEnsureMemberParams {
pub mob_id: String,
pub spec: MobMemberSpecWire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(transparent)]
pub struct WireMemberRef(String);
impl WireMemberRef {
#[must_use]
pub fn encode(mob_id: &str, agent_identity: &str) -> Self {
let payload = serde_json::json!({ "m": mob_id, "a": agent_identity });
Self(base64_url_encode(payload.to_string().as_bytes()))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
#[must_use]
pub fn from_token(token: impl Into<String>) -> Self {
Self(token.into())
}
pub fn decode(&self) -> Result<(String, String), WireMemberRefError> {
let bytes = base64_url_decode(&self.0).map_err(|_| WireMemberRefError::Malformed)?;
let value: Value =
serde_json::from_slice(&bytes).map_err(|_| WireMemberRefError::Malformed)?;
let mob_id = value
.get("m")
.and_then(Value::as_str)
.ok_or(WireMemberRefError::Malformed)?;
let agent_identity = value
.get("a")
.and_then(Value::as_str)
.ok_or(WireMemberRefError::Malformed)?;
Ok((mob_id.to_string(), agent_identity.to_string()))
}
}
#[derive(Debug, thiserror::Error)]
pub enum WireMemberRefError {
#[error("malformed member ref token")]
Malformed,
}
fn base64_url_encode(bytes: &[u8]) -> String {
use base64::Engine as _;
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(bytes)
}
fn base64_url_decode(input: &str) -> Result<Vec<u8>, base64::DecodeError> {
use base64::Engine as _;
base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(input)
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobSpawnReceiptWire {
pub agent_identity: String,
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobMemberStatus {
Active,
Retiring,
Broken,
Completed,
Unknown,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMemberListEntryWire {
pub agent_identity: String,
pub member_ref: WireMemberRef,
pub role: String,
pub runtime_mode: WireMobRuntimeMode,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub wired_to: Vec<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub labels: BTreeMap<String, String>,
pub status: WireMobMemberStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
pub is_final: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub enum MobEnsureMemberOutcomeWire {
#[serde(rename = "spawned")]
Spawned(MobSpawnReceiptWire),
#[serde(rename = "existed")]
Existed(MobMemberListEntryWire),
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobEnsureMemberResult {
pub outcome: MobEnsureMemberOutcomeWire,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobReconcileOptionsWire {
#[serde(default)]
pub retire_stale: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobReconcileStage {
Spawn,
Retire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobReconcileParams {
pub mob_id: String,
#[serde(default)]
pub desired: Vec<MobMemberSpecWire>,
#[serde(default)]
pub options: MobReconcileOptionsWire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireMobError {
pub code: MobSpawnManyFailureCause,
pub message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobReconcileFailureWire {
pub agent_identity: String,
pub stage: WireMobReconcileStage,
pub error: WireMobError,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobReconcileReportWire {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub desired: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub retained: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub spawned: Vec<MobSpawnReceiptWire>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub retired: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub failures: Vec<MobReconcileFailureWire>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobReconcileResult {
pub report: MobReconcileReportWire,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobLifecycleAction {
Stop,
Resume,
Complete,
Reset,
Destroy,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobWireAction {
Wire,
Unwire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobLifecycleParams {
pub mob_id: String,
pub action: WireMobLifecycleAction,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobLifecycleResult {
pub mob_id: String,
pub action: WireMobLifecycleAction,
pub ok: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub destroy_report: Option<Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobAppendSystemContextParams {
pub mob_id: String,
pub agent_identity: String,
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireAppendSystemContextStatus {
Applied,
Staged,
Duplicate,
}
impl From<meerkat_core::AppendSystemContextStatus> for WireAppendSystemContextStatus {
fn from(status: meerkat_core::AppendSystemContextStatus) -> Self {
match status {
meerkat_core::AppendSystemContextStatus::Applied => Self::Applied,
meerkat_core::AppendSystemContextStatus::Staged => Self::Staged,
meerkat_core::AppendSystemContextStatus::Duplicate => Self::Duplicate,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobAppendSystemContextResult {
pub mob_id: String,
pub agent_identity: String,
pub status: WireAppendSystemContextStatus,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobFlowsResult {
pub mob_id: String,
pub flows: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFlowRunParams {
pub mob_id: String,
pub flow_id: String,
#[serde(default)]
pub params: Value,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRunParams {
pub mob_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub flow_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt: Option<String>,
#[serde(default)]
pub params: Value,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobFlowRunResult {
pub run_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFlowStatusParams {
pub mob_id: String,
pub run_id: String,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMobRunStatus {
Pending,
Running,
Completed,
Failed,
Canceled,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireMobRun {
pub run_id: String,
pub mob_id: String,
pub flow_id: String,
pub status: WireMobRunStatus,
#[serde(flatten)]
pub kernel: serde_json::Map<String, Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobFlowStatusResult {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run: Option<WireMobRun>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRunResultParams {
pub mob_id: String,
pub run_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireMobRunResultEnvelope {
pub run_id: String,
pub mob_id: String,
pub flow_id: String,
pub status: WireMobRunStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub result: Option<Value>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub outputs: BTreeMap<String, Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRunResult {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run: Option<WireMobRunResultEnvelope>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobFlowCancelParams {
pub mob_id: String,
pub run_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobFlowCancelResult {
pub canceled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSpawnHelperParams {
pub mob_id: String,
pub prompt: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_identity: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_override: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobForkHelperParams {
pub mob_id: String,
pub source_member_id: String,
pub prompt: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_identity: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_override: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub fork_context: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<WireMobBackendKind>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobHelperResult {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<String>,
pub tokens_used: u64,
pub agent_identity: String,
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobForceCancelResult {
pub cancelled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobTurnStartParams {
pub mob_id: String,
pub agent_identity: String,
pub prompt: WireContentInput,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skill_refs: Option<Vec<meerkat_core::skills::SkillRef>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub turn_tool_overlay: Option<meerkat_core::service::PublicTurnToolOverlay>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub keep_alive: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub self_hosted_server_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_tokens: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_prompt: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_schema: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub structured_output_retries: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider_params:
Option<WireTurnMetadataOverride<crate::wire::runtime::WireProviderParamsOverride>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireTurnMetadataOverride<WireAuthBindingRef>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub injected_context: Option<Vec<WireContentInput>>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireUnreachablePeer {
pub peer: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WirePeerConnectivitySnapshot {
pub reachable_peer_count: usize,
pub unknown_peer_count: usize,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub unreachable_peers: Vec<WireUnreachablePeer>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum WirePeerConnectivity {
NotApplicable,
ProbeTimedOut,
Known {
snapshot: WirePeerConnectivitySnapshot,
},
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMemberStatusResult {
pub status: WireMobMemberStatus,
pub member_ref: WireMemberRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_preview: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
pub tokens_used: u64,
pub is_final: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub peer_connectivity: Option<WirePeerConnectivity>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub kickoff: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub external_member: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resolved_capabilities: Option<crate::wire::WireResolvedModelCapabilities>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub progress: Option<WireMemberProgressSnapshot>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub placement: Option<WireHostRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub control_reachability: Option<WireReachability>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub comms_reachability: Option<WireReachability>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_seen_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub freshness_reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub lifecycle_capabilities: Option<WireMemberLifecycleCapabilities>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub non_portable_disabled: Option<Vec<super::portable_spec::WireNonPortableResourceKind>>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMemberRunState {
Idle,
RunOpen,
Unknown,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMemberHealthClass {
Healthy,
Degraded,
Wedged,
Unknown,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireMemberProgressEvent {
ExecutionAdvanced,
BecameIdle,
Unchanged,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireMemberProgressSnapshot {
pub run_state: WireMemberRunState,
pub in_flight_work: u64,
pub last_progress_at_ms: u64,
pub last_progress_event: WireMemberProgressEvent,
pub health: WireMemberHealthClass,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobSnapshotResult {
pub mob_id: String,
pub status: WireMobLifecycleStatus,
pub members: Vec<MobMemberListEntryWire>,
}
#[cfg(test)]
mod member_status_capability_tests {
use super::*;
#[test]
fn member_status_result_round_trips_resolved_capabilities() -> Result<(), serde_json::Error> {
let capabilities = crate::wire::WireResolvedModelCapabilities {
vision: true,
image_input: true,
image_tool_results: false,
inline_video: false,
realtime: true,
web_search: true,
image_generation: true,
};
let result = MobMemberStatusResult {
status: WireMobMemberStatus::Active,
member_ref: WireMemberRef::encode("mob-1", "worker-1"),
output_preview: None,
error: None,
tokens_used: 0,
is_final: false,
current_session_id: Some("session-1".to_string()),
peer_connectivity: Some(WirePeerConnectivity::Known {
snapshot: WirePeerConnectivitySnapshot {
reachable_peer_count: 1,
unknown_peer_count: 0,
unreachable_peers: Vec::new(),
},
}),
kickoff: None,
external_member: None,
resolved_capabilities: Some(capabilities.clone()),
progress: None,
placement: None,
control_reachability: None,
comms_reachability: None,
last_seen_ms: None,
freshness_reason: None,
lifecycle_capabilities: None,
non_portable_disabled: None,
};
let json = serde_json::to_string(&result)?;
assert!(json.contains("\"resolved_capabilities\""));
let parsed: MobMemberStatusResult = serde_json::from_str(&json)?;
assert_eq!(parsed.resolved_capabilities, Some(capabilities));
Ok(())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobDestroyResult {
pub mob_id: String,
pub ok: bool,
pub destroy_report: Value,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRotateSupervisorResult {
pub mob_id: String,
pub ok: bool,
pub report: SupervisorRotationReportWire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct SupervisorRotationReportWire {
pub previous_epoch: u64,
pub current_epoch: u64,
pub public_peer_id: String,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum SupervisorRotationIncompleteKind {
SupervisorRotationIncomplete,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum SupervisorRotationRetryAuthority {
PendingRotation,
PreRotation,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum SupervisorRotationRetryScope {
Durable,
PreRotation,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub struct SupervisorRotationIncompleteDetailsWire {
pub kind: SupervisorRotationIncompleteKind,
pub previous_epoch: u64,
pub attempted_epoch: u64,
pub attempted_public_peer_id: String,
pub rotated_peer_count: usize,
pub rollback_succeeded: bool,
pub pending_authority_recorded: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rollback_error: Option<String>,
pub retry_authority: SupervisorRotationRetryAuthority,
pub retry_scope: SupervisorRotationRetryScope,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub struct SupervisorRotationIncompleteDataWire {
pub code: String,
pub message: String,
pub details: SupervisorRotationIncompleteDetailsWire,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobWaitParams {
pub mob_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub member_ids: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_ms: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobWaitMembersResult {
pub members: Vec<Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobCancelWorkResult {
pub mob_id: String,
pub ok: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobCancelAllWorkResult {
pub mob_id: String,
pub ok: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobProfileCreateParams {
pub name: String,
pub profile: MobProfileInput,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobProfileNameParams {
pub name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobProfileUpdateParams {
pub name: String,
pub profile: MobProfileInput,
pub expected_revision: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobProfileDeleteParams {
pub name: String,
pub expected_revision: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobProfileLookupResult {
#[serde(default)]
pub not_found: bool,
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<WireMobProfile>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub created_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub updated_at: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobProfileListResult {
pub profiles: Vec<MobProfileLookupResult>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobProfileDeleteResult {
pub name: String,
pub deleted_revision: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobStreamOpenParams {
pub mob_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_identity: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobStreamOpenResult {
pub stream_id: String,
pub opened: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobStreamCloseParams {
pub stream_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobStreamCloseResult {
pub stream_id: String,
pub closed: bool,
pub already_closed: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireWorkOrigin {
#[default]
External,
Internal,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobSubmitWorkParams {
pub member_ref: WireMemberRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub work_ref: Option<String>,
pub content: WireContentInput,
#[serde(default)]
pub origin: WireWorkOrigin,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub injected_context: Option<Vec<WireContentInput>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub objective_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobSubmitWorkResult {
pub mob_id: String,
pub work_ref: String,
pub member_ref: WireMemberRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub objective_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobConcludeObjectiveParams {
pub member_ref: WireMemberRef,
pub objective_id: String,
pub outcome: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobConcludeObjectiveResult {
pub member_ref: WireMemberRef,
pub objective_id: String,
pub concluded: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobCancelWorkParams {
pub mob_id: String,
pub work_ref: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobCancelAllWorkParams {
pub member_ref: WireMemberRef,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberFilterWire {
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub labels: BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub status: Option<WireMobMemberStatus>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobListMembersMatchingParams {
pub mob_id: String,
#[serde(default)]
pub filter: MobMemberFilterWire,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobListMembersMatchingResult {
#[serde(default)]
pub members: Vec<Value>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireControlScope {
List,
ReadHistory,
SubscribeEvents,
SendCommand,
Cancel,
Retire,
WireTopology,
Live,
AdminHost,
AdminGrants,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireReachability {
Reachable,
Stale,
Unreachable,
Unknown,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(transparent)]
pub struct WireHostRef(pub String);
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireMemberLifecycleCapabilities {
pub transcript_edits: bool,
pub revisions: bool,
pub resume_after_restart: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireHostBindPhase {
Requested,
Bound,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireHostCapabilityFlags {
pub protocol_min: u64,
pub protocol_max: u64,
pub engine_version: String,
pub durable_sessions: bool,
pub autonomous_members: bool,
pub hard_cancel_member: bool,
#[serde(default)]
pub tracked_input_cancel: bool,
pub memory_store: bool,
pub mcp: bool,
#[serde(default, skip_serializing_if = "std::collections::BTreeSet::is_empty")]
pub resolvable_providers: std::collections::BTreeSet<String>,
pub approval_forwarding: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub live_endpoint: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobHostStatus {
pub host_id: WireHostRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub endpoint: Option<String>,
pub bind_phase: WireHostBindPhase,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub authority_epoch: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub capabilities: Option<WireHostCapabilityFlags>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub control_reachability: Option<WireReachability>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_seen_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub freshness_reason: Option<String>,
pub materialized_member_count: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobHostsResult {
pub hosts: Vec<MobHostStatus>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireRouteInstallObligation {
pub edge_a: String,
pub edge_b: String,
pub host: WireHostRef,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRouteInstallsResult {
pub outstanding: Vec<WireRouteInstallObligation>,
pub complete: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "snake_case")]
pub enum WireProjectionProvenance {
HostClaimed,
ControllingHostVerified,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(transparent)]
pub struct WireHistoryRow(pub super::session::WireSessionMessage);
impl PartialEq for WireHistoryRow {
fn eq(&self, other: &Self) -> bool {
match (
serde_json::to_value(&self.0),
serde_json::to_value(&other.0),
) {
(Ok(a), Ok(b)) => a == b,
_ => false,
}
}
}
impl Eq for WireHistoryRow {}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireMemberHistoryPageBody {
pub from_index: u64,
pub messages: Vec<WireHistoryRow>,
pub message_count: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_index: Option<u64>,
pub complete: bool,
}
impl WireMemberHistoryPageBody {
pub fn try_from_history_page(
page: &meerkat_core::service::SessionHistoryPage,
) -> Result<Self, super::error::WireConversionError> {
let invalid =
|reason: String| super::error::WireConversionError::MemberHistoryPage { debug: reason };
let message_count = u64::try_from(page.message_count).map_err(|_| {
invalid(format!(
"message_count {} exceeds the u64 wire domain",
page.message_count
))
})?;
let from_index = u64::try_from(page.offset).map_err(|_| {
invalid(format!(
"offset {} exceeds the u64 wire domain",
page.offset
))
})?;
let served = u64::try_from(page.messages.len()).map_err(|_| {
invalid(format!(
"served row count {} exceeds the u64 wire domain",
page.messages.len()
))
})?;
let next_index = if page.has_more {
if served == 0 {
return Err(invalid(format!(
"page at offset {from_index} claims more rows but serves none"
)));
}
Some(from_index.checked_add(served).ok_or_else(|| {
invalid(format!(
"offset {from_index} plus served row count {served} exhausts the u64 cursor domain"
))
})?)
} else {
None
};
Ok(Self {
from_index,
messages: page
.messages
.iter()
.map(|message| {
WireHistoryRow(super::session::WireSessionMessage::from(message.clone()))
})
.collect(),
message_count,
next_index,
complete: !page.has_more,
})
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberHistoryParams {
pub mob_id: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub from_index: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub limit: Option<u32>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobMemberHistoryResult {
pub page: WireMemberHistoryPageBody,
pub generation: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub placement: Option<WireHostRef>,
pub provenance: WireProjectionProvenance,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobBindHostParams {
pub mob_id: String,
pub descriptor: super::supervisor_bridge::WireHostBindingDescriptor,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobBindHostResult {
pub host_id: WireHostRef,
pub capabilities: WireHostCapabilityFlags,
pub authority_epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRevokeHostParams {
pub mob_id: String,
pub host_id: WireHostRef,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRevokeHostResult {
pub host_id: WireHostRef,
pub released_members: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct WireGrantRecord {
pub principal: String,
pub scopes: Vec<WireControlScope>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expires_at_ms: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobGrantScopesParams {
pub mob_id: String,
pub principal: String,
pub scopes: Vec<WireControlScope>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expires_at_ms: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobGrantScopesResult {
pub record: WireGrantRecord,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobRevokeScopesParams {
pub mob_id: String,
pub principal: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scopes: Option<Vec<WireControlScope>>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobRevokeScopesResult {
pub removed: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobGrantsResult {
pub grants: Vec<WireGrantRecord>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct WireScopeDeniedDetail {
pub required: WireControlScope,
pub presented: Vec<WireControlScope>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberLiveOpenParams {
pub mob_id: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub turning_mode: Option<super::realtime::RealtimeTurningMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transport: Option<super::live::LiveOpenTransport>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberLiveChannelParams {
pub mob_id: String,
pub agent_identity: String,
pub channel_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberLiveStatusParams {
pub mob_id: String,
pub agent_identity: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub channel_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobHardCancelParams {
pub mob_id: String,
pub agent_identity: String,
pub reason: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
pub struct MobHardCancelResult {
pub cancelled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct MobMemberLiveControlParams {
pub mob_id: String,
pub agent_identity: String,
pub channel_id: String,
pub verb: super::supervisor_bridge::BridgeLiveControlVerb,
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
#[test]
fn mob_helper_params_carry_structural_auth_binding() {
let parsed: MobSpawnHelperParams = serde_json::from_value(serde_json::json!({
"mob_id": "mob-1",
"prompt": "help",
"agent_identity": "helper",
"auth_binding": {
"realm": "dev",
"binding": "default_anthropic",
"profile": "console"
}
}))
.expect("spawn helper params parse");
let auth_binding = parsed.auth_binding.expect("auth_binding should parse");
assert_eq!(auth_binding.realm.as_str(), "dev");
assert_eq!(auth_binding.binding.as_str(), "default_anthropic");
assert_eq!(
auth_binding
.profile
.as_ref()
.map(|profile| profile.as_str()),
Some("console")
);
let parsed: MobForkHelperParams = serde_json::from_value(serde_json::json!({
"mob_id": "mob-1",
"source_member_id": "source",
"prompt": "help",
"agent_identity": "helper",
"auth_binding": {
"realm": "dev",
"binding": "default_anthropic"
}
}))
.expect("fork helper params parse");
let auth_binding = parsed.auth_binding.expect("auth_binding should parse");
assert_eq!(auth_binding.realm.as_str(), "dev");
assert_eq!(auth_binding.binding.as_str(), "default_anthropic");
assert!(auth_binding.profile.is_none());
}
#[test]
fn wire_mob_profile_parses_provider_fields_fail_closed() {
let legacy: WireMobProfile =
serde_json::from_str(r#"{"model":"claude-opus-4-8"}"#).expect("legacy profile parses");
assert_eq!(legacy.provider, None);
assert!(legacy.resume_overrides.is_empty());
let full: WireMobProfile = serde_json::from_str(
r#"{
"model": "claude-internal-preview",
"provider": "anthropic",
"image_generation_provider": "gemini",
"auto_compact_threshold": 60000,
"resume_overrides": ["model", "provider"]
}"#,
)
.expect("typed profile parses");
assert_eq!(full.provider, Some(meerkat_core::Provider::Anthropic));
assert_eq!(
full.image_generation_provider,
Some(meerkat_core::Provider::Gemini)
);
assert_eq!(
full.resume_overrides,
vec![
WireMobResumeOverrideField::Model,
WireMobResumeOverrideField::Provider
]
);
assert!(
serde_json::from_str::<WireMobProfile>(r#"{"model":"m","provider":"not-a-provider"}"#)
.is_err(),
"unknown provider names must fail closed at the wire boundary"
);
assert!(
serde_json::from_str::<WireMobProfile>(r#"{"model":"m","auto_compact_threshold":0}"#)
.is_err(),
"zero auto_compact_threshold must fail closed at the wire boundary"
);
assert!(
serde_json::from_str::<WireMobProfile>(
r#"{"model":"m","resume_overrides":["everything"]}"#
)
.is_err(),
"resume_overrides vocabulary is closed"
);
}
#[test]
fn mob_definition_input_parses_custom_models() {
let input: MobDefinitionInput = serde_json::from_str(
r#"{
"id": "m",
"profiles": {"worker": {"model": "claude-internal-preview"}},
"models": {
"claude-internal-preview": {
"provider": "anthropic",
"context_window": 500000,
"vision": true
}
},
"image_generation_provider": "openai"
}"#,
)
.expect("definition with custom models parses");
let model = input
.models
.get("claude-internal-preview")
.expect("custom model present");
assert_eq!(model.provider, meerkat_core::Provider::Anthropic);
assert_eq!(model.context_window, Some(500_000));
assert_eq!(model.vision, Some(true));
assert_eq!(
input.image_generation_provider,
Some(meerkat_core::Provider::OpenAI)
);
}
#[test]
fn wire_member_ref_round_trips_through_encode_decode() {
let token = WireMemberRef::encode("mob-42", "worker-1");
let (mob_id, agent_identity) = token.decode().expect("decode round-trips");
assert_eq!(mob_id, "mob-42");
assert_eq!(agent_identity, "worker-1");
}
#[test]
fn wire_member_ref_rejects_malformed_token() {
let err = WireMemberRef::from_token("not-a-token-payload")
.decode()
.expect_err("malformed tokens must fail to decode");
assert!(matches!(err, WireMemberRefError::Malformed));
}
#[test]
fn mob_spawn_many_spec_placement_is_optional_and_round_trips() {
let placed: MobSpawnSpecParams = serde_json::from_value(serde_json::json!({
"profile": "worker",
"agent_identity": "w1",
"placement": "host-b-peer"
}))
.expect("placed spawn-many spec parses");
assert_eq!(
placed.placement.as_ref().map(|host| host.0.as_str()),
Some("host-b-peer")
);
assert_eq!(
serde_json::to_value(&placed).expect("placed spawn-many spec serializes")["placement"],
"host-b-peer"
);
let local: MobSpawnSpecParams = serde_json::from_value(serde_json::json!({
"profile": "worker",
"agent_identity": "w2"
}))
.expect("local spawn-many spec parses");
assert!(local.placement.is_none());
assert!(
serde_json::to_value(&local)
.expect("local spawn-many spec serializes")
.get("placement")
.is_none(),
"absent placement must remain omitted for source and wire compatibility"
);
}
#[test]
fn mob_member_spec_placement_is_optional_and_round_trips() {
let placed: MobMemberSpecWire = serde_json::from_value(serde_json::json!({
"profile": "worker",
"agent_identity": "w1",
"placement": "host-b-peer"
}))
.expect("placed declarative member spec parses");
assert_eq!(
placed.placement.as_ref().map(|host| host.0.as_str()),
Some("host-b-peer")
);
let local: MobMemberSpecWire = serde_json::from_value(serde_json::json!({
"profile": "worker",
"agent_identity": "w2"
}))
.expect("local declarative member spec parses");
assert!(local.placement.is_none());
}
#[test]
fn mob_member_spec_exposes_shared_surface_metadata() {
let spec = MobMemberSpecWire {
profile: "worker".into(),
agent_identity: "w1".into(),
initial_message: None,
runtime_mode: None,
backend: None,
placement: None,
binding: None,
context: Some(serde_json::json!({"client_ref": "member-card"})),
labels: Some(BTreeMap::from([("client.member_id".into(), "w1".into())])),
additional_instructions: None,
auto_wire_parent: None,
};
let metadata = spec.surface_metadata();
assert_eq!(
metadata.labels.get("client.member_id").map(String::as_str),
Some("w1")
);
assert_eq!(
metadata.app_context,
Some(serde_json::json!({"client_ref": "member-card"}))
);
}
#[test]
fn mob_member_spec_surface_metadata_rejects_reserved_keys() {
let spec = MobMemberSpecWire {
profile: "worker".into(),
agent_identity: "w1".into(),
initial_message: None,
runtime_mode: None,
backend: None,
placement: None,
binding: None,
context: None,
labels: Some(BTreeMap::from([("mob_id".into(), "spoof".into())])),
additional_instructions: None,
auto_wire_parent: None,
};
assert!(spec.validate_public_surface_metadata().is_err());
}
#[test]
fn mob_reconcile_failure_stage_is_typed_wire_enum() {
let failure = MobReconcileFailureWire {
agent_identity: "worker-1".into(),
stage: WireMobReconcileStage::Spawn,
error: WireMobError {
code: MobSpawnManyFailureCause::ProfileNotFound,
message: "spawn failed".into(),
},
};
let json = serde_json::to_value(&failure).expect("serialize failure");
assert_eq!(json["stage"], "spawn");
assert_eq!(json["error"]["code"], "profile_not_found");
assert_eq!(json["error"]["message"], "spawn failed");
let round_trip: MobReconcileFailureWire =
serde_json::from_value(json).expect("deserialize failure");
assert_eq!(round_trip.stage, WireMobReconcileStage::Spawn);
assert_eq!(
round_trip.error.code,
MobSpawnManyFailureCause::ProfileNotFound
);
let err = serde_json::from_value::<MobReconcileFailureWire>(serde_json::json!({
"agent_identity": "worker-1",
"stage": "restart",
"error": { "code": "profile_not_found", "message": "bad stage" }
}))
.expect_err("unknown reconcile stage must be rejected");
assert!(err.to_string().contains("unknown variant"));
}
#[test]
fn mob_lifecycle_params_reject_unknown_action_string() {
let err = serde_json::from_value::<MobLifecycleParams>(serde_json::json!({
"mob_id": "mob-1",
"action": "explode"
}))
.expect_err("unknown lifecycle actions must fail at the typed wire boundary");
assert!(
err.to_string().contains("unknown variant"),
"unexpected error: {err}"
);
}
#[test]
fn mob_lifecycle_result_round_trips_typed_action() {
let result = MobLifecycleResult {
mob_id: "mob-1".into(),
action: WireMobLifecycleAction::Complete,
ok: true,
destroy_report: None,
};
let json = serde_json::to_value(&result).expect("serialize lifecycle result");
assert_eq!(json["action"], "complete");
let round_trip: MobLifecycleResult =
serde_json::from_value(json).expect("deserialize lifecycle result");
assert_eq!(round_trip.action, WireMobLifecycleAction::Complete);
}
#[test]
fn mob_wire_members_batch_contract_is_local_edge_native() {
let params: MobWireMembersBatchParams = serde_json::from_value(serde_json::json!({
"mob_id": "mob-1",
"edges": [
{ "a": "lead", "b": "worker-b" },
{ "a": "worker-a", "b": "lead" }
]
}))
.expect("batch wire params deserialize");
assert_eq!(params.mob_id, "mob-1");
assert_eq!(params.edges.len(), 2);
assert_eq!(params.edges[0].a, "lead");
assert_eq!(params.edges[0].b, "worker-b");
let result = MobWireMembersBatchResult {
requested: 2,
wired: vec![MobWireMembersBatchEdge {
a: "lead".into(),
b: "worker-a".into(),
}],
already_wired: vec![MobWireMembersBatchEdge {
a: "lead".into(),
b: "worker-b".into(),
}],
};
let json = serde_json::to_value(&result).expect("serialize batch wire result");
assert_eq!(json["requested"], 2);
assert_eq!(json["wired"][0]["a"], "lead");
assert_eq!(json["already_wired"][0]["b"], "worker-b");
let err = serde_json::from_value::<MobWireMembersBatchParams>(serde_json::json!({
"mob_id": "mob-1",
"edges": [{ "member": "lead", "peer": "worker-a" }]
}))
.expect_err("mixed local/external mob/wire shape must not deserialize");
let message = err.to_string();
assert!(
message.contains("unknown field `member`") || message.contains("missing field `a`"),
"unexpected error: {message}"
);
}
#[test]
fn mob_spawn_many_result_entry_uses_typed_status_result_envelope() {
let member_ref = WireMemberRef::encode("mob-1", "worker-1");
let entry = MobSpawnManyResultEntry::spawned("worker-1", member_ref.clone());
let json = serde_json::to_value(&entry).expect("serialize typed spawn_many row");
assert_eq!(json["status"], "spawned");
assert_eq!(json["result"]["agent_identity"], "worker-1");
assert_eq!(json["result"]["member_ref"], member_ref.as_str());
assert!(json.get("ok").is_none());
assert!(json.get("error").is_none());
let round_trip: MobSpawnManyResultEntry =
serde_json::from_value(json).expect("deserialize typed spawn_many row");
assert_eq!(round_trip, entry);
let failed = MobSpawnManyResultEntry::failed(
MobSpawnManyFailureCause::ProfileNotFound,
"profile missing",
);
let json = serde_json::to_value(&failed).expect("serialize typed failed spawn_many row");
assert_eq!(json["status"], "failed");
assert_eq!(json["result"]["cause"], "profile_not_found");
assert_eq!(json["result"]["message"], "profile missing");
assert!(json.get("ok").is_none());
assert!(json.get("error").is_none());
let round_trip: MobSpawnManyResultEntry =
serde_json::from_value(json).expect("deserialize typed failed spawn_many row");
assert_eq!(round_trip, failed);
}
#[test]
fn mob_spawn_many_result_entry_rejects_legacy_or_malformed_envelopes() {
let legacy = serde_json::json!({
"ok": true,
"agent_identity": "worker-1",
"member_ref": WireMemberRef::encode("mob-1", "worker-1"),
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(legacy)
.expect_err("legacy ok carrier must not deserialize");
assert!(
err.to_string().contains("missing field `status`")
|| err.to_string().contains("unknown field"),
"unexpected error: {err}"
);
let missing_result = serde_json::json!({
"status": "spawned"
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(missing_result)
.expect_err("missing typed result must fail closed");
assert!(
err.to_string().contains("missing field `result`"),
"unexpected error: {err}"
);
let unknown_status = serde_json::json!({
"status": "ok",
"result": {
"agent_identity": "worker-1",
"member_ref": WireMemberRef::encode("mob-1", "worker-1"),
}
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(unknown_status)
.expect_err("unknown typed status must fail closed");
assert!(
err.to_string().contains("unknown variant"),
"unexpected error: {err}"
);
let mismatched = serde_json::json!({
"status": "spawned",
"result": {
"cause": "profile_not_found",
"message": "profile missing"
}
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(mismatched)
.expect_err("status/result mismatch must fail closed");
assert!(
err.to_string()
.contains("status spawned requires spawned result"),
"unexpected error: {err}"
);
let message_only_failure = serde_json::json!({
"status": "failed",
"result": {
"message": "profile missing"
}
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(message_only_failure)
.expect_err("string-only failure result must fail closed");
assert!(
err.to_string().contains("data did not match any variant")
|| err.to_string().contains("missing field `cause`"),
"unexpected error: {err}"
);
let unknown_failure_cause = serde_json::json!({
"status": "failed",
"result": {
"cause": "future_failure",
"message": "future failure"
}
});
let err = serde_json::from_value::<MobSpawnManyResultEntry>(unknown_failure_cause)
.expect_err("unknown failure cause must fail closed");
assert!(
err.to_string().contains("data did not match any variant")
|| err.to_string().contains("unknown variant"),
"unexpected error: {err}"
);
}
#[test]
fn mob_wire_params_reject_legacy_local_target_shape() {
let err = serde_json::from_value::<MobWireParams>(serde_json::json!({
"mob_id": "mob-1",
"local": "member-a",
"target": { "local": "member-b" }
}))
.expect_err("legacy local/target shape must be rejected");
let msg = err.to_string();
assert!(
msg.contains("unknown field `local`") || msg.contains("missing field `member`"),
"unexpected error: {msg}"
);
}
#[test]
fn mob_wire_params_accept_canonical_external_peer_identity() {
let params = serde_json::from_value::<MobWireParams>(serde_json::json!({
"mob_id": "mob-1",
"member": "member-a",
"peer": {
"external": {
"name": "external-worker",
"address": "inproc://external-worker",
"identity": {
"kind": "ed25519_public_key",
"public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
}
}
}
}))
.expect("canonical external peer identity should deserialize");
let MobPeerTarget::External(spec) = params.peer else {
panic!("expected external peer target");
};
assert_eq!(spec.name, "external-worker");
}
#[test]
fn mob_wire_params_reject_raw_external_peer_id_shape() {
let err = serde_json::from_value::<MobWireParams>(serde_json::json!({
"mob_id": "mob-1",
"member": "member-a",
"peer": {
"external": {
"name": "external-worker",
"peer_id": meerkat_core::comms::PeerId::from_ed25519_pubkey(&[7u8; 32]).to_string(),
"address": "inproc://external-worker",
"pubkey": vec![7u8; 32]
}
}
}))
.expect_err("raw peer_id/pubkey external peer shape must be rejected");
let msg = err.to_string();
assert!(
msg.contains("peer_id") || msg.contains("identity"),
"unexpected error: {msg}"
);
}
#[test]
fn mob_wire_params_reject_missing_external_peer_pubkey_material() {
let err = serde_json::from_value::<MobWireParams>(serde_json::json!({
"mob_id": "mob-1",
"member": "member-a",
"peer": {
"external": {
"name": "external-worker",
"address": "inproc://external-worker",
"identity": {
"kind": "ed25519_public_key"
}
}
}
}))
.expect_err("missing external peer pubkey material must fail closed");
let msg = err.to_string();
assert!(
msg.contains("public_key") || msg.contains("identity"),
"unexpected error: {msg}"
);
}
#[test]
fn runtime_binding_accepts_canonical_external_peer_identity() {
let binding = serde_json::from_value::<WireRuntimeBinding>(serde_json::json!({
"kind": "external",
"address": "inproc://external-worker",
"identity": {
"kind": "ed25519_public_key",
"public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
}
}))
.expect("canonical external runtime binding identity should deserialize");
let WireRuntimeBinding::External {
identity, address, ..
} = binding
else {
panic!("expected external runtime binding");
};
assert_eq!(address, "inproc://external-worker");
assert_eq!(
identity.resolve().expect("identity resolves").pubkey,
[7u8; 32]
);
}
#[test]
fn runtime_binding_rejects_raw_external_peer_id_shape() {
let err = serde_json::from_value::<WireRuntimeBinding>(serde_json::json!({
"kind": "external",
"peer_id": meerkat_core::comms::PeerId::from_ed25519_pubkey(&[7u8; 32]).to_string(),
"address": "inproc://external-worker",
"pubkey": vec![7u8; 32]
}))
.expect_err("raw peer_id/pubkey external runtime binding shape must be rejected");
let msg = err.to_string();
assert!(
msg.contains("peer_id") || msg.contains("identity"),
"unexpected error: {msg}"
);
}
#[test]
fn runtime_binding_rejects_missing_external_peer_pubkey_material() {
let err = serde_json::from_value::<WireRuntimeBinding>(serde_json::json!({
"kind": "external",
"address": "inproc://external-worker",
"identity": {
"kind": "ed25519_public_key"
}
}))
.expect_err("missing external runtime binding pubkey material must fail closed");
let msg = err.to_string();
assert!(
msg.contains("public_key") || msg.contains("identity"),
"unexpected error: {msg}"
);
}
#[test]
fn mob_turn_start_params_capture_turn_override_fields() {
let params = serde_json::from_value::<MobTurnStartParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
"prompt": "continue",
"output_schema": { "type": "object" },
"structured_output_retries": 2
}))
.expect("turn_start should accept explicit turn override fields");
assert_eq!(params.mob_id, "mob-1");
assert_eq!(params.agent_identity, "worker");
assert_eq!(params.prompt, WireContentInput::Text("continue".into()));
assert_eq!(
params.output_schema,
Some(serde_json::json!({ "type": "object" }))
);
assert_eq!(params.structured_output_retries, Some(2));
let err = serde_json::from_value::<MobTurnStartParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
"prompt": "continue",
"unknown_override": true
}))
.expect_err("turn_start must reject unknown override fields");
assert!(
err.to_string().contains("unknown field"),
"unexpected error: {err}"
);
}
#[test]
fn mob_create_params_reject_reserved_runtime_lifecycle_fields() {
let err = serde_json::from_value::<MobCreateParams>(serde_json::json!({
"definition": {
"id": "mob-1",
"owner_runtime_binding": "runtime:worker:0",
"profiles": {
"worker": { "model": "claude-sonnet-4-6" }
}
}
}))
.expect_err("reserved runtime lifecycle fields must be rejected");
assert!(
err.to_string()
.contains("unknown field `owner_runtime_binding`"),
"unexpected error: {err}"
);
}
#[test]
fn mob_create_params_reject_reserved_runtime_bridge_owner_field() {
let err = serde_json::from_value::<MobCreateParams>(serde_json::json!({
"definition": {
"id": "mob-1",
"owner_transport_binding": "transport:worker:0",
"profiles": {
"worker": { "model": "claude-sonnet-4-6" }
}
}
}))
.expect_err("reserved runtime bridge owner field must be rejected");
assert!(
err.to_string()
.contains("unknown field `owner_transport_binding`"),
"unexpected error: {err}"
);
}
#[test]
fn mob_create_params_reject_internal_profile_tool_bundles() {
let err = serde_json::from_value::<MobCreateParams>(serde_json::json!({
"definition": {
"id": "mob-1",
"profiles": {
"worker": {
"model": "claude-sonnet-4-6",
"tools": {
"rust_bundles": ["internal-only"]
}
}
}
}
}))
.expect_err("internal rust tool bundles must be rejected");
assert!(
err.to_string().contains("did not match any variant")
|| err.to_string().contains("unknown field `rust_bundles`"),
"unexpected error: {err}"
);
}
#[test]
fn mob_create_params_accept_typed_nested_flow_definition() {
let params = serde_json::from_value::<MobCreateParams>(serde_json::json!({
"definition": {
"id": "mob-1",
"profiles": {
"worker": { "model": "claude-sonnet-4-6" }
},
"flows": {
"review": {
"description": "review flow",
"steps": {
"draft": {
"role": "worker",
"message": "draft it"
}
}
}
}
}
}))
.expect("typed nested flow definition should parse");
assert_eq!(
params.definition.flows["review"].steps["draft"].role,
"worker"
);
}
#[test]
fn mob_spawn_params_reject_deleted_budget_split_policy() {
let err = serde_json::from_value::<MobSpawnParams>(serde_json::json!({
"mob_id": "mob-1",
"profile": "worker",
"agent_identity": "worker-1",
"budget_split_policy": { "type": "equal" }
}))
.expect_err("deleted budget_split_policy must fail closed at the wire boundary");
assert!(
err.to_string()
.contains("unknown field `budget_split_policy`"),
"unexpected error: {err}"
);
}
#[test]
fn mob_spawn_params_placement_round_trips_and_defaults_absent() {
let params = serde_json::from_value::<MobSpawnParams>(serde_json::json!({
"mob_id": "mob-1",
"profile": "worker",
"agent_identity": "worker-1"
}))
.expect("placement-less params must decode");
assert_eq!(params.placement, None);
let encoded = serde_json::to_value(¶ms).expect("serialize params");
assert!(
encoded.get("placement").is_none(),
"absent placement must not serialize: {encoded}"
);
let params = serde_json::from_value::<MobSpawnParams>(serde_json::json!({
"mob_id": "mob-1",
"profile": "worker",
"agent_identity": "worker-1",
"placement": "host-peer-b"
}))
.expect("placed params must decode");
assert_eq!(params.placement.as_deref(), Some("host-peer-b"));
let encoded = serde_json::to_value(¶ms).expect("serialize params");
assert_eq!(encoded["placement"], serde_json::json!("host-peer-b"));
}
fn minimal_member_status() -> MobMemberStatusResult {
MobMemberStatusResult {
status: WireMobMemberStatus::Active,
member_ref: WireMemberRef::encode("mob-1", "worker-1"),
output_preview: None,
error: None,
tokens_used: 0,
is_final: false,
current_session_id: None,
peer_connectivity: None,
kickoff: None,
external_member: None,
resolved_capabilities: None,
progress: None,
placement: None,
control_reachability: None,
comms_reachability: None,
last_seen_ms: None,
freshness_reason: None,
lifecycle_capabilities: None,
non_portable_disabled: None,
}
}
#[test]
fn member_status_multi_host_fields_skip_when_absent_and_decode_legacy() {
let value = serde_json::to_value(minimal_member_status()).expect("serialize member status");
for absent in [
"placement",
"control_reachability",
"comms_reachability",
"last_seen_ms",
"freshness_reason",
"lifecycle_capabilities",
"non_portable_disabled",
] {
assert!(
value.get(absent).is_none(),
"absent {absent} must be omitted from the wire form: {value}"
);
}
let legacy = serde_json::json!({
"status": "active",
"member_ref": WireMemberRef::encode("mob-1", "worker-1"),
"tokens_used": 3,
"is_final": false,
});
let decoded: MobMemberStatusResult =
serde_json::from_value(legacy).expect("legacy member status decodes");
assert!(decoded.placement.is_none());
assert!(decoded.lifecycle_capabilities.is_none());
}
#[test]
fn member_status_carries_placement_only_at_typed_keys() {
let mut status = minimal_member_status();
status.placement = Some(WireHostRef("host-b-peer".to_string()));
status.control_reachability = Some(WireReachability::Stale);
status.comms_reachability = Some(WireReachability::Reachable);
status.last_seen_ms = Some(1_234);
status.freshness_reason = Some("pump idle".to_string());
status.lifecycle_capabilities = Some(WireMemberLifecycleCapabilities {
transcript_edits: false,
revisions: false,
resume_after_restart: true,
});
status.non_portable_disabled = Some(vec![
super::super::portable_spec::WireNonPortableResourceKind::WorkgraphTools,
]);
status.external_member = Some(serde_json::json!({"endpoint": "tcp://10.0.0.2:7101"}));
let value = serde_json::to_value(&status).expect("serialize member status");
assert_eq!(value["placement"], serde_json::json!("host-b-peer"));
assert_eq!(value["control_reachability"], serde_json::json!("stale"));
assert_eq!(
value["non_portable_disabled"],
serde_json::json!(["workgraph_tools"])
);
assert!(
value["external_member"].get("placement").is_none(),
"placement must not ride the opaque external_member value (SD-5)"
);
let decoded: MobMemberStatusResult =
serde_json::from_value(value).expect("decode member status");
assert_eq!(decoded.placement, status.placement);
assert_eq!(decoded.control_reachability, status.control_reachability);
}
#[test]
fn control_scope_and_reachability_round_trip_snake_case() {
let scopes: &[(WireControlScope, &str)] = &[
(WireControlScope::List, "list"),
(WireControlScope::ReadHistory, "read_history"),
(WireControlScope::SubscribeEvents, "subscribe_events"),
(WireControlScope::SendCommand, "send_command"),
(WireControlScope::Cancel, "cancel"),
(WireControlScope::Retire, "retire"),
(WireControlScope::WireTopology, "wire_topology"),
(WireControlScope::Live, "live"),
(WireControlScope::AdminHost, "admin_host"),
(WireControlScope::AdminGrants, "admin_grants"),
];
for (scope, expected) in scopes {
let value = serde_json::to_value(scope).expect("serialize scope");
assert_eq!(value, serde_json::json!(expected));
let decoded: WireControlScope = serde_json::from_value(value).expect("decode scope");
assert_eq!(decoded, *scope);
}
assert!(
serde_json::from_value::<WireControlScope>(serde_json::json!("admin")).is_err(),
"unknown scopes must fail decode (closed vocabulary)"
);
let classes: &[(WireReachability, &str)] = &[
(WireReachability::Reachable, "reachable"),
(WireReachability::Stale, "stale"),
(WireReachability::Unreachable, "unreachable"),
(WireReachability::Unknown, "unknown"),
];
for (class, expected) in classes {
let value = serde_json::to_value(class).expect("serialize reachability");
assert_eq!(value, serde_json::json!(expected));
let decoded: WireReachability =
serde_json::from_value(value).expect("decode reachability");
assert_eq!(decoded, *class);
}
}
#[test]
fn scope_denied_detail_round_trips_snake_case_and_denies_unknown_fields() {
let detail = WireScopeDeniedDetail {
required: WireControlScope::AdminGrants,
presented: vec![WireControlScope::List, WireControlScope::SendCommand],
};
let value = serde_json::to_value(&detail).expect("serialize detail");
assert_eq!(
value,
serde_json::json!({
"required": "admin_grants",
"presented": ["list", "send_command"],
})
);
let decoded: WireScopeDeniedDetail = serde_json::from_value(value).expect("decode detail");
assert_eq!(decoded, detail);
assert!(
serde_json::from_value::<WireScopeDeniedDetail>(serde_json::json!({
"required": "admin_grants",
"presented": [],
"reason": "extra",
}))
.is_err(),
"unknown fields must be rejected (deny_unknown_fields)"
);
}
#[test]
fn host_status_and_grants_round_trip() {
let host = MobHostStatus {
host_id: WireHostRef("host-b-peer".to_string()),
endpoint: Some("tcp://10.0.0.2:7100".to_string()),
bind_phase: WireHostBindPhase::Bound,
authority_epoch: Some(4),
capabilities: Some(WireHostCapabilityFlags {
protocol_min: 2,
protocol_max: 4,
engine_version: "0.7.22".to_string(),
durable_sessions: true,
autonomous_members: true,
hard_cancel_member: false,
tracked_input_cancel: false,
memory_store: false,
mcp: true,
resolvable_providers: std::collections::BTreeSet::from(["anthropic".to_string()]),
approval_forwarding: false,
live_endpoint: None,
}),
control_reachability: Some(WireReachability::Reachable),
last_seen_ms: Some(250),
freshness_reason: None,
materialized_member_count: 2,
};
let result = MobHostsResult { hosts: vec![host] };
let value = serde_json::to_value(&result).expect("serialize hosts");
assert_eq!(value["hosts"][0]["bind_phase"], serde_json::json!("bound"));
let decoded: MobHostsResult = serde_json::from_value(value).expect("decode hosts");
assert_eq!(decoded, result);
let requested = MobHostStatus {
host_id: WireHostRef("host-c-peer".to_string()),
endpoint: None,
bind_phase: WireHostBindPhase::Requested,
authority_epoch: None,
capabilities: None,
control_reachability: None,
last_seen_ms: None,
freshness_reason: None,
materialized_member_count: 0,
};
let value = serde_json::to_value(&requested).expect("serialize requested host");
assert_eq!(value["bind_phase"], serde_json::json!("requested"));
assert!(value.get("endpoint").is_none());
assert!(value.get("authority_epoch").is_none());
assert!(value.get("capabilities").is_none());
let decoded: MobHostStatus = serde_json::from_value(value).expect("decode requested host");
assert_eq!(decoded, requested);
let record = WireGrantRecord {
principal: "console:luka".to_string(),
scopes: vec![WireControlScope::List, WireControlScope::Live],
expires_at_ms: None,
};
let value = serde_json::to_value(&record).expect("serialize grant");
assert!(
value.get("expires_at_ms").is_none(),
"absent expiry must be omitted"
);
let decoded: WireGrantRecord = serde_json::from_value(value).expect("decode grant");
assert_eq!(decoded, record);
}
#[test]
fn member_history_result_round_trips_with_provenance() {
let result = MobMemberHistoryResult {
page: WireMemberHistoryPageBody {
from_index: 5,
messages: Vec::new(),
message_count: 12,
next_index: Some(10),
complete: false,
},
generation: 2,
placement: Some(WireHostRef("host-b-peer".to_string())),
provenance: WireProjectionProvenance::HostClaimed,
};
let value = serde_json::to_value(&result).expect("serialize history result");
assert_eq!(value["provenance"], serde_json::json!("host_claimed"));
assert_eq!(value["page"]["message_count"], serde_json::json!(12));
let decoded: MobMemberHistoryResult =
serde_json::from_value(value.clone()).expect("decode history result");
let reencoded = serde_json::to_value(&decoded).expect("reserialize history result");
assert_eq!(value, reencoded);
}
#[test]
fn member_history_projection_rejects_non_advancing_page() {
let page = meerkat_core::service::SessionHistoryPage {
session_id: meerkat_core::SessionId::new(),
message_count: 1,
offset: 0,
limit: Some(1),
has_more: true,
messages: Vec::new(),
};
let error = WireMemberHistoryPageBody::try_from_history_page(&page)
.expect_err("a page that claims more rows must advance its cursor");
assert!(matches!(
error,
crate::wire::error::WireConversionError::MemberHistoryPage { debug }
if debug.contains("serves none")
));
}
#[cfg(target_pointer_width = "64")]
#[test]
fn member_history_projection_rejects_exhausted_cursor() {
let page = meerkat_core::service::SessionHistoryPage {
session_id: meerkat_core::SessionId::new(),
message_count: usize::MAX,
offset: usize::MAX,
limit: Some(2),
has_more: true,
messages: vec![
meerkat_core::types::Message::User(meerkat_core::types::UserMessage::text("first")),
meerkat_core::types::Message::User(meerkat_core::types::UserMessage::text(
"second",
)),
],
};
let error = WireMemberHistoryPageBody::try_from_history_page(&page)
.expect_err("MAX has no representable member-history successor");
assert!(matches!(
error,
crate::wire::error::WireConversionError::MemberHistoryPage { debug }
if debug.contains("exhausts the u64 cursor domain")
));
}
#[test]
fn hard_cancel_params_round_trip_and_reject_unknown_fields() {
let params = MobHardCancelParams {
mob_id: "mob-1".to_string(),
agent_identity: "worker".to_string(),
reason: "operator interrupt".to_string(),
};
let value = serde_json::to_value(¶ms).expect("serialize hard-cancel params");
let decoded: MobHardCancelParams =
serde_json::from_value(value).expect("decode hard-cancel params");
assert_eq!(decoded, params);
serde_json::from_value::<MobHardCancelParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
"reason": "x",
"force": true,
}))
.expect_err("unknown field must be rejected");
serde_json::from_value::<MobHardCancelParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
}))
.expect_err("missing reason must be rejected");
let result = MobHardCancelResult { cancelled: true };
let value = serde_json::to_value(&result).expect("serialize hard-cancel result");
assert_eq!(value, serde_json::json!({ "cancelled": true }));
}
#[test]
fn member_live_status_params_keep_the_discovery_read() {
let discovery: MobMemberLiveStatusParams = serde_json::from_value(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
}))
.expect("discovery status params parse");
assert_eq!(discovery.channel_id, None);
let value = serde_json::to_value(&discovery).expect("serialize discovery params");
assert!(
value.get("channel_id").is_none(),
"absent channel_id must be omitted"
);
let named: MobMemberLiveStatusParams = serde_json::from_value(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
"channel_id": "chan-7",
}))
.expect("named status params parse");
assert_eq!(named.channel_id.as_deref(), Some("chan-7"));
serde_json::from_value::<MobMemberLiveStatusParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
"chan": "chan-7",
}))
.expect_err("unknown field must be rejected");
serde_json::from_value::<MobMemberLiveChannelParams>(serde_json::json!({
"mob_id": "mob-1",
"agent_identity": "worker",
}))
.expect_err("close without channel_id must be rejected");
}
}