#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use meerkat_client::{LlmClient, LlmError, LlmEvent, LlmRequest};
use meerkat_core::skills::{SkillKey, SkillName, SkillRef, SourceUuid};
use meerkat_core::types::{HandlingMode, TranscriptUserRole};
use meerkat_core::{ContentInput, Message, Provider, SessionId};
use meerkat_mob::{MobDefinition, MobDeliveryIdentity, WorkOrigin, WorkSpec};
use serde_json::json;
use crate::console_aggregator::{
AllowAllConsoleVisibilityPolicy, ConsoleFrame, ConsoleFrameSourceKind, ConsoleFrameStatus,
ConsoleInteractionAccepted, ConsoleSendRequest, ConsoleTimelineQuery, MobKitConsoleAggregator,
};
use crate::identity_first::agent_memory::{
AgentMemoryConfig, AgentMemoryError, AgentMemoryPerTurnInjection, AgentMemoryProvider,
AgentMemoryRecallRequest, AgentMemoryRecord, AgentMemoryRuntimeInjector,
};
use crate::identity_first::bridge::{HostHumanInputError, MobSessionBridge};
use crate::identity_first::orchestrator::restore_flow;
use crate::identity_first::*;
use crate::unified_runtime::{ConsoleEventStore, UnifiedRuntime, UnifiedRuntimeBuilder};
const WAIT: Duration = Duration::from_secs(30);
const HUMAN: &str = "CURRENT_VALUE is cobalt; remember my exact choice.";
#[derive(Clone)]
struct RecordingClient {
requests: Arc<Mutex<Vec<LlmRequest>>>,
gate: tokio::sync::watch::Sender<bool>,
}
impl RecordingClient {
fn new() -> Self {
Self {
requests: Arc::default(),
gate: tokio::sync::watch::channel(true).0,
}
}
}
#[async_trait::async_trait]
impl LlmClient for RecordingClient {
fn project_replay_messages(&self, messages: &[Message]) -> Result<Vec<Message>, LlmError> {
Ok(messages.to_vec())
}
fn stream<'a>(&'a self, request: &'a LlmRequest) -> meerkat_client::types::LlmStream<'a> {
self.requests.lock().unwrap().push(request.clone());
let mut gate = self.gate.subscribe();
Box::pin(async_stream::try_stream! {
while !*gate.borrow_and_update() {
gate.changed().await.unwrap();
}
yield LlmEvent::TextDelta { delta: "Your choice is cobalt.".to_string(), meta: None };
yield LlmEvent::UsageUpdate {
usage: meerkat_core::TurnUsage::host_declared(
Provider::OpenAI, &request.model, meerkat_core::Usage::default(),
),
};
yield LlmEvent::Done {
outcome: meerkat_client::LlmDoneOutcome::Success {
stop_reason: meerkat_core::StopReason::EndTurn,
},
};
})
}
fn provider(&self) -> Provider {
Provider::OpenAI
}
async fn health_check(&self) -> Result<(), LlmError> {
Ok(())
}
}
struct Harness {
unified: UnifiedRuntime,
identity_runtime: Arc<IdentityRuntime>,
bridge: Arc<MobSessionBridge>,
aggregator: MobKitConsoleAggregator,
identity: AgentIdentity,
runtime_id: AgentRuntimeId,
session_id: SessionId,
client: RecordingClient,
_state: tempfile::TempDir,
}
impl Harness {
async fn new(budget: Duration) -> Self {
Self::with_kickoff(budget, false).await
}
async fn with_kickoff(budget: Duration, kickoff: bool) -> Self {
Self::build(budget, kickoff, None).await
}
async fn with_meerkat_config(config: meerkat::Config) -> Self {
Self::build(WAIT, false, Some(config)).await
}
async fn build(budget: Duration, kickoff: bool, config: Option<meerkat::Config>) -> Self {
let state = tempfile::tempdir_in(std::env::current_dir().unwrap()).unwrap();
let client = RecordingClient::new();
let builder = match config {
Some(config) => UnifiedRuntimeBuilder::default().meerkat_config(config),
None => UnifiedRuntimeBuilder::default(),
};
let unified = builder
.definition(
MobDefinition::from_toml(&format!(
r#"
[mob]
id = "console-human-{}"
[profiles.human]
model = "gpt-5.5"
external_addressable = true
runtime_mode = "autonomous_host"
[profiles.human.tools]
comms = true
"#,
uuid::Uuid::new_v4()
))
.unwrap(),
)
.persistent_state(state.path())
.comms(true)
.default_llm_client(Arc::new(client.clone()))
.build()
.await
.unwrap();
let bridge = Arc::new(
MobSessionBridge::with_session_service(
unified.mob_handle(),
unified.mob_runtime().session_service().cloned().unwrap(),
)
.with_actor_admission_budget(budget),
);
let identity_runtime = Arc::new(
IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory().unwrap()),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "human-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(bridge.clone()),
default_timeout: None,
})
.with_runtime_services(AgentRuntimeServices::new(unified.mob_handle())),
);
let identity = AgentIdentity::parse("human").unwrap();
let spec = DurableAgentSpec {
identity: identity.clone(),
profile: "human".into(),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: kickoff.then(|| "fixture kickoff".into()),
runtime_mode_override: Some(meerkat_mob::MobRuntimeMode::AutonomousHost),
backend: None,
binding: None,
placement: None,
};
restore_flow(&identity_runtime, &[spec], None, None)
.await
.unwrap();
let status = identity_runtime.status(&identity).await.unwrap();
if kickoff {
let member_id = crate::member_comms_id::mob_member_id(identity.as_str());
tokio::time::timeout(WAIT, async {
loop {
let entry = unified.mob_handle().get_member(&member_id).await.unwrap();
if entry
.as_ref()
.and_then(|entry| entry.kickoff.as_ref())
.is_some_and(|kickoff| {
kickoff.phase == meerkat_mob::MobMemberKickoffPhase::Started
})
{
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.unwrap();
}
let aggregator = MobKitConsoleAggregator::in_memory();
aggregator.register_runtime_handles_with_policy(
"human",
"",
unified.mob_runtime().clone(),
Some(identity_runtime.clone()),
ConsoleEventStore::new(),
Arc::new(AllowAllConsoleVisibilityPolicy),
);
Self {
unified,
identity_runtime,
bridge,
aggregator,
identity,
runtime_id: status.agent_runtime_id.unwrap(),
session_id: status.session_id.unwrap(),
client,
_state: state,
}
}
fn request(&self, key: &str, content: &str) -> ConsoleSendRequest {
ConsoleSendRequest {
identity: self.identity.to_string(),
content: json!(content),
origin: "console".to_string(),
idempotency_key: key.to_string(),
handling_mode: None,
origin_kind: None,
skill_refs: Vec::new(),
}
}
async fn send(
&self,
request: ConsoleSendRequest,
) -> Result<ConsoleInteractionAccepted, crate::console_aggregator::ConsoleSendError> {
super::console_send_identity_first(
&self.aggregator,
self.identity_runtime.clone(),
None,
request,
)
.await
}
async fn wait_status(
&self,
accepted: &ConsoleInteractionAccepted,
expected: ConsoleFrameStatus,
) {
tokio::time::timeout(WAIT, async {
loop {
let page = self
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(accepted.identity.clone()),
..Default::default()
})
.await
.unwrap();
if page
.frames
.iter()
.any(|frame| frame.id == accepted.input_frame_id && frame.status == expected)
{
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("console dispatch status");
}
async fn finish_human(&self, accepted: &ConsoleInteractionAccepted, content: &str) {
self.finish_human_with_mode(accepted, content, HandlingMode::Queue)
.await;
}
async fn finish_human_with_mode(
&self,
accepted: &ConsoleInteractionAccepted,
content: &str,
mode: HandlingMode,
) {
self.finish_human_with_context(accepted, content, mode, Vec::new())
.await;
}
async fn finish_human_with_context(
&self,
accepted: &ConsoleInteractionAccepted,
content: &str,
mode: HandlingMode,
injected_context: Vec<ContentInput>,
) {
self.wait_status(accepted, ConsoleFrameStatus::Delivered)
.await;
let handle = self.unified.mob_handle();
let member = handle
.get_member(&crate::member_comms_id::mob_member_id(&accepted.identity))
.await
.unwrap()
.unwrap();
let interaction = accepted.interaction_id.parse::<uuid::Uuid>().unwrap();
let spec = WorkSpec::new(content, WorkOrigin::Internal)
.with_injected_context(injected_context)
.with_interaction_id(meerkat_core::interaction::InteractionId(interaction));
let turn = handle
.start_host_human_input_bounded(
member.agent_runtime_id,
member.fence_token,
spec,
mode,
MobDeliveryIdentity::new(&accepted.input_frame_id, interaction.to_string())
.unwrap(),
std::time::Instant::now() + WAIT,
)
.await
.unwrap();
tokio::time::timeout(WAIT, turn.wait())
.await
.unwrap()
.unwrap();
}
async fn finish_human_with_skills(
&self,
accepted: &ConsoleInteractionAccepted,
content: &str,
skills: Vec<SkillKey>,
) {
self.wait_status(accepted, ConsoleFrameStatus::Delivered)
.await;
let handle = self.unified.mob_handle();
let member = handle
.get_member(&crate::member_comms_id::mob_member_id(&accepted.identity))
.await
.unwrap()
.unwrap();
let interaction = accepted.interaction_id.parse::<uuid::Uuid>().unwrap();
let spec = WorkSpec::new(content, WorkOrigin::Internal)
.with_interaction_id(meerkat_core::interaction::InteractionId(interaction));
let turn = handle
.start_host_human_input_with_options_bounded(
member.agent_runtime_id,
member.fence_token,
spec,
HandlingMode::Queue,
meerkat_mob::MemberTurnOptions::new().with_skill_references(skills),
MobDeliveryIdentity::new(&accepted.input_frame_id, interaction.to_string())
.unwrap(),
std::time::Instant::now() + WAIT,
)
.await
.unwrap();
tokio::time::timeout(WAIT, turn.wait())
.await
.unwrap()
.unwrap();
}
async fn user_input_frames(&self, identity: &str) -> Vec<(String, ConsoleFrameStatus)> {
self.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(identity.to_string()),
..Default::default()
})
.await
.unwrap()
.frames
.into_iter()
.filter(|frame| frame.kind == "user_input")
.map(|frame| (frame.id, frame.status))
.collect()
}
async fn history(&self) -> Vec<Message> {
self.unified
.mob_runtime()
.read_session_history(&self.session_id.to_string(), 0, None)
.await
.unwrap()
.messages
}
async fn stop(self) {
match self.unified.stop_mob_for_teardown().await {
crate::unified_runtime::MobStopOutcome::Stopped
| crate::unified_runtime::MobStopOutcome::ProceededWithoutInterrupt { .. } => {}
crate::unified_runtime::MobStopOutcome::Failed(error) => {
panic!("mob teardown failed: {error}");
}
}
}
}
fn assert_canonical_human_frame(
frame: &ConsoleFrame,
accepted: &ConsoleInteractionAccepted,
session_id: &SessionId,
canonical: &[Message],
) {
let matching: Vec<_> = canonical
.iter()
.enumerate()
.filter(|(_, message)| {
matches!(message,
Message::User(user)
if user.transcript_role == TranscriptUserRole::Conversational
&& user.identity.interaction_id.as_ref().is_some_and(|id|
id.0.to_string() == accepted.interaction_id))
})
.collect();
assert_eq!(
matching.len(),
1,
"one committed human input per interaction"
);
let (offset, message) = matching[0];
let Message::User(user) = message else {
unreachable!()
};
assert_eq!(user.text_content(), HUMAN);
assert_eq!(frame.kind, "user_input");
assert_eq!(frame.source.kind, ConsoleFrameSourceKind::SessionHistory);
assert_eq!(
frame.session_id.as_deref(),
Some(session_id.to_string().as_str())
);
assert_eq!(
frame.interaction_id.as_deref(),
Some(accepted.interaction_id.as_str())
);
assert_eq!(
frame.run_id,
user.identity.run_id.as_ref().map(|id| id.0.to_string())
);
assert!(
frame.run_id.is_some(),
"committed human input retains its runtime run"
);
assert_eq!(
frame.source.source_cursor,
Some(format!("{session_id}:{offset}"))
);
assert_eq!(
frame.payload["message"],
serde_json::to_value(message).unwrap()
);
assert_eq!(
frame.payload["content"],
serde_json::to_value(&user.content).unwrap()
);
}
fn assert_human_input_source_pair(
frames: &[ConsoleFrame],
accepted: &ConsoleInteractionAccepted,
session_id: &SessionId,
canonical: &[Message],
) -> ConsoleFrame {
let inputs: Vec<_> = frames
.iter()
.filter(|frame| frame.kind == "user_input")
.collect();
assert_eq!(
inputs.len(),
2,
"exactly one admission and one canonical source: {frames:#?}"
);
let reservations: Vec<_> = inputs
.iter()
.copied()
.filter(|frame| frame.source.kind == ConsoleFrameSourceKind::Send)
.collect();
assert_eq!(reservations.len(), 1);
let reserved = reservations[0];
assert_eq!(reserved.id, accepted.input_frame_id);
assert_eq!(reserved.cursor, accepted.cursor);
assert_eq!(
reserved.session_id.as_deref(),
Some(session_id.to_string().as_str())
);
assert_eq!(
reserved.interaction_id.as_deref(),
Some(accepted.interaction_id.as_str())
);
assert_eq!(reserved.payload["content"], HUMAN);
let history: Vec<_> = inputs
.iter()
.copied()
.filter(|frame| frame.source.kind == ConsoleFrameSourceKind::SessionHistory)
.collect();
assert_eq!(history.len(), 1);
assert_canonical_human_frame(history[0], accepted, session_id, canonical);
history[0].clone()
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_is_canonical_once_and_generic_identity_work_stays_external_event() {
let h = Harness::new(WAIT).await;
let request = h.request("canonical", HUMAN);
let accepted = h.send(request.clone()).await.unwrap();
h.finish_human(&accepted, HUMAN).await;
let history = h.history().await;
let interaction =
meerkat_core::interaction::InteractionId(accepted.interaction_id.parse().unwrap());
let users: Vec<_> = history
.iter()
.filter_map(|message| match message {
Message::User(user) if user.identity.interaction_id == Some(interaction) => Some(user),
_ => None,
})
.collect();
assert_eq!(users.len(), 1);
assert_eq!(users[0].text_content(), HUMAN);
assert_eq!(users[0].transcript_role, TranscriptUserRole::Conversational);
assert!(history.iter().any(|message| matches!(message,
Message::BlockAssistant(answer) if answer.identity.interaction_id == Some(interaction))));
let calls = h.client.requests.lock().unwrap().len();
let replay = h.send(request.clone()).await.unwrap();
assert_eq!(replay.input_frame_id, accepted.input_frame_id);
assert_eq!(h.client.requests.lock().unwrap().len(), calls);
let mut changed = request;
changed.content = json!("changed user intent");
assert!(matches!(
h.send(changed).await,
Err(crate::console_aggregator::ConsoleSendError::IdempotencyConflict(_))
));
h.aggregator.refresh_session_history().await.unwrap();
let page = h
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(h.identity.to_string()),
..Default::default()
})
.await
.unwrap();
assert_human_input_source_pair(&page.frames, &accepted, &h.session_id, &history);
h.identity_runtime
.send(
&h.identity,
&ContentInput::Text("GENERIC_SYSTEM_WORK".to_string()),
)
.await
.unwrap();
let generic = tokio::time::timeout(WAIT, async {
loop {
let messages = h.history().await;
if messages.iter().any(|message| {
matches!(message, Message::SystemNotice(_))
&& serde_json::to_string(message)
.unwrap()
.contains("GENERIC_SYSTEM_WORK")
}) {
break messages;
}
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert!(generic.iter().any(|message| {
matches!(message, Message::SystemNotice(_))
&& serde_json::to_string(message)
.unwrap()
.contains("GENERIC_SYSTEM_WORK")
}));
assert!(!generic.iter().any(|message|
matches!(message, Message::User(user) if user.text_content() == "GENERIC_SYSTEM_WORK")));
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_recall_stays_canonical_context_not_refreshed_or_paged_human_input() {
const MEMORY: &str = "Recalled reference: the cobalt sample is stored in locker seven.";
struct RecallProvider;
#[async_trait::async_trait]
impl AgentMemoryProvider for RecallProvider {
async fn recall(
&self,
_request: AgentMemoryRecallRequest,
) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
Ok(vec![AgentMemoryRecord {
memory_id: "cobalt-location".to_string(),
title: "Cobalt reference".to_string(),
body: MEMORY.to_string(),
tags: Vec::new(),
created_at_ms: 1,
updated_at_ms: 1,
}])
}
}
let h = Harness::new(WAIT).await;
h.identity_runtime
.set_agent_memory(Some(AgentMemoryRuntimeInjector::new(
Arc::new(RecallProvider),
AgentMemoryConfig {
per_turn_injection: AgentMemoryPerTurnInjection::Budgeted,
..Default::default()
},
)))
.await;
let accepted = h.send(h.request("recall-projection", HUMAN)).await.unwrap();
let interaction =
meerkat_core::interaction::InteractionId(accepted.interaction_id.parse().unwrap());
let injected = tokio::time::timeout(WAIT, async {
loop {
let injected: Vec<_> = h
.history()
.await
.into_iter()
.filter_map(|message| match message {
Message::User(user)
if user.identity.interaction_id == Some(interaction)
&& user.transcript_role == TranscriptUserRole::InjectedContext =>
{
Some(ContentInput::Text(user.text_content()))
}
_ => None,
})
.collect();
if !injected.is_empty() {
break injected;
}
tokio::task::yield_now().await;
}
})
.await
.expect("configured recall must reach the canonical injected-context slot");
assert!(serde_json::to_string(&injected).unwrap().contains(MEMORY));
h.finish_human_with_context(&accepted, HUMAN, HandlingMode::Queue, injected)
.await;
let canonical = h.history().await;
let mut retained_canonical = None;
for _ in 0..2 {
h.aggregator.refresh_session_history().await.unwrap();
let full = h
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(h.identity.to_string()),
..Default::default()
})
.await
.unwrap();
let source =
assert_human_input_source_pair(&full.frames, &accepted, &h.session_id, &canonical);
if let Some(previous) = &retained_canonical {
assert_eq!(
&source, previous,
"refresh cannot duplicate or rewrite the canonical row"
);
} else {
retained_canonical = Some(source);
}
assert!(
!serde_json::to_string(&full.frames)
.unwrap()
.contains(MEMORY)
);
}
let retained_canonical = retained_canonical.unwrap();
let mut cursor = accepted.cursor.clone();
let mut reached_end = false;
let mut paged_inputs = Vec::new();
for _ in 0..20 {
let page = h
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(h.identity.to_string()),
after: Some(cursor.clone()),
limit: 1,
..Default::default()
})
.await
.unwrap();
for frame in page
.frames
.iter()
.filter(|frame| frame.kind == "user_input")
{
assert_canonical_human_frame(frame, &accepted, &h.session_id, &canonical);
assert_eq!(frame, &retained_canonical);
paged_inputs.push(frame.clone());
}
assert!(
!serde_json::to_string(&page.frames)
.unwrap()
.contains(MEMORY)
);
let Some(last) = page.frames.last() else {
reached_end = true;
break;
};
cursor = last.cursor.clone();
}
assert!(
reached_end,
"bounded pagination must exhaust the transcript"
);
assert_eq!(
paged_inputs,
vec![retained_canonical],
"paging retains the committed human row exactly once, without another admission or recall row"
);
let history_only = MobKitConsoleAggregator::in_memory();
history_only.register_runtime_handles_with_policy(
"human",
"",
h.unified.mob_runtime().clone(),
Some(h.identity_runtime.clone()),
ConsoleEventStore::new(),
Arc::new(AllowAllConsoleVisibilityPolicy),
);
history_only.refresh_session_history().await.unwrap();
let page = history_only
.query_timeline(ConsoleTimelineQuery {
identity: Some(h.identity.to_string()),
..Default::default()
})
.await
.unwrap();
assert_eq!(
page.frames
.iter()
.filter(|frame| frame.kind == "user_input")
.count(),
1
);
assert_canonical_human_frame(
page.frames
.iter()
.find(|frame| frame.kind == "user_input")
.unwrap(),
&accepted,
&h.session_id,
&canonical,
);
assert!(
!serde_json::to_string(&page.frames)
.unwrap()
.contains(MEMORY)
);
let canonical = h.history().await;
assert!(canonical.iter().any(|message| matches!(message,
Message::User(user) if user.transcript_role == TranscriptUserRole::InjectedContext
&& user.identity.interaction_id == Some(interaction)
&& user.text_content().contains(MEMORY))));
assert_eq!(canonical.iter().filter(|message| matches!(message,
Message::User(user) if user.transcript_role == TranscriptUserRole::Conversational
&& user.identity.interaction_id == Some(interaction) && user.text_content() == HUMAN
)).count(), 1);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_timeout_replay_preserves_failed_receipt_without_repreparing() {
let h = Harness::new(Duration::ZERO).await;
let request = h.request("unknown", HUMAN);
let accepted = h.send(request.clone()).await.unwrap();
h.wait_status(&accepted, ConsoleFrameStatus::DeliveryFailed)
.await;
let replay = h.send(request).await.unwrap_err();
let data = super::console_send_error_data(&replay).unwrap();
assert_eq!(data["kind"], "host_human_input_replay_unavailable");
assert_eq!(data["input_frame_id"], accepted.input_frame_id);
assert_eq!(data["execution_fate"], "unknown");
assert_eq!(data["retry_submitted"], false);
assert!(h.client.requests.lock().unwrap().is_empty());
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_lost_admission_observation_does_not_reexecute_committed_work() {
let h = Harness::new(WAIT).await;
let request = h.request("lost-response", HUMAN);
let accepted = h.send(request.clone()).await.unwrap();
h.finish_human(&accepted, HUMAN).await;
let calls = h.client.requests.lock().unwrap().len();
h.aggregator
.mark_interaction_delivery_failed(&accepted.input_frame_id)
.await
.unwrap();
let replay = h.send(request).await.unwrap_err();
assert_eq!(
super::console_send_error_data(&replay).unwrap()["kind"],
"host_human_input_replay_unavailable",
);
assert_eq!(h.client.requests.lock().unwrap().len(), calls);
assert_eq!(
h.history()
.await
.iter()
.filter(
|message| matches!(message, Message::User(user) if user.text_content() == HUMAN)
)
.count(),
1
);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_bridge_slots_strict_replay_and_stale_session_refusal() {
let h = Harness::new(WAIT).await;
let interaction = uuid::Uuid::new_v4();
let mut delivery = BridgeDelivery::new(HUMAN.into(), HandlingMode::Queue);
delivery.interaction_id = Some(interaction.to_string());
delivery.delivery_identity =
Some(MobDeliveryIdentity::new("slots", interaction.to_string()).unwrap());
delivery.system_prompt = Some("PER_TURN_SYSTEM".to_string());
delivery.injected_context = vec!["AMBIENT_MEMORY".into()];
h.bridge
.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery.clone())
.await
.unwrap();
let member = h
.unified
.mob_handle()
.get_member(&crate::member_comms_id::mob_member_id(h.identity.as_str()))
.await
.unwrap()
.unwrap();
let spec = WorkSpec::new(HUMAN, WorkOrigin::Internal)
.with_system_prompt("PER_TURN_SYSTEM")
.with_injected_context(vec!["AMBIENT_MEMORY".into()])
.with_interaction_id(meerkat_core::interaction::InteractionId(interaction));
let turn = h
.unified
.mob_handle()
.start_host_human_input_bounded(
member.agent_runtime_id,
member.fence_token,
spec,
HandlingMode::Queue,
delivery.delivery_identity.clone().unwrap(),
std::time::Instant::now() + WAIT,
)
.await
.unwrap();
tokio::time::timeout(WAIT, turn.wait())
.await
.unwrap()
.unwrap();
let history = h.history().await;
assert!(history.iter().any(|message| matches!(message,
Message::User(user) if user.transcript_role == TranscriptUserRole::InjectedContext
&& user.text_content() == "AMBIENT_MEMORY")));
assert!(history.iter().any(|message| matches!(message,
Message::System(system) if system.content == "PER_TURN_SYSTEM")));
let calls_after_turn = h.client.requests.lock().unwrap().len();
h.bridge
.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery.clone())
.await
.unwrap();
delivery.injected_context = vec!["CHANGED_MEMORY".into()];
assert!(matches!(
h.bridge.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery.clone()).await,
Err(BridgeError::HostHumanInput(HostHumanInputError::Mob(error)))
if matches!(*error, meerkat_mob::MobError::WorkInputIdempotencyConflict { .. })
));
assert!(matches!(
h.bridge
.deliver_host_human_input(&h.runtime_id, &SessionId::new(), delivery.clone())
.await,
Err(BridgeError::HostHumanInput(
HostHumanInputError::BindingChanged
))
));
delivery.handling_mode = HandlingMode::Steer;
let interaction = uuid::Uuid::new_v4();
delivery.interaction_id = Some(interaction.to_string());
delivery.delivery_identity =
Some(MobDeliveryIdentity::new("steer-slots", interaction.to_string()).unwrap());
let refusal = h
.bridge
.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery.clone())
.await;
assert!(
matches!(
&refusal,
Err(BridgeError::HostHumanInput(HostHumanInputError::Mob(error)))
if matches!(error.as_ref(), meerkat_mob::MobError::InjectedContextUndeliverable { .. })
),
"{refusal:?}"
);
delivery.injected_context.clear();
let refusal = h
.bridge
.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery)
.await;
assert!(
matches!(
&refusal,
Err(BridgeError::HostHumanInput(HostHumanInputError::Mob(error)))
if matches!(error.as_ref(), meerkat_mob::MobError::UnsupportedForMode { .. })
),
"{refusal:?}"
);
assert_eq!(h.client.requests.lock().unwrap().len(), calls_after_turn);
assert_eq!(
h.history()
.await
.iter()
.filter(
|message| matches!(message, Message::User(user) if user.text_content() == HUMAN)
)
.count(),
1
);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_member_only_worker_uses_the_same_canonical_user_seam() {
let h = Harness::new(WAIT).await;
h.unified
.mob_handle()
.spawn_spec(meerkat_mob::SpawnMemberSpec::from_wire(
"human".to_string(),
"worker".to_string(),
None,
None,
None,
))
.await
.unwrap();
let request = ConsoleSendRequest {
identity: "worker".to_string(),
..h.request("worker", HUMAN)
};
let accepted = h.aggregator.send(request.clone()).await.unwrap();
h.finish_human(&accepted, HUMAN).await;
let messages = h
.unified
.mob_runtime()
.read_session_history(accepted.session_id.as_deref().unwrap(), 0, None)
.await
.unwrap()
.messages;
assert_eq!(
messages
.iter()
.filter(
|message| matches!(message, Message::User(user) if user.text_content() == HUMAN)
)
.count(),
1
);
let replay = h.aggregator.send(request).await.unwrap();
assert_eq!(replay.input_frame_id, accepted.input_frame_id);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_active_steer_stays_steer_and_commits_once() {
let h = Harness::with_kickoff(WAIT, true).await;
let calls_before = h.client.requests.lock().unwrap().len();
h.client.gate.send_replace(false);
let active = h
.send(h.request("active", "active queue work"))
.await
.unwrap();
tokio::time::timeout(WAIT, async {
while h.client.requests.lock().unwrap().len() <= calls_before {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
let request = ConsoleSendRequest {
handling_mode: Some("steer".to_string()),
origin_kind: None,
..h.request("steer", HUMAN)
};
let steer = tokio::time::timeout(WAIT, h.send(request.clone()))
.await
.unwrap()
.unwrap();
assert!(
!*h.client.gate.borrow(),
"admission cannot wait for the active LLM"
);
h.client.gate.send_replace(true);
h.finish_human(&active, "active queue work").await;
h.finish_human_with_mode(&steer, HUMAN, HandlingMode::Steer)
.await;
assert_eq!(
h.history()
.await
.iter()
.filter(
|message| matches!(message, Message::User(user) if user.text_content() == HUMAN)
)
.count(),
1
);
assert_eq!(
h.send(request).await.unwrap().input_frame_id,
steer.input_frame_id
);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_human_reset_refuses_old_prepared_session_and_old_worker_fence() {
let h = Harness::new(WAIT).await;
let request = h.request("before-reset", HUMAN);
let accepted = h.send(request.clone()).await.unwrap();
h.finish_human(&accepted, HUMAN).await;
let old_member = h
.unified
.mob_handle()
.list_members()
.await
.into_iter()
.find(|member| member.agent_identity.as_str() == h.identity.as_str())
.unwrap();
let successor = h.identity_runtime.reset(&h.identity).await.unwrap();
assert_ne!(successor.session_id, h.session_id);
let replay = h.send(request).await.unwrap();
assert_eq!(
replay.session_id.as_deref(),
Some(h.session_id.to_string().as_str())
);
let mut delivery = BridgeDelivery::new(HUMAN.into(), HandlingMode::Queue);
delivery.interaction_id = Some(accepted.interaction_id.clone());
delivery.delivery_identity =
Some(MobDeliveryIdentity::new(&accepted.input_frame_id, &accepted.interaction_id).unwrap());
assert!(matches!(
h.bridge
.deliver_host_human_input(&h.runtime_id, &h.session_id, delivery)
.await,
Err(BridgeError::HostHumanInput(
HostHumanInputError::BindingChanged
))
));
assert!(
crate::mob_handle_runtime::send_console_human_on_mob(
&h.unified.mob_handle(),
&old_member,
HUMAN.into(),
&[],
HandlingMode::Queue,
&accepted,
)
.await
.is_err()
);
let replacement_history = h
.unified
.mob_runtime()
.read_session_history(&successor.session_id.to_string(), 0, None)
.await
.unwrap();
assert!(
!replacement_history
.messages
.iter()
.any(|message| matches!(message, Message::User(user) if user.text_content() == HUMAN))
);
h.stop().await;
}
const SKILL_BODY: &str = "Answer in the chosen house style: numbered, terse, no hedging.";
const OTHER_SKILL_BODY: &str = "Answer in the other house style: long prose with caveats.";
fn house_style_skill_config(root: &std::path::Path) -> (meerkat::Config, SkillKey, SkillKey) {
let name = SkillName::parse("house-style").unwrap();
let mut config = meerkat::Config::default();
let mut keys = Vec::new();
for (index, (source, body)) in [
("5b0c7c2e-4c61-4f2d-9d0e-2f5a8c6b1e11", SKILL_BODY),
("9e4d2a10-7b3c-4c8e-a1f5-6d2b0c9e8f77", OTHER_SKILL_BODY),
]
.into_iter()
.enumerate()
{
let source_uuid = SourceUuid::parse(source).unwrap();
let repository = root.join(format!("skills-{index}"));
let skill_dir = repository.join(name.as_str());
std::fs::create_dir_all(&skill_dir).unwrap();
std::fs::write(
skill_dir.join("SKILL.md"),
format!("---\nname: {name}\ndescription: fixture skill\n---\n{body}\n"),
)
.unwrap();
config
.skills
.repositories
.push(meerkat_core::skills_config::SkillRepositoryConfig {
name: format!("fixture-skills-{index}"),
source_uuid: source_uuid.clone(),
transport: meerkat_core::skills_config::SkillRepoTransport::Filesystem {
path: repository.to_string_lossy().into_owned(),
},
});
keys.push(SkillKey::new(source_uuid, name.clone()));
}
let other = keys.pop().unwrap();
let chosen = keys.pop().unwrap();
(config, chosen, other)
}
fn skill_context_rows(messages: &[Message], interaction_id: &str) -> Vec<(SkillKey, String)> {
let interaction =
meerkat_core::interaction::InteractionId(interaction_id.parse::<uuid::Uuid>().unwrap());
messages
.iter()
.filter_map(|message| match message {
Message::User(user) if user.identity.interaction_id == Some(interaction) => Some(user),
_ => None,
})
.flat_map(|user| user.content.iter())
.filter_map(|block| match block {
meerkat_core::types::ContentBlock::SkillContext { skill_key, text } => {
Some((skill_key.clone(), text.clone()))
}
_ => None,
})
.collect()
}
#[tokio::test(flavor = "multi_thread")]
async fn console_send_selected_skills_resolve_natively_and_replay_exactly() {
use futures::StreamExt;
let skills = tempfile::tempdir_in(std::env::current_dir().unwrap()).unwrap();
let (config, chosen, other) = house_style_skill_config(skills.path());
let h = Harness::with_meerkat_config(config).await;
let mut events = h
.unified
.mob_handle()
.subscribe_agent_events(&crate::member_comms_id::mob_member_id(h.identity.as_str()))
.await
.unwrap();
let request = ConsoleSendRequest {
skill_refs: vec![SkillRef::Structured(chosen.clone())],
..h.request("selected-skill", HUMAN)
};
let accepted = h.send(request.clone()).await.unwrap();
h.finish_human_with_skills(&accepted, HUMAN, vec![chosen.clone()])
.await;
let mut resolved = Vec::new();
tokio::time::timeout(WAIT, async {
while let Some(envelope) = events.next().await {
match envelope.payload {
meerkat_core::AgentEvent::SkillsResolved { skills, .. } => resolved.push(skills),
meerkat_core::AgentEvent::SkillResolutionFailed { reason, .. } => {
panic!("selected skill must resolve, got {reason}")
}
meerkat_core::AgentEvent::RunCompleted { .. } if !resolved.is_empty() => break,
_ => {}
}
}
})
.await
.expect("the selected turn's events arrive");
assert_eq!(
resolved,
vec![vec![chosen.clone()]],
"one native activation"
);
let rows = skill_context_rows(&h.history().await, &accepted.interaction_id);
assert_eq!(rows.len(), 1, "exactly one durable SkillContext");
assert_eq!(rows[0].0, chosen);
assert!(rows[0].1.contains(SKILL_BODY), "native rendered skill body");
let requests = h.client.requests.lock().unwrap().clone();
let bodies = requests
.iter()
.map(|request| serde_json::to_string(&request.messages).unwrap())
.collect::<Vec<_>>();
assert_eq!(
bodies
.iter()
.filter(|body| body.contains(SKILL_BODY))
.count(),
1,
"exactly one provider call carries the selected body"
);
assert!(
bodies.iter().all(|body| !body.contains(OTHER_SKILL_BODY)),
"the same-name skill is never selected"
);
let calls_after_turn = requests.len();
let frames_before_replay = h.user_input_frames(&accepted.identity).await;
let replay = h.send(request.clone()).await.unwrap();
assert_eq!(replay.input_frame_id, accepted.input_frame_id);
assert_eq!(replay.interaction_id, accepted.interaction_id);
let changed = ConsoleSendRequest {
skill_refs: vec![SkillRef::Structured(other)],
..request
};
assert!(matches!(
h.send(changed).await,
Err(crate::console_aggregator::ConsoleSendError::IdempotencyConflict(_))
));
assert_eq!(
h.user_input_frames(&accepted.identity).await,
frames_before_replay,
"no new console input was reserved or dispatched"
);
h.finish_human_with_skills(&accepted, HUMAN, vec![chosen.clone()])
.await;
assert_eq!(
skill_context_rows(&h.history().await, &accepted.interaction_id).len(),
1
);
assert_eq!(h.client.requests.lock().unwrap().len(), calls_after_turn);
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_send_selected_skills_reach_a_member_only_worker() {
let skills = tempfile::tempdir_in(std::env::current_dir().unwrap()).unwrap();
let (config, key, _) = house_style_skill_config(skills.path());
let h = Harness::with_meerkat_config(config).await;
h.unified
.mob_handle()
.spawn_spec(meerkat_mob::SpawnMemberSpec::from_wire(
"human".to_string(),
"worker".to_string(),
None,
None,
None,
))
.await
.unwrap();
let request = ConsoleSendRequest {
identity: "worker".to_string(),
skill_refs: vec![SkillRef::Structured(key.clone())],
..h.request("worker-skill", HUMAN)
};
let accepted = h.aggregator.send(request.clone()).await.unwrap();
h.finish_human_with_skills(&accepted, HUMAN, vec![key.clone()])
.await;
let messages = h
.unified
.mob_runtime()
.read_session_history(accepted.session_id.as_deref().unwrap(), 0, None)
.await
.unwrap()
.messages;
let rows = skill_context_rows(&messages, &accepted.interaction_id);
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].0, key);
assert!(rows[0].1.contains(SKILL_BODY));
let frames_before_replay = h.user_input_frames(&accepted.identity).await;
let replay = h.aggregator.send(request).await.unwrap();
assert_eq!(replay.input_frame_id, accepted.input_frame_id);
assert_eq!(
h.user_input_frames(&accepted.identity).await,
frames_before_replay,
"the replay reserved and dispatched nothing new"
);
h.stop().await;
}
#[test]
fn console_send_request_skill_refs_are_typed_and_optional() {
let bare: ConsoleSendRequest = serde_json::from_value(json!({
"identity": "human",
"content": "hello",
"origin": "console",
"idempotency_key": "k",
}))
.unwrap();
assert!(bare.skill_refs.is_empty());
assert!(
serde_json::to_value(&bare)
.unwrap()
.get("skill_refs")
.is_none()
);
let typed: ConsoleSendRequest = serde_json::from_value(json!({
"identity": "human",
"content": "hello",
"origin": "console",
"idempotency_key": "k",
"skill_refs": [{
"kind": "structured",
"source_uuid": "5b0c7c2e-4c61-4f2d-9d0e-2f5a8c6b1e11",
"skill_name": "house-style",
}],
}))
.unwrap();
assert_eq!(
typed.selected_skill_keys(),
vec![SkillKey::new(
SourceUuid::parse("5b0c7c2e-4c61-4f2d-9d0e-2f5a8c6b1e11").unwrap(),
SkillName::parse("house-style").unwrap(),
)]
);
assert!(
serde_json::from_value::<ConsoleSendRequest>(json!({
"identity": "human",
"content": "hello",
"origin": "console",
"idempotency_key": "k",
"skill_refs": ["house-style"],
}))
.is_err(),
"untyped legacy skill strings are refused, not folded"
);
}
fn dormant_spec(identity: &str) -> DurableAgentSpec {
DurableAgentSpec {
identity: AgentIdentity::parse(identity).unwrap(),
profile: "human".into(),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: Some(meerkat_mob::MobRuntimeMode::AutonomousHost),
backend: None,
binding: None,
placement: None,
}
}
fn selected_skill() -> SkillRef {
SkillRef::Structured(SkillKey::new(
SourceUuid::parse("5b0c7c2e-4c61-4f2d-9d0e-2f5a8c6b1e11").unwrap(),
SkillName::parse("house-style").unwrap(),
))
}
#[tokio::test(flavor = "multi_thread")]
async fn console_send_selected_skills_refused_before_reservation_for_external_or_placed_identity() {
let h = Harness::new(WAIT).await;
let mut external = dormant_spec("external-target");
external.backend = Some(meerkat_mob::MobBackendKind::External);
let mut placed = dormant_spec("placed-target");
placed.placement = Some(meerkat_contracts::WireHostRef("remote-host".to_string()));
for spec in [external, placed] {
let identity = spec.identity.to_string();
h.identity_runtime
.register(spec, IdentityLifecycleState::Dormant, None, None)
.await;
let request = ConsoleSendRequest {
identity: identity.clone(),
skill_refs: vec![selected_skill()],
..h.request("refused-selection", HUMAN)
};
let refusal = h.send(request).await;
assert!(
matches!(
&refusal,
Err(crate::console_aggregator::ConsoleSendError::InvalidRequest(
_
))
),
"{identity}: {refusal:?}"
);
let frames = h
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(identity.clone()),
..Default::default()
})
.await
.unwrap()
.frames;
assert!(
frames.is_empty(),
"{identity}: nothing was reserved: {frames:?}"
);
}
h.stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn console_send_selected_skills_member_lane_admission_refusal_surfaces_typed() {
let h = Harness::new(WAIT).await;
let calls_before = h.client.requests.lock().unwrap().len();
h.client.gate.send_replace(false);
h.unified
.mob_handle()
.spawn_spec(meerkat_mob::SpawnMemberSpec::from_wire(
"human".to_string(),
"worker".to_string(),
Some("worker kickoff".into()),
None,
None,
))
.await
.unwrap();
tokio::time::timeout(WAIT, async {
while h.client.requests.lock().unwrap().len() <= calls_before {
tokio::task::yield_now().await;
}
})
.await
.expect("the worker kickoff reaches the provider");
let request = ConsoleSendRequest {
identity: "worker".to_string(),
handling_mode: Some("steer".to_string()),
skill_refs: vec![selected_skill()],
..h.request("member-lane-refusal", HUMAN)
};
let accepted = h.aggregator.send(request).await.unwrap();
h.wait_status(&accepted, ConsoleFrameStatus::DeliveryFailed)
.await;
let frames = h
.aggregator
.query_timeline(ConsoleTimelineQuery {
identity: Some(accepted.identity.clone()),
..Default::default()
})
.await
.unwrap()
.frames;
let failure = frames
.iter()
.find(|frame| {
frame.kind == "message_delivery_failed"
&& frame.parent_frame_id.as_deref() == Some(accepted.input_frame_id.as_str())
})
.expect("a typed delivery failure frame");
assert_eq!(
failure.payload["data"]["kind"], "host_human_input_unsupported",
"{:?}",
failure.payload
);
h.client.gate.send_replace(true);
h.stop().await;
}