use super::flow_frame_engine::FlowFrameLoopStorePlan;
use super::terminalization::{TerminalizationOutcome, TerminalizationTarget};
use super::*;
use crate::machines::mob_machine as mob_dsl;
use crate::run::{MobMachineFlowRunCommand, MobRun, flow_run};
#[cfg(target_arch = "wasm32")]
use crate::tokio;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum MobState {
Creating = 0,
Running = 1,
Stopped = 2,
Completed = 3,
Destroyed = 4,
}
impl MobState {
pub fn as_str(self) -> &'static str {
match self {
Self::Creating => "Creating",
Self::Running => "Running",
Self::Stopped => "Stopped",
Self::Completed => "Completed",
Self::Destroyed => "Destroyed",
}
}
}
impl std::fmt::Display for MobState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MobOrchestratorSnapshot {
pub phase: MobState,
pub coordinator_bound: bool,
pub pending_spawn_count: u32,
pub active_flow_count: u32,
pub topology_revision: u32,
pub supervisor_active: bool,
}
impl Default for MobOrchestratorSnapshot {
fn default() -> Self {
Self {
phase: MobState::Creating,
coordinator_bound: false,
pending_spawn_count: 0,
active_flow_count: 0,
topology_revision: 0,
supervisor_active: false,
}
}
}
#[cfg(test)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct MobLifecycleSnapshot {
pub phase: MobState,
pub active_run_count: u32,
pub cleanup_pending: bool,
}
#[cfg(test)]
impl Default for MobLifecycleSnapshot {
fn default() -> Self {
Self {
phase: MobState::Running,
active_run_count: 0,
cleanup_pending: false,
}
}
}
#[cfg(test)]
#[derive(Debug, Clone, Default)]
pub(crate) struct MobDslT2Snapshot {
pub destroy_admitted: bool,
pub flow_authority_schema_version: u64,
pub owner_bridge_session_id: Option<crate::machines::mob_machine::SessionId>,
pub owner_bridge_destroy_on_archive: bool,
pub implicit_delegation_mob: bool,
pub supervisor_authority_peer_id: Option<crate::machines::mob_machine::PeerId>,
pub supervisor_authority_signing_key: Option<crate::machines::mob_machine::PeerSigningKey>,
pub supervisor_authority_epoch: Option<u64>,
pub supervisor_authority_protocol_version:
Option<crate::machines::mob_machine::SupervisorProtocolVersion>,
pub supervisor_pending_authority_peer_id: Option<crate::machines::mob_machine::PeerId>,
pub supervisor_pending_authority_signing_key:
Option<crate::machines::mob_machine::PeerSigningKey>,
pub supervisor_pending_authority_epoch: Option<u64>,
pub supervisor_pending_authority_protocol_version:
Option<crate::machines::mob_machine::SupervisorProtocolVersion>,
pub supervisor_pending_authority_operation_id: Option<String>,
pub supervisor_pending_authority_member_target_names:
std::collections::BTreeMap<crate::machines::mob_machine::PeerId, String>,
pub supervisor_pending_authority_member_target_addresses:
std::collections::BTreeMap<crate::machines::mob_machine::PeerId, String>,
pub supervisor_pending_authority_accepted_peer_ids:
std::collections::BTreeSet<crate::machines::mob_machine::PeerId>,
pub pending_recipient_trust: std::collections::BTreeSet<crate::machines::mob_machine::PeerId>,
pub member_state_markers: std::collections::BTreeMap<
crate::machines::mob_machine::AgentRuntimeId,
crate::machines::mob_machine::MobMemberState,
>,
pub wiring_edges: std::collections::BTreeSet<crate::machines::mob_machine::WiringEdge>,
pub external_peer_edges:
std::collections::BTreeSet<crate::machines::mob_machine::ExternalPeerEdge>,
pub external_peer_edges_by_key: std::collections::BTreeMap<
crate::machines::mob_machine::ExternalPeerKey,
crate::machines::mob_machine::ExternalPeerEdge,
>,
pub identity_to_runtime: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::AgentRuntimeId,
>,
pub identity_runtime_generations: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::Generation,
>,
pub identity_runtime_fence_tokens: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::FenceToken,
>,
pub member_profile_names:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub member_runtime_modes: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::SpawnPolicyRuntimeMode,
>,
pub member_peer_ids: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::PeerId,
>,
pub member_peer_endpoints: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::MemberPeerEndpoint,
>,
pub member_prior_peer_endpoints: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
std::collections::BTreeSet<crate::machines::mob_machine::MemberPeerEndpoint>,
>,
pub member_restore_failures:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub member_restore_failure_codes:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub runtime_retire_refusal_codes:
std::collections::BTreeMap<crate::machines::mob_machine::AgentRuntimeId, String>,
pub runtime_retire_refusal_reasons:
std::collections::BTreeMap<crate::machines::mob_machine::AgentRuntimeId, String>,
pub runtime_retire_pending_sessions: std::collections::BTreeMap<
crate::machines::mob_machine::AgentRuntimeId,
crate::machines::mob_machine::SessionId,
>,
pub remote_runtime_retired_ids:
std::collections::BTreeSet<crate::machines::mob_machine::AgentRuntimeId>,
pub remote_supervisor_revoked_ids:
std::collections::BTreeSet<crate::machines::mob_machine::AgentRuntimeId>,
pub member_revival_pending:
std::collections::BTreeSet<crate::machines::mob_machine::AgentIdentity>,
pub member_session_bindings: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::SessionId,
>,
pub spawn_exec_phase: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::SpawnExecPhase,
>,
pub pending_spawn_sessions: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::SessionId,
>,
pub pending_session_ingress_detach_runtime_ids:
std::collections::BTreeSet<crate::machines::mob_machine::AgentRuntimeId>,
pub topology_epoch: u64,
pub spawn_policy_enabled: bool,
pub spawn_policy_revision: u64,
pub spawn_policy_resolution_revision:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, u64>,
pub spawn_policy_resolution_profiles:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_policy_resolution_runtime_modes: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
Option<crate::machines::mob_machine::SpawnPolicyRuntimeMode>,
>,
pub spawn_policy_resolution_absent:
std::collections::BTreeSet<crate::machines::mob_machine::AgentIdentity>,
pub spawn_profile_authority_profile_names:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_profile_authority_models:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_profile_authority_material_digests:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_profile_authority_tool_config_digests:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_profile_authority_skills_digests:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, String>,
pub spawn_profile_authority_provider_params_digests:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, Option<String>>,
pub spawn_profile_authority_output_schema_digests:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, Option<String>>,
pub spawn_profile_authority_external_addressable:
std::collections::BTreeMap<crate::machines::mob_machine::AgentIdentity, bool>,
pub orphan_budget: u64,
pub topology_default_policy: crate::machines::mob_machine::PolicyDecision,
pub external_member_rebind_capability: std::collections::BTreeMap<
crate::machines::mob_machine::AgentIdentity,
crate::machines::mob_machine::ExternalMemberRebindCapability,
>,
pub desired_members: std::collections::BTreeSet<crate::machines::mob_machine::AgentIdentity>,
pub members_to_spawn: std::collections::BTreeSet<crate::machines::mob_machine::AgentIdentity>,
pub members_to_retire: std::collections::BTreeSet<crate::machines::mob_machine::AgentIdentity>,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct MobStartupKickoffSnapshot {
pub pending_kickoff_member_ids: std::collections::BTreeSet<String>,
pub ready_runtime_ids: std::collections::BTreeSet<String>,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct MobMemberMachineProjection {
pub runtime_id: Option<crate::machines::mob_machine::AgentRuntimeId>,
pub state_marker: Option<crate::machines::mob_machine::MobMemberState>,
pub live_runtime: bool,
pub bound_session_id: Option<crate::machines::mob_machine::SessionId>,
}
pub(super) struct SubmitWorkPayload {
pub runtime_id: AgentRuntimeId,
pub fence_token: FenceToken,
pub work_ref: WorkRef,
pub content: ContentInput,
pub origin: WorkOrigin,
pub injected_context: Vec<ContentInput>,
pub interaction_id: Option<meerkat_core::interaction::InteractionId>,
pub handling_mode: meerkat_core::types::HandlingMode,
pub render_metadata: Option<meerkat_core::types::RenderMetadata>,
pub ack_mode: crate::mob_machine::SubmitWorkAckMode,
}
pub(super) enum MobCommand {
Spawn {
spec: Box<super::handle::SpawnMemberSpec>,
spawn_source: super::handle::SpawnSource,
owner_bridge_session_id: Option<SessionId>,
ops_registry: Option<Arc<dyn meerkat_core::ops_lifecycle::OpsLifecycleRegistry>>,
reply_tx: oneshot::Sender<Result<super::handle::MemberSpawnReceipt, MobError>>,
},
SpawnProvisioned {
spawn_ticket: u64,
result: Result<super::handle::MemberSpawnReceipt, MobError>,
},
Retire {
agent_identity: AgentIdentity,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
Respawn {
agent_identity: AgentIdentity,
initial_message: Option<ContentInput>,
reply_tx: oneshot::Sender<
Result<super::handle::MemberRespawnReceipt, super::handle::MobRespawnError>,
>,
},
RetireAll {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
SubmitWork {
payload: Box<SubmitWorkPayload>,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
SendPeerMessage {
from: AgentIdentity,
to: AgentIdentity,
content: ContentInput,
handling_mode: meerkat_core::types::HandlingMode,
reply_tx: oneshot::Sender<Result<meerkat_core::comms::SendReceipt, MobError>>,
},
DeclareMemberOutboundTaint {
identity: AgentIdentity,
taint: Option<meerkat_core::comms::SenderContentTaint>,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
CancelAllWork {
runtime_id: AgentRuntimeId,
fence_token: FenceToken,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
#[cfg(feature = "runtime-adapter")]
KickoffOutcomeResolved {
agent_identity: AgentIdentity,
outcome: Result<
meerkat_runtime::completion::CompletionOutcome,
meerkat_runtime::completion::CompletionWaitError,
>,
ack_tx: oneshot::Sender<()>,
},
RunFlow {
flow_id: FlowId,
activation_params: serde_json::Value,
scoped_event_tx: Option<tokio::sync::mpsc::Sender<meerkat_core::ScopedAgentEvent>>,
reply_tx: oneshot::Sender<Result<RunId, MobError>>,
},
PreviewRunFlowAdmission {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
CancelFlow {
run_id: RunId,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
FlowStatus {
run_id: RunId,
reply_tx: oneshot::Sender<Result<Option<MobRun>, MobError>>,
},
CommitFlowRunCommand {
run_id: RunId,
command: Box<MobMachineFlowRunCommand>,
context: &'static str,
reply_tx: oneshot::Sender<Result<Option<Vec<flow_run::Effect>>, MobError>>,
},
CommitFlowTerminalization {
run_id: RunId,
flow_id: FlowId,
target: TerminalizationTarget,
command: Box<MobMachineFlowRunCommand>,
context: &'static str,
reply_tx: oneshot::Sender<Result<TerminalizationOutcome, MobError>>,
},
CommitFlowFrameStorePlan {
run_id: RunId,
plan: Box<FlowFrameLoopStorePlan>,
reply_tx: oneshot::Sender<Result<bool, MobError>>,
},
ProjectMachineInput {
input: Box<mob_dsl::MobMachineInput>,
reply_tx: oneshot::Sender<Result<mob_dsl::MobMachineState, MobError>>,
},
ApplyMachineInputEffects {
input: Box<mob_dsl::MobMachineInput>,
reply_tx: oneshot::Sender<Result<Vec<mob_dsl::MobMachineEffect>, MobError>>,
},
PreviewMachineInput {
input: Box<mob_dsl::MobMachineInput>,
reply_tx: oneshot::Sender<Result<mob_dsl::MobMachineState, MobError>>,
},
QueryMachineState {
reply_tx: oneshot::Sender<mob_dsl::MobMachineState>,
},
#[cfg(test)]
AuthorizeMemberTrustCleanupForTest {
edge: mob_dsl::WiringEdge,
reply_tx: oneshot::Sender<
Result<
crate::generated::protocol_mob_member_trust_unwiring::MobMemberTrustUnwiringObligation,
MobError,
>,
>,
},
ApplyExternalPeerReciprocalTrust {
key: mob_dsl::ExternalPeerKey,
target_comms: std::sync::Arc<dyn meerkat_core::agent::CommsRuntime>,
peer: meerkat_core::comms::TrustedPeerDescriptor,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
ProjectMachineSignal {
signal: mob_dsl::MobMachineSignal,
},
RecordMissingMemberBridgeSession {
agent_identity: crate::ids::AgentIdentity,
bridge_session_id: SessionId,
},
FlowFinished {
run_id: RunId,
},
FlowCanceledCleanup {
run_id: RunId,
terminalized: bool,
},
#[cfg(test)]
FlowTrackerCounts {
reply_tx: oneshot::Sender<(usize, usize)>,
},
#[cfg(test)]
OrchestratorSnapshot {
reply_tx: oneshot::Sender<MobOrchestratorSnapshot>,
},
#[cfg(test)]
LifecycleSnapshot {
reply_tx: oneshot::Sender<MobLifecycleSnapshot>,
},
#[cfg(test)]
LifecycleNotificationBurst {
count: usize,
message: String,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
#[cfg(test)]
DslT2Snapshot {
reply_tx: oneshot::Sender<MobDslT2Snapshot>,
},
StartupKickoffSnapshot {
reply_tx: oneshot::Sender<MobStartupKickoffSnapshot>,
},
ProjectMemberList {
include_retiring: bool,
reply_tx: oneshot::Sender<Vec<super::MobMemberListEntry>>,
},
ProjectMemberStatus {
agent_identity: crate::ids::AgentIdentity,
reply_tx: oneshot::Sender<Result<super::MobMemberSnapshot, crate::MobError>>,
},
MemberMachineProjection {
agent_identity: crate::ids::AgentIdentity,
reply_tx: oneshot::Sender<MobMemberMachineProjection>,
},
Stop {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
ResumeLifecycle {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
Complete {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
Destroy {
reply_tx: oneshot::Sender<
Result<super::handle::MobDestroyReport, super::handle::MobDestroyError>,
>,
},
Reset {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
RotateSupervisor {
reply_tx: oneshot::Sender<Result<super::handle::SupervisorRotationReport, MobError>>,
},
PollEvents {
after_cursor: u64,
limit: usize,
reply_tx: oneshot::Sender<Result<Vec<crate::event::MobEvent>, MobError>>,
},
ReplayAllEvents {
reply_tx: oneshot::Sender<Result<Vec<crate::event::MobEvent>, MobError>>,
},
RecordOperatorActionProvenance {
tool_name: String,
authority_context: meerkat_core::service::MobToolAuthorityContext,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
ForceCancel {
agent_identity: AgentIdentity,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
Wire {
local: AgentIdentity,
target: super::handle::PeerTarget,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
WireMembersBatch {
edges: Vec<(AgentIdentity, AgentIdentity)>,
reply_tx: oneshot::Sender<Result<super::handle::MobWireMembersBatchReport, MobError>>,
},
Unwire {
local: AgentIdentity,
target: super::handle::PeerTarget,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
SetSpawnPolicy {
policy: Option<Arc<dyn super::spawn_policy::SpawnPolicy>>,
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
Shutdown {
reply_tx: oneshot::Sender<Result<(), MobError>>,
},
QueryPhase {
reply_tx: oneshot::Sender<MobState>,
},
}
impl MobCommand {
pub(super) fn kind(&self) -> &'static str {
match self {
Self::Spawn { .. } => "Spawn",
Self::SpawnProvisioned { .. } => "SpawnProvisioned",
Self::Retire { .. } => "Retire",
Self::Respawn { .. } => "Respawn",
Self::RetireAll { .. } => "RetireAll",
Self::SubmitWork { .. } => "SubmitWork",
Self::SendPeerMessage { .. } => "SendPeerMessage",
Self::DeclareMemberOutboundTaint { .. } => "DeclareMemberOutboundTaint",
Self::CancelAllWork { .. } => "CancelAllWork",
#[cfg(feature = "runtime-adapter")]
Self::KickoffOutcomeResolved { .. } => "KickoffOutcomeResolved",
Self::RunFlow { .. } => "RunFlow",
Self::PreviewRunFlowAdmission { .. } => "PreviewRunFlowAdmission",
Self::CancelFlow { .. } => "CancelFlow",
Self::FlowStatus { .. } => "FlowStatus",
Self::CommitFlowRunCommand { .. } => "CommitFlowRunCommand",
Self::CommitFlowTerminalization { .. } => "CommitFlowTerminalization",
Self::CommitFlowFrameStorePlan { .. } => "CommitFlowFrameStorePlan",
Self::ProjectMachineInput { .. } => "ProjectMachineInput",
Self::ApplyMachineInputEffects { .. } => "ApplyMachineInputEffects",
Self::PreviewMachineInput { .. } => "PreviewMachineInput",
Self::QueryMachineState { .. } => "QueryMachineState",
#[cfg(test)]
Self::AuthorizeMemberTrustCleanupForTest { .. } => "AuthorizeMemberTrustCleanupForTest",
Self::ApplyExternalPeerReciprocalTrust { .. } => "ApplyExternalPeerReciprocalTrust",
Self::ProjectMachineSignal { .. } => "ProjectMachineSignal",
Self::RecordMissingMemberBridgeSession { .. } => "RecordMissingMemberBridgeSession",
Self::FlowFinished { .. } => "FlowFinished",
Self::FlowCanceledCleanup { .. } => "FlowCanceledCleanup",
#[cfg(test)]
Self::FlowTrackerCounts { .. } => "FlowTrackerCounts",
#[cfg(test)]
Self::OrchestratorSnapshot { .. } => "OrchestratorSnapshot",
#[cfg(test)]
Self::LifecycleSnapshot { .. } => "LifecycleSnapshot",
#[cfg(test)]
Self::LifecycleNotificationBurst { .. } => "LifecycleNotificationBurst",
#[cfg(test)]
Self::DslT2Snapshot { .. } => "DslT2Snapshot",
Self::StartupKickoffSnapshot { .. } => "StartupKickoffSnapshot",
Self::ProjectMemberList { .. } => "ProjectMemberList",
Self::ProjectMemberStatus { .. } => "ProjectMemberStatus",
Self::MemberMachineProjection { .. } => "MemberMachineProjection",
Self::Stop { .. } => "Stop",
Self::ResumeLifecycle { .. } => "ResumeLifecycle",
Self::Complete { .. } => "Complete",
Self::Destroy { .. } => "Destroy",
Self::Reset { .. } => "Reset",
Self::RotateSupervisor { .. } => "RotateSupervisor",
Self::PollEvents { .. } => "PollEvents",
Self::ReplayAllEvents { .. } => "ReplayAllEvents",
Self::RecordOperatorActionProvenance { .. } => "RecordOperatorActionProvenance",
Self::ForceCancel { .. } => "ForceCancel",
Self::Wire { .. } => "Wire",
Self::WireMembersBatch { .. } => "WireMembersBatch",
Self::Unwire { .. } => "Unwire",
Self::SetSpawnPolicy { .. } => "SetSpawnPolicy",
Self::Shutdown { .. } => "Shutdown",
Self::QueryPhase { .. } => "QueryPhase",
}
}
}