#![cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
#![allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
use std::sync::Arc;
use meerkat::surface::{
build_runtime_backed_service_with_default_reconfigure_host, default_persistent_executor,
};
use meerkat::{
AgentFactory, Config, CreateSessionRequest, FactoryAgentBuilder, PersistentSessionService,
};
use meerkat_client::TestClient;
use meerkat_core::SessionBuildOptions;
use meerkat_core::image_generation::{
SwitchTurnDuration, SwitchTurnIntent, SwitchTurnOrigin, SwitchTurnReasonTextDisposition,
SwitchTurnRequestId,
};
use meerkat_core::lifecycle::RunId;
use meerkat_core::lifecycle::run_primitive::ModelId;
use meerkat_core::session::model_routing_control::SessionModelRoutingControlRecord;
use meerkat_runtime::completion::CompletionOutcome;
use meerkat_runtime::{Input, MeerkatMachine, PromptInput};
use tokio::time::Duration;
async fn build_service(
root: &std::path::Path,
backend: meerkat_store::RealmBackend,
) -> (
Arc<PersistentSessionService<FactoryAgentBuilder>>,
Arc<MeerkatMachine>,
) {
let (_manifest, persistence) = meerkat::open_realm_persistence_in(
root,
"brain-swap-realm",
Some(backend),
Some(meerkat_store::RealmOrigin::Explicit),
)
.await
.expect("open realm persistence");
let factory = AgentFactory::new(root.join("sessions"));
let mut builder = FactoryAgentBuilder::new(factory, Config::default());
builder.default_llm_client = Some(Arc::new(TestClient::for_provider(
meerkat_core::Provider::OpenAI,
)));
let (service, adapter) = build_runtime_backed_service_with_default_reconfigure_host(
builder,
4,
persistence,
root.join("config_state.json"),
);
(service, adapter)
}
fn create_request() -> CreateSessionRequest {
CreateSessionRequest {
injected_context: Vec::new(),
model: "gpt-5.4".to_string(),
prompt: meerkat_core::ContentInput::Text(String::new()),
system_prompt: meerkat::SystemPromptOverride::Set("pre-dequeue contract".to_string()),
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions::default()),
labels: None,
}
}
async fn run_prompt(
adapter: &Arc<MeerkatMachine>,
session_id: &meerkat::SessionId,
prompt: &str,
) -> CompletionOutcome {
let (_outcome, handle) = adapter
.accept_input_with_completion(session_id, Input::Prompt(PromptInput::new(prompt, None)))
.await
.expect("accept prompt input");
let handle = handle.expect("completion handle");
tokio::time::timeout(Duration::from_secs(20), handle.wait())
.await
.expect("prompt should settle in time")
.expect("completion waiter should resolve")
}
fn until_changed_model_intent(target: &str) -> SwitchTurnIntent {
SwitchTurnIntent {
target_model: ModelId::new(target),
duration: SwitchTurnDuration::UntilChanged,
origin: SwitchTurnOrigin::Model {
reason: SwitchTurnReasonTextDisposition::NotProvided,
},
}
}
async fn uncommitted_origin_blocks_the_next_input(backend: meerkat_store::RealmBackend) {
let temp = tempfile::tempdir().expect("tempdir");
let (service, adapter) = build_service(temp.path(), backend).await;
let mut session = meerkat::Session::new();
let session_id = session.id().clone();
session
.append_model_routing_control_record(
SessionModelRoutingControlRecord::request(
SwitchTurnRequestId::new(uuid::Uuid::from_bytes([42u8; 16])),
RunId::new(),
until_changed_model_intent("gpt-5.5"),
)
.expect("representable durable request"),
)
.expect("seed committed request");
let service_for_executor = Arc::clone(&service);
let adapter_for_executor = Arc::clone(&adapter);
Box::pin(meerkat::surface::materialize_session(
&service,
&adapter,
session,
create_request(),
move |session_id| {
default_persistent_executor(service_for_executor, adapter_for_executor, session_id)
},
))
.await
.expect("materialize session with committed unprovable request");
let second = run_prompt(&adapter, &session_id, "second").await;
assert!(
!matches!(second, CompletionOutcome::Completed(_)),
"an unprovable committed handoff must block the next input, got {second:?}"
);
}
#[tokio::test]
async fn head_canonical_uncommitted_origin_blocks_the_next_input() {
uncommitted_origin_blocks_the_next_input(meerkat_store::RealmBackend::Sqlite).await;
}