use super::*;
use crate::identity_first::AgentIdentity;
use meerkat_core::SessionError;
use meerkat_core::types::{Message, SystemNoticeKind, SystemNoticeMessage};
use std::num::NonZeroU64;
pub const OPERATOR_SESSION_BUSY_CODE: i64 = -32015;
pub const OPERATOR_VERB_UNAVAILABLE_CODE: i64 = -32016;
const DEFAULT_COMPACT_FLOOR_TOKENS: u64 = 1024;
const DEFAULT_COMPACT_TIMEOUT_MS: u64 = 60_000;
const COMPACT_TIMEOUT_EVIDENCE_BUDGET: Duration = Duration::from_secs(5);
const DEFAULT_BOUND_KEEP_LAST: usize = 50;
fn rpc_error(response_id: Value, code: i64, message: String) -> JsonRpcResponse {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: None,
error: Some(JsonRpcError {
code,
message,
data: None,
}),
}
}
fn rpc_result(response_id: Value, result: Value) -> JsonRpcResponse {
JsonRpcResponse {
jsonrpc: JSONRPC_VERSION.to_string(),
id: response_id,
result: Some(result),
error: None,
}
}
fn optional_u64_param(params: &Value, field: &str) -> Result<Option<u64>, String> {
match params.get(field) {
None | Some(Value::Null) => Ok(None),
Some(value) => match value.as_u64() {
Some(parsed) if parsed > 0 => Ok(Some(parsed)),
_ => Err(format!("{field} must be a positive integer")),
},
}
}
pub(crate) fn pair_safe_cut_index(messages: &[Message], keep_last: usize) -> usize {
let mut cut = messages.len().saturating_sub(keep_last);
while cut > 0 && matches!(messages.get(cut), Some(Message::ToolResults { .. })) {
cut -= 1;
}
cut
}
struct ResolvedOperatorTarget {
identity: AgentIdentity,
expected_alias: Option<String>,
}
async fn resolve_operator_target(
runtime: &UnifiedRuntime,
identity_rt: &crate::identity_first::IdentityRuntime,
params: &Value,
response_id: &Value,
) -> Result<ResolvedOperatorTarget, Box<JsonRpcResponse>> {
let identity_str = params
.get("identity")
.and_then(|v| v.as_str())
.unwrap_or("");
let target = match resolve_rpc_identity_control_target(runtime, identity_rt, identity_str).await
{
Ok(target) => target,
Err(e) => {
return Err(Box::new(rpc_error(
response_id.clone(),
-32602,
format!("invalid identity: {e}"),
)));
}
};
if let Some(response) =
rpc_stale_live_alias_error_response(identity_rt, &target, response_id.clone()).await
{
return Err(Box::new(response));
}
let expected_alias = crate::member_comms_id::is_reserved_generated_alias(identity_str)
.then(|| identity_str.to_string());
Ok(ResolvedOperatorTarget {
identity: target.identity,
expected_alias,
})
}
async fn read_transcript_facts(
service: &Arc<dyn crate::memory::hygienist::TranscriptEditSessionService>,
session_id: &meerkat_core::types::SessionId,
) -> Result<(usize, Option<String>, Option<String>), SessionError> {
let page = service
.read_history(
session_id,
meerkat_core::service::SessionHistoryQuery {
offset: 0,
limit: None,
},
)
.await?;
let (head_revision, last_reason) = match service
.list_transcript_revisions(
session_id,
meerkat_core::service::SessionTranscriptRevisionListQuery {
limit: None,
offset: None,
},
)
.await
{
Ok(list) => (
Some(list.head_revision),
list.entries.last().map(|entry| entry.reason.clone()),
),
Err(SessionError::Unsupported(_)) => (None, None),
Err(err) => return Err(err),
};
Ok((page.messages.len(), head_revision, last_reason))
}
pub(super) async fn handle_compact_member(
runtime: &UnifiedRuntime,
ctx: &IdentityFirstContext,
params: &Value,
response_id: Value,
) -> JsonRpcResponse {
let identity_rt = &ctx.runtime;
let Some(floors) = ctx.compaction_floors.as_ref() else {
return rpc_error(
response_id,
OPERATOR_VERB_UNAVAILABLE_CODE,
"compact_member is not wired on this gateway: the identity bridge's compaction-floor \
registry was not threaded into the RPC context"
.to_string(),
);
};
let floor_tokens = match optional_u64_param(params, "floor_tokens") {
Ok(value) => value.unwrap_or(DEFAULT_COMPACT_FLOOR_TOKENS),
Err(message) => return rpc_error(response_id, -32602, message),
};
let Some(floor) = NonZeroU64::new(floor_tokens) else {
return rpc_error(
response_id,
-32602,
"floor_tokens must be greater than 0".to_string(),
);
};
let timeout_ms = match optional_u64_param(params, "timeout_ms") {
Ok(value) => value.unwrap_or(DEFAULT_COMPACT_TIMEOUT_MS),
Err(message) => return rpc_error(response_id, -32602, message),
};
let target = match resolve_operator_target(runtime, identity_rt, params, &response_id).await {
Ok(target) => target,
Err(response) => return *response,
};
let identity = target.identity;
let Some(spec) = identity_rt
.roster_inspect()
.await
.remove(&identity)
.map(|(spec, _)| spec)
else {
return identity_error_response(
response_id,
&crate::identity_first::IdentityRuntimeError::UnknownIdentity(identity),
);
};
floors.set(&identity, floor);
let session_id = match rebuild_member_for_fresh_build(
ctx,
&identity,
target.expected_alias.as_deref(),
spec.clone(),
)
.await
{
Ok(session_id) => session_id,
Err(detail) => {
floors.clear(&identity);
return rpc_error(
response_id,
-32000,
format!("compact_member could not rebuild the member with the floor: {detail}"),
);
}
};
let before = match ctx.transcript_edit_service.as_ref() {
Some(service) => match read_transcript_facts(service, &session_id).await {
Ok(facts) => Some(facts),
Err(err) => {
floors.clear(&identity);
let (_, restore) = restore_after_floor(ctx, &identity, spec.clone()).await;
return rpc_error(
response_id,
-32000,
format!(
"compact_member aborted reading the transcript before compaction: {err}\
{restore}"
),
);
}
},
None => None,
};
let nudge = meerkat_core::ContentInput::Text(
"[mobkit-gateway operator verb compact_member] Maintenance turn: transcript compaction \
was forced for this turn. Reply with a brief acknowledgement only."
.to_string(),
);
let admission = match identity_rt
.send_admission_tracked(
&identity,
None,
&nudge,
meerkat_core::types::HandlingMode::Queue,
None,
)
.await
{
Ok(admission) => admission,
Err(err) => {
floors.clear(&identity);
let (_, restore) = restore_after_floor(ctx, &identity, spec.clone()).await;
return rpc_error(
response_id,
-32000,
format!("compact_member maintenance turn was not admitted: {err}{restore}"),
);
}
};
if let Err(err) = identity_rt
.wait_for_completion(
&identity,
admission.completion_baseline,
Duration::from_millis(timeout_ms),
)
.await
{
floors.clear(&identity);
let (rolled_back, restore) = restore_after_floor(ctx, &identity, spec.clone()).await;
let turn_fate = if rolled_back {
"; the in-flight maintenance turn was quiesced by the rollback rebuild \
(mob retirement cancels the active runtime turn before retiring)"
} else {
"; the maintenance turn may still be running on the floored build, so this member \
can stay briefly unresponsive and later reads on it may queue"
};
let compaction_evidence = match (before.as_ref(), ctx.transcript_edit_service.as_ref()) {
(Some(before), Some(service)) => {
match tokio::time::timeout(
COMPACT_TIMEOUT_EVIDENCE_BUDGET,
read_transcript_facts(service, &session_id),
)
.await
{
Ok(Ok(after)) if &after != before => {
"; the forced compaction rewrite IS durably applied \
(transcript facts changed before the timeout)"
}
Ok(Ok(_)) => {
"; the forced compaction rewrite had not durably applied at \
timeout"
}
Ok(Err(_)) | Err(_) => "",
}
}
_ => "",
};
return rpc_error(
response_id,
-32000,
format!(
"compact_member maintenance turn did not complete within {timeout_ms}ms: \
{err}{turn_fate}{compaction_evidence}{restore}"
),
);
}
floors.clear(&identity);
if let Err(detail) = rebuild_member_for_fresh_build(ctx, &identity, None, spec).await {
return rpc_error(
response_id,
-32000,
format!(
"compact_member forced the compaction but the restore rebuild failed: {detail}; \
the member keeps the temporary floor ({floor} tokens) until its next rebuild"
),
);
}
let after = match ctx.transcript_edit_service.as_ref() {
Some(service) => match read_transcript_facts(service, &session_id).await {
Ok(facts) => Some(facts),
Err(err) => {
return rpc_error(
response_id,
-32000,
format!(
"compact_member completed but reading the post-compaction transcript \
failed: {err}"
),
);
}
},
None => None,
};
let messages_before = before.as_ref().map(|(count, _, _)| *count);
let messages_after = after.as_ref().map(|(count, _, _)| *count);
let compaction_applied = match (messages_before, messages_after) {
(Some(before), Some(after)) => Some(after < before),
_ => None,
};
rpc_result(
response_id,
serde_json::json!({
"identity": identity.as_str(),
"session_id": session_id.to_string(),
"floor_tokens": floor.get(),
"messages_before": messages_before,
"messages_after": messages_after,
"compaction_applied": compaction_applied,
"head_revision": after.as_ref().and_then(|(_, head, _)| head.clone()),
"last_rewrite_reason": after.as_ref().and_then(|(_, _, reason)| reason.clone()),
}),
)
}
async fn rebuild_member_for_fresh_build(
ctx: &IdentityFirstContext,
identity: &AgentIdentity,
expected_alias: Option<&str>,
spec: crate::identity_first::DurableAgentSpec,
) -> Result<meerkat_core::types::SessionId, String> {
use crate::identity_first::IdentityLifecycleState;
let identity_rt = &ctx.runtime;
let state = identity_rt
.status(identity)
.await
.map_err(|err| format!("status before rebuild: {err}"))?
.state;
if state == IdentityLifecycleState::Active {
match expected_alias {
Some(alias) => identity_rt
.retire_member_alias_tracked(identity, alias)
.await
.map(|_| ()),
None => identity_rt.retire_tracked(identity).await.map(|_| ()),
}
.map_err(|err| format!("quiesce retire: {err}"))?;
}
let resolved = identity_rt
.continuity_store()
.resolve_many(std::slice::from_ref(identity))
.await
.map_err(|err| format!("continuity resolve after quiesce: {err}"))?;
let record = match resolved.get(identity) {
Some(crate::identity_first::ContinuityResolveState::Ready { record }) => record.clone(),
Some(crate::identity_first::ContinuityResolveState::Broken { failure }) => {
return Err(format!(
"identity {} has broken continuity ({}); cannot rebuild onto the same durable \
session",
identity.as_str(),
failure.detail
));
}
Some(crate::identity_first::ContinuityResolveState::Uninitialized) | None => {
return Err(format!(
"identity {} has no durable continuity record; cannot rebuild onto the same \
durable session",
identity.as_str()
));
}
};
identity_rt
.register(spec, IdentityLifecycleState::Dormant, Some(record), None)
.await;
identity_rt
.materialize_tracked(identity)
.await
.map(|record| record.session_id)
.map_err(|err| format!("re-materialization: {err}"))
}
async fn restore_after_floor(
ctx: &IdentityFirstContext,
identity: &AgentIdentity,
spec: crate::identity_first::DurableAgentSpec,
) -> (bool, String) {
match rebuild_member_for_fresh_build(ctx, identity, None, spec).await {
Ok(_) => (
true,
"; the temporary floor was rolled back (member rebuilt at its original threshold)"
.to_string(),
),
Err(detail) => (
false,
format!(
"; rollback rebuild also failed ({detail}) - the member keeps the temporary \
floor until its next rebuild"
),
),
}
}
pub(super) async fn handle_bound_member_transcript(
runtime: &UnifiedRuntime,
ctx: &IdentityFirstContext,
params: &Value,
response_id: Value,
) -> JsonRpcResponse {
let identity_rt = &ctx.runtime;
let Some(service) = ctx.transcript_edit_service.as_ref() else {
return rpc_error(
response_id,
OPERATOR_VERB_UNAVAILABLE_CODE,
"bound_member_transcript is not wired on this gateway: the concrete persistent \
session service was not threaded into the RPC context (the erased MobSessionService \
cannot reach SessionServiceTranscriptEditExt)"
.to_string(),
);
};
let keep_last = match optional_u64_param(params, "keep_last") {
Ok(value) => value.map_or(DEFAULT_BOUND_KEEP_LAST, |parsed| parsed as usize),
Err(message) => return rpc_error(response_id, -32602, message),
};
let note = match params.get("note") {
None | Some(Value::Null) => None,
Some(Value::String(note)) => Some(note.clone()),
Some(_) => {
return rpc_error(
response_id,
-32602,
"note must be a string when provided".to_string(),
);
}
};
let target = match resolve_operator_target(runtime, identity_rt, params, &response_id).await {
Ok(target) => target,
Err(response) => return *response,
};
let identity = target.identity;
let status = match identity_rt.status(&identity).await {
Ok(status) => status,
Err(err) => return identity_error_response(response_id, &err),
};
let Some(session_id) = status.session_id else {
return rpc_error(
response_id,
-32000,
format!(
"identity {} has no current session to bound",
identity.as_str()
),
);
};
let messages = match service
.read_history(
&session_id,
meerkat_core::service::SessionHistoryQuery {
offset: 0,
limit: None,
},
)
.await
{
Ok(page) => page.messages,
Err(err) => {
return rpc_error(
response_id,
-32000,
format!("bound_member_transcript failed to read the transcript: {err}"),
);
}
};
let expected_parent_revision = match service
.list_transcript_revisions(
&session_id,
meerkat_core::service::SessionTranscriptRevisionListQuery {
limit: Some(0),
offset: None,
},
)
.await
{
Ok(list) => Some(list.head_revision),
Err(SessionError::Unsupported(_)) => None,
Err(err) => {
return rpc_error(
response_id,
-32000,
format!("bound_member_transcript failed to read the head revision: {err}"),
);
}
};
let cut = pair_safe_cut_index(&messages, keep_last);
if cut == 0 {
return rpc_result(
response_id,
serde_json::json!({
"identity": identity.as_str(),
"session_id": session_id.to_string(),
"bounded": false,
"removed": 0,
"message_count": messages.len(),
}),
);
}
let marker = Message::SystemNotice(SystemNoticeMessage::new(
SystemNoticeKind::Generic,
format!(
"[operator] transcript bounded: {cut} earlier message(s) were removed by \
bound_member_transcript; the conversation continues from the {} most recent \
message(s).",
messages.len() - cut
),
));
let mut reason = meerkat_core::TranscriptRewriteReason::new("operator_bound_transcript");
reason.note = Some(note.unwrap_or_else(|| {
format!("mobkit-gateway operator verb bound_member_transcript keep_last={keep_last}")
}));
let request = meerkat_core::service::SessionTranscriptRewriteRequest {
selection: meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: cut },
replacement: vec![marker],
reason,
actor: Some("mobkit-gateway operator verb".to_string()),
expected_parent_revision,
running_behavior: meerkat_core::TranscriptEditRunningBehavior::default(),
};
match service
.rewrite_session_transcript(&session_id, request)
.await
{
Ok(result) => rpc_result(
response_id,
serde_json::json!({
"identity": identity.as_str(),
"session_id": result.session_id.to_string(),
"bounded": true,
"removed": cut,
"kept": messages.len() - cut,
"message_count": result.message_count,
"parent_revision": result.parent_revision,
"revision": result.revision,
}),
),
Err(SessionError::Busy { id }) => rpc_error(
response_id,
OPERATOR_SESSION_BUSY_CODE,
format!(
"session {id} is live/running; transcript surgery only supports Reject while \
work is active - quiesce the member first (retire or park it), then retry"
),
),
Err(err) => rpc_error(
response_id,
-32000,
format!("bound_member_transcript rewrite failed: {err}"),
),
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use meerkat_core::types::{AssistantBlock, BlockAssistantMessage, ToolResult, UserMessage};
fn user(text: &str) -> Message {
Message::User(UserMessage::text(text))
}
fn tool_use_pair(id: &str) -> (Message, Message) {
let assistant = Message::BlockAssistant(BlockAssistantMessage::new(
vec![AssistantBlock::ToolUse {
id: id.to_string(),
name: "probe".to_string(),
args: serde_json::value::RawValue::from_string("{}".to_string()).expect("raw args"),
meta: None,
}],
meerkat_core::types::StopReason::ToolUse,
));
let results = Message::tool_results(vec![ToolResult {
tool_use_id: id.to_string(),
content: vec![],
is_error: false,
}]);
(assistant, results)
}
#[test]
fn cut_keeps_whole_transcript_when_keep_last_covers_it() {
let messages = vec![user("a"), user("b")];
assert_eq!(pair_safe_cut_index(&messages, 2), 0);
assert_eq!(pair_safe_cut_index(&messages, 10), 0);
}
#[test]
fn cut_lands_on_plain_message_boundary() {
let messages = vec![user("a"), user("b"), user("c"), user("d")];
assert_eq!(pair_safe_cut_index(&messages, 2), 2);
}
#[test]
fn cut_never_orphans_tool_results_from_their_assistant_message() {
let (assistant, results) = tool_use_pair("call-1");
let messages = vec![user("a"), assistant, results, user("tail")];
assert_eq!(pair_safe_cut_index(&messages, 2), 1);
}
#[test]
fn cut_after_tool_pair_is_untouched() {
let (assistant, results) = tool_use_pair("call-1");
let messages = vec![user("a"), assistant, results, user("tail")];
assert_eq!(pair_safe_cut_index(&messages, 1), 3);
}
#[test]
fn cut_walks_back_to_zero_when_transcript_leads_with_pairs() {
let (assistant, results) = tool_use_pair("call-1");
let messages = vec![assistant, results, user("tail")];
assert_eq!(pair_safe_cut_index(&messages, 2), 0);
}
use crate::identity_first::{
AgentAddressability, ContinuityGeneration, ContinuityRecord, ContinuityStore,
DurabilityPolicy, DurableAgentSpec, IdentityLifecycleState, IdentityRuntime,
IdentityRuntimeConfig, LeaseAcquireResult, LeaseProvider, LocalContinuityStore,
LocalLeaseProvider, MobSessionBridge, RosterContext, RosterError, RosterProvider,
};
use async_trait::async_trait;
use std::sync::atomic::{AtomicBool, Ordering};
fn worker_spec(identity: &AgentIdentity) -> DurableAgentSpec {
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: std::collections::BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
placement: None,
}
}
struct EmptyRoster;
#[async_trait]
impl RosterProvider for EmptyRoster {
async fn roster(
&self,
_context: &RosterContext,
) -> Result<Vec<DurableAgentSpec>, RosterError> {
Ok(Vec::new())
}
}
struct GatedUsageLlmClient {
input_tokens: u64,
gate_armed: Arc<AtomicBool>,
in_call: Arc<AtomicBool>,
release: Arc<tokio::sync::Notify>,
}
impl meerkat_client::LlmClient for GatedUsageLlmClient {
fn project_replay_messages(
&self,
messages: &[meerkat_core::Message],
) -> Result<Vec<meerkat_core::Message>, meerkat_client::LlmError> {
Ok(messages.to_vec())
}
fn stream<'a>(
&'a self,
request: &'a meerkat_client::LlmRequest,
) -> std::pin::Pin<
Box<
dyn futures::Stream<
Item = Result<meerkat_client::LlmEvent, meerkat_client::LlmError>,
> + Send
+ 'a,
>,
> {
use futures::StreamExt;
let gate_armed = self.gate_armed.load(Ordering::SeqCst);
let in_call = Arc::clone(&self.in_call);
let release = Arc::clone(&self.release);
let input_tokens = self.input_tokens;
Box::pin(
futures::stream::once(async move {
in_call.store(true, Ordering::SeqCst);
if gate_armed {
release.notified().await;
}
in_call.store(false, Ordering::SeqCst);
let [usage, done] =
crate::mob_handle_runtime::test_llm_usage::usage_then_done_with(
request,
meerkat_core::Provider::OpenAI,
meerkat_core::types::Usage {
input_tokens,
..Default::default()
},
meerkat_core::types::StopReason::EndTurn,
);
futures::stream::iter(vec![
Ok(meerkat_client::LlmEvent::TextDelta {
delta: "ack".to_string(),
meta: None,
}),
Ok(usage),
Ok(done),
])
})
.flatten(),
)
}
fn provider(&self) -> meerkat_core::Provider {
meerkat_core::Provider::OpenAI
}
fn health_check<'life0, 'async_trait>(
&'life0 self,
) -> std::pin::Pin<
Box<
dyn std::future::Future<Output = Result<(), meerkat_client::LlmError>>
+ Send
+ 'async_trait,
>,
>
where
'life0: 'async_trait,
Self: 'async_trait,
{
Box::pin(async { Ok(()) })
}
}
struct OperatorVerbHarness {
_temp_dir: tempfile::TempDir,
runtime: crate::UnifiedRuntime,
concrete: Arc<meerkat_session::PersistentSessionService<meerkat::FactoryAgentBuilder>>,
identity_runtime: Arc<IdentityRuntime>,
floors: Arc<crate::identity_first::CompactionFloorRegistry>,
identity: AgentIdentity,
member_alias: String,
gate_armed: Arc<AtomicBool>,
in_call: Arc<AtomicBool>,
release: Arc<tokio::sync::Notify>,
}
impl OperatorVerbHarness {
fn identity_ctx(&self) -> IdentityFirstContext {
IdentityFirstContext {
runtime: Arc::clone(&self.identity_runtime),
roster_provider: Arc::new(EmptyRoster),
topology_provider: None,
customizer: None,
agent_memory_provider: None,
mob_definition: Some(self.runtime.mob_handle().definition().clone()),
transcript_edit_service: Some(Arc::clone(&self.concrete) as _),
compaction_floors: Some(Arc::clone(&self.floors)),
}
}
async fn run_turn(&self, text: String) {
let admission = self
.identity_runtime
.send_admission_tracked(
&self.identity,
None,
&meerkat_core::ContentInput::Text(text),
meerkat_core::types::HandlingMode::Queue,
None,
)
.await
.expect("seed turn admitted");
self.identity_runtime
.wait_for_completion(
&self.identity,
admission.completion_baseline,
Duration::from_secs(30),
)
.await
.expect("seed turn completed");
}
async fn transcript_facts(&self) -> (usize, Option<String>) {
let service: Arc<dyn crate::memory::hygienist::TranscriptEditSessionService> =
Arc::clone(&self.concrete) as _;
let session_id = self
.identity_runtime
.status(&self.identity)
.await
.expect("identity status")
.session_id
.expect("identity session");
let (count, _head, last_reason) = read_transcript_facts(&service, &session_id)
.await
.expect("transcript facts");
(count, last_reason)
}
async fn settled_transcript_count(&self) -> usize {
let service: Arc<dyn crate::memory::hygienist::TranscriptEditSessionService> =
Arc::clone(&self.concrete) as _;
let session_id = self
.identity_runtime
.status(&self.identity)
.await
.expect("identity status")
.session_id
.expect("identity session");
settled_transcript_len(&service, &session_id).await
}
}
async fn transcript_len(
service: &Arc<dyn crate::memory::hygienist::TranscriptEditSessionService>,
session_id: &meerkat_core::types::SessionId,
) -> usize {
service
.read_history(
session_id,
meerkat_core::service::SessionHistoryQuery {
offset: 0,
limit: None,
},
)
.await
.expect("read durable transcript")
.messages
.len()
}
async fn settled_transcript_len(
service: &Arc<dyn crate::memory::hygienist::TranscriptEditSessionService>,
session_id: &meerkat_core::types::SessionId,
) -> usize {
const REQUIRED_STABLE_OBSERVATIONS: usize = 3;
let deadline = std::time::Instant::now() + Duration::from_mins(1);
let mut settled = transcript_len(service, session_id).await;
let mut stable = 1;
while stable < REQUIRED_STABLE_OBSERVATIONS {
tokio::time::sleep(Duration::from_millis(100)).await;
let observed = transcript_len(service, session_id).await;
if observed == settled {
stable += 1;
} else {
settled = observed;
stable = 1;
}
assert!(
std::time::Instant::now() < deadline,
"durable transcript never settled: still growing (last observed {observed} \
messages), so no index derived from it can be trusted"
);
}
settled
}
async fn operator_verb_harness(member: &str, mob_id: &str) -> OperatorVerbHarness {
let temp_dir = tempfile::tempdir().expect("temp dir");
let state = temp_dir.path().join("state");
std::fs::create_dir_all(&state).expect("state dir");
let session_store: Arc<dyn meerkat::SessionStore> = Arc::new(
meerkat_store::SqliteSessionStore::open(state.join("sessions.db"))
.expect("session store"),
);
let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> = Arc::new(
meerkat_runtime::store::SqliteRuntimeStore::new(state.join("runtime.sqlite"))
.expect("runtime store"),
);
let blob_store: Arc<dyn meerkat_core::BlobStore> =
Arc::new(meerkat_store::MemoryBlobStore::new());
let factory = meerkat::AgentFactory::new(&state).comms(true);
let mut inner_builder =
meerkat::FactoryAgentBuilder::new(factory, meerkat::Config::default());
inner_builder.default_session_store = Some(Arc::new(meerkat_store::StoreAdapter::new(
session_store.clone(),
)));
inner_builder.default_blob_store = Some(blob_store.clone());
let adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
Arc::clone(&runtime_store),
Arc::clone(&blob_store),
));
let concrete = Arc::new(meerkat_session::PersistentSessionService::new(
inner_builder,
16,
session_store.clone(),
runtime_store,
blob_store,
));
let gate_armed = Arc::new(AtomicBool::new(false));
let in_call = Arc::new(AtomicBool::new(false));
let release = Arc::new(tokio::sync::Notify::new());
let definition = meerkat_mob::MobDefinition::from_toml(&format!(
r#"
[mob]
id = "{mob_id}"
[profiles.worker]
model = "gpt-5.5"
[profiles.worker.tools]
comms = true
"#
))
.expect("mob definition");
let mob_spec = crate::mob_handle_runtime::MobBootstrapSpec::new(
definition,
meerkat_mob::MobStorage::in_memory(),
concrete.clone(),
)
.with_session_runtime_adapter(adapter)
.with_options(crate::mob_handle_runtime::MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: Some(Arc::new(GatedUsageLlmClient {
input_tokens: 5_000,
gate_armed: Arc::clone(&gate_armed),
in_call: Arc::clone(&in_call),
release: Arc::clone(&release),
})),
});
let mut runtime = crate::UnifiedRuntime::bootstrap(
mob_spec,
crate::MobKitConfig {
modules: vec![],
discovery: crate::DiscoverySpec {
namespace: mob_id.to_string(),
modules: vec![],
},
pre_spawn: vec![],
},
Duration::from_secs(5),
)
.await
.expect("bootstrap unified runtime");
let handle = runtime.mob_handle();
let roster_id = crate::member_comms_id::mob_member_id_str(member).into_owned();
let roster_identity = meerkat_mob::ids::AgentIdentity::from(roster_id.clone());
let durable_identity = IdentityRuntime::identity_for_generated_member_alias(member)
.expect("generated alias identity");
let mut member_labels = std::collections::BTreeMap::new();
member_labels.insert(
"agent_identity".to_string(),
durable_identity.as_str().to_string(),
);
handle
.ensure_member(
meerkat_mob::SpawnMemberSpec::new(
meerkat_mob::ProfileName::from("worker"),
roster_identity.clone(),
)
.with_labels(member_labels),
)
.await
.expect("spawn identity-first member");
handle
.wait_for_members_kickoff_complete(
std::slice::from_ref(&roster_identity),
Some(Duration::from_secs(5)),
)
.await
.expect("member kickoff settled");
let member_session = handle
.resolve_bridge_session_id_observation(&roster_identity)
.await
.expect("member session id");
let public_member_alias =
crate::member_comms_id::runtime_alias_str(&roster_id).into_owned();
let identity = durable_identity;
let continuity_store =
Arc::new(LocalContinuityStore::in_memory().expect("continuity store"));
let lease_provider = Arc::new(LocalLeaseProvider::new());
let lease_results = LeaseProvider::acquire_leases(
lease_provider.as_ref(),
std::slice::from_ref(&identity),
"operator-verb-test",
)
.await
.expect("acquire identity lease");
let lease = match lease_results.get(&identity) {
Some(LeaseAcquireResult::Acquired(lease)) => lease.clone(),
other => panic!("expected acquired identity lease, got {other:?}"),
};
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: crate::identity_first::AgentRuntimeId::parse(&public_member_alias)
.expect("runtime alias"),
session_id: member_session,
generation: ContinuityGeneration::new(0),
checkpoint_version: crate::identity_first::CheckpointVersion::new(0),
};
ContinuityStore::upsert_continuity_record(
continuity_store.as_ref(),
&record,
lease.fencing_token,
)
.await
.expect("persist identity continuity");
let bridge = MobSessionBridge::with_session_service(handle.clone(), concrete.clone());
let floors = bridge.compaction_floors();
let identity_runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store,
lease_provider,
runtime_instance_id: "operator-verb-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(Arc::new(bridge)),
default_timeout: None,
}));
identity_runtime
.register(
worker_spec(&identity),
IdentityLifecycleState::Active,
Some(record),
Some(lease),
)
.await;
runtime.attach_identity_first_context(Arc::new(
crate::identity_first::IdentityFirstRuntimeContext::new(
Arc::clone(&identity_runtime),
Arc::new(EmptyRoster),
None,
None,
Some(handle.definition().clone()),
),
));
OperatorVerbHarness {
_temp_dir: temp_dir,
runtime,
concrete,
identity_runtime,
floors,
identity,
member_alias: public_member_alias,
gate_armed,
in_call,
release,
}
}
async fn rpc(harness: &OperatorVerbHarness, method: &str, params: Value) -> Value {
let ctx = harness.identity_ctx();
let raw = handle_unified_rpc_json(
&harness.runtime,
&serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": method,
"params": params,
})
.to_string(),
Duration::from_mins(1),
None,
Some(&ctx),
)
.await;
serde_json::from_str(&raw).expect("json-rpc response")
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_member_forces_one_compaction_and_restores_the_profile() {
let harness = operator_verb_harness("rt:worker:main:0", "operator-compact-verb").await;
let fat = "seeded transcript ballast ".repeat(160);
for turn in 0..8 {
harness.run_turn(format!("turn {turn}: {fat}")).await;
}
let messages_before = harness.settled_transcript_count().await;
assert!(
messages_before >= 16,
"seed must materialize a fat transcript, got {messages_before}"
);
let response = rpc(
&harness,
"mobkit/compact_member",
serde_json::json!({
"identity": harness.member_alias,
"floor_tokens": 256,
"timeout_ms": 30_000,
}),
)
.await;
assert!(
response["error"].is_null(),
"compact_member must succeed: {response:#?}"
);
let result = &response["result"];
assert_eq!(
result["compaction_applied"],
Value::Bool(true),
"{result:#?}"
);
let reported_before = result["messages_before"].as_u64().expect("messages_before");
let reported_after = result["messages_after"].as_u64().expect("messages_after");
assert!(
reported_after < reported_before,
"compaction must shrink the transcript: {result:#?}"
);
assert!(
result["last_rewrite_reason"]
.as_str()
.is_some_and(|reason| reason.to_lowercase().contains("compact")),
"the landed rewrite must be compaction-semantic: {result:#?}"
);
assert!(
harness.floors.get(&harness.identity).is_none(),
"the floor registry must be disarmed after the verb"
);
let count_after_verb = harness.settled_transcript_count().await;
harness.run_turn("post-verb probe turn".to_string()).await;
crate::test_wait::poll_until(
&format!(
"the restored build appended to the durable transcript instead of re-compacting \
(still {count_after_verb} messages)"
),
crate::test_wait::STRUCTURAL_BACKSTOP,
async || harness.transcript_facts().await.0 > count_after_verb,
)
.await;
let response = rpc(
&harness,
"mobkit/compact_member",
serde_json::json!({ "identity": "nobody:here" }),
)
.await;
assert_eq!(
response["error"]["code"],
serde_json::json!(-32001),
"an unowned identity must surface the typed unknown-identity refusal: {response:#?}"
);
let _ = harness.runtime.mob_handle().stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_member_timeout_names_the_turn_fate_and_restores_the_member() {
let harness = operator_verb_harness("rt:worker:main:0", "operator-compact-timeout").await;
let fat = "seeded transcript ballast ".repeat(160);
for turn in 0..4 {
harness.run_turn(format!("turn {turn}: {fat}")).await;
}
let response = rpc(
&harness,
"mobkit/compact_member",
serde_json::json!({
"identity": harness.member_alias,
"floor_tokens": 256,
"timeout_ms": 1,
}),
)
.await;
let message = response["error"]["message"]
.as_str()
.unwrap_or_else(|| panic!("timeout must surface a typed error: {response:#?}"));
assert!(
message.contains("did not complete within 1ms"),
"the error must name the exhausted deadline: {message}"
);
assert!(
message.contains("quiesced by the rollback rebuild")
|| message.contains("may still be running on the floored build"),
"the error must state the in-flight turn's actual fate: {message}"
);
assert!(
!message.contains("did not complete: "),
"the pre-fix message shape (bare wait error, no turn fate) must be gone: {message}"
);
assert!(
harness.floors.get(&harness.identity).is_none(),
"the floor registry must be disarmed after a timed-out verb"
);
let count_before_probe = harness.settled_transcript_count().await;
harness
.run_turn("post-timeout probe turn".to_string())
.await;
crate::test_wait::poll_until(
&format!(
"the rolled-back member accepted a turn and appended to the durable transcript \
(still {count_before_probe} messages)"
),
crate::test_wait::STRUCTURAL_BACKSTOP,
async || harness.transcript_facts().await.0 > count_before_probe,
)
.await;
let _ = harness.runtime.mob_handle().stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn bound_member_transcript_commits_a_pair_safe_cut() {
let harness = operator_verb_harness("rt:worker:main:0", "operator-bound-verb").await;
let service: Arc<dyn crate::memory::hygienist::TranscriptEditSessionService> =
Arc::clone(&harness.concrete) as _;
harness
.run_turn("seed one committed turn".to_string())
.await;
let session_id = harness
.identity_runtime
.status(&harness.identity)
.await
.expect("identity status")
.session_id
.expect("identity session");
let seeded_len = settled_transcript_len(&service, &session_id).await;
let (assistant, results) = tool_use_pair("call-straddle");
let fixture = vec![user("older"), assistant, results, user("tail")];
let fixture_len = fixture.len();
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
let request = meerkat_core::service::SessionTranscriptRewriteRequest {
selection: meerkat_core::TranscriptRewriteSelection::MessageRange {
start: seeded_len,
end: seeded_len,
},
replacement: fixture.clone(),
reason: meerkat_core::TranscriptRewriteReason::new("test_seed"),
actor: Some("operator-verb-test".to_string()),
expected_parent_revision: None,
running_behavior: meerkat_core::TranscriptEditRunningBehavior::default(),
};
match service
.rewrite_session_transcript(&session_id, request)
.await
{
Ok(_) => break,
Err(SessionError::Busy { .. }) if std::time::Instant::now() < deadline => {
tokio::time::sleep(Duration::from_millis(50)).await;
}
Err(err) => panic!("fixture seed rewrite failed: {err}"),
}
}
let staged_len = transcript_len(&service, &session_id).await;
assert_eq!(
staged_len,
seeded_len + fixture_len,
"the seeded fixture must be the whole of the transcript growth: expected the \
{seeded_len} settled rows plus the {fixture_len} fixture rows, saw {staged_len}. A \
different count means rows landed alongside the fixture and every index below is \
derived from a stale premise."
);
let response = rpc(
&harness,
"mobkit/bound_member_transcript",
serde_json::json!({
"identity": harness.member_alias,
"keep_last": 2,
}),
)
.await;
assert!(
response["error"].is_null(),
"bound_member_transcript must succeed on an idle session: {response:#?}"
);
let result = &response["result"];
assert_eq!(result["bounded"], Value::Bool(true), "{result:#?}");
assert_eq!(
result["removed"],
serde_json::json!(seeded_len + 1),
"the cut must walk back off the tool_results row: {result:#?}"
);
assert!(result["revision"].as_str().is_some(), "{result:#?}");
let page = service
.read_history(
&session_id,
meerkat_core::service::SessionHistoryQuery {
offset: 0,
limit: None,
},
)
.await
.expect("read bounded transcript");
assert_eq!(page.messages.len(), 4, "marker + intact pair + tail");
assert!(
matches!(page.messages[0], Message::SystemNotice(_)),
"bounded transcript must lead with the operator marker"
);
assert!(
matches!(page.messages[1], Message::BlockAssistant(_))
&& matches!(page.messages[2], Message::ToolResults { .. }),
"the straddled tool pair must survive whole"
);
let _ = harness.runtime.mob_handle().stop().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn bound_member_transcript_refuses_running_sessions_typed() {
let harness = operator_verb_harness("rt:worker:main:0", "operator-bound-busy").await;
harness
.run_turn("seed one committed turn".to_string())
.await;
harness.gate_armed.store(true, Ordering::SeqCst);
let admission = harness
.identity_runtime
.send_admission_tracked(
&harness.identity,
None,
&meerkat_core::ContentInput::Text("held turn".to_string()),
meerkat_core::types::HandlingMode::Queue,
None,
)
.await
.expect("held turn admitted");
let deadline = std::time::Instant::now() + Duration::from_mins(1);
while !harness.in_call.load(Ordering::SeqCst) {
assert!(
std::time::Instant::now() < deadline,
"the held turn never reached the LLM call: the gate was armed and the turn was \
admitted, but `in_call` was never set"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
let response = rpc(
&harness,
"mobkit/bound_member_transcript",
serde_json::json!({
"identity": harness.member_alias,
"keep_last": 1,
}),
)
.await;
assert_eq!(
response["error"]["code"],
serde_json::json!(OPERATOR_SESSION_BUSY_CODE),
"a running session must surface the typed Busy refusal: {response:#?}"
);
assert!(
response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("quiesce")),
"the refusal must document quiesce-first: {response:#?}"
);
harness.gate_armed.store(false, Ordering::SeqCst);
harness.release.notify_waiters();
harness
.identity_runtime
.wait_for_completion(
&harness.identity,
admission.completion_baseline,
Duration::from_secs(30),
)
.await
.expect("held turn completed after release");
let _ = harness.runtime.mob_handle().stop().await;
}
#[test]
fn optional_u64_param_rejects_zero_and_non_integers() {
let params = serde_json::json!({ "floor_tokens": 0 });
assert!(optional_u64_param(¶ms, "floor_tokens").is_err());
let params = serde_json::json!({ "floor_tokens": "many" });
assert!(optional_u64_param(¶ms, "floor_tokens").is_err());
let params = serde_json::json!({});
assert_eq!(optional_u64_param(¶ms, "floor_tokens").unwrap(), None);
let params = serde_json::json!({ "floor_tokens": 2048 });
assert_eq!(
optional_u64_param(¶ms, "floor_tokens").unwrap(),
Some(2048)
);
}
}