use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use indexmap::IndexMap;
use meerkat_core::ToolCallId;
use meerkat_core::ToolDispatchOutcome;
use meerkat_core::ToolError;
use meerkat_core::live_adapter::{
LiveAdapter, LiveAdapterCommand, LiveAdapterError, LiveAdapterErrorCode,
LiveAdapterObservation, LiveAdapterStatus, LiveInputChunk, LiveToolResult,
};
use meerkat_core::types::{SessionId, StopReason, ToolCall, ToolName, ToolResult, Usage};
use meerkat_core::{
RealtimeTranscriptApplyOutcome, RealtimeTranscriptEvent, RealtimeUserContentApplyOutcome,
};
use tokio::sync::Mutex;
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct LiveChannelId(String);
impl LiveChannelId {
#[must_use]
pub fn new(id: impl Into<String>) -> Self {
Self(id.into())
}
#[must_use]
pub fn random_uuid() -> Self {
Self(uuid::Uuid::new_v4().to_string())
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LiveRefreshQueueAcceptance {
channel_id: String,
acceptance_sequence: u64,
}
impl LiveRefreshQueueAcceptance {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn acceptance_sequence(&self) -> u64 {
self.acceptance_sequence
}
fn from_host_queue_acceptance(
channel_id: impl Into<String>,
acceptance_sequence: u64,
) -> Option<Self> {
let channel_id = channel_id.into();
if channel_id.is_empty() || acceptance_sequence == 0 {
return None;
}
Some(Self {
channel_id,
acceptance_sequence,
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum LiveCommandAcceptanceKind {
SendInput,
CommitInput,
Interrupt,
TruncateAssistantOutput,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LiveCommandQueueAcceptance {
channel_id: String,
kind: LiveCommandAcceptanceKind,
acceptance_sequence: u64,
}
impl LiveCommandQueueAcceptance {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn kind(&self) -> LiveCommandAcceptanceKind {
self.kind
}
#[must_use]
pub fn acceptance_sequence(&self) -> u64 {
self.acceptance_sequence
}
fn from_host_queue_acceptance(
channel_id: impl Into<String>,
kind: LiveCommandAcceptanceKind,
acceptance_sequence: u64,
) -> Option<Self> {
let channel_id = channel_id.into();
if channel_id.is_empty() || acceptance_sequence == 0 {
return None;
}
Some(Self {
channel_id,
kind,
acceptance_sequence,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LiveChannelCloseObservation {
channel_id: String,
close_sequence: u64,
}
impl LiveChannelCloseObservation {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn close_sequence(&self) -> u64 {
self.close_sequence
}
fn from_host_close_observation(
channel_id: impl Into<String>,
close_sequence: u64,
) -> Option<Self> {
let channel_id = channel_id.into();
if channel_id.is_empty() || close_sequence == 0 {
return None;
}
Some(Self {
channel_id,
close_sequence,
})
}
}
#[derive(Debug, Clone)]
pub struct LiveChannelCloseCommitAuthority {
channel_id: String,
close_sequence: u64,
consumed: Arc<AtomicBool>,
}
impl LiveChannelCloseCommitAuthority {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn close_sequence(&self) -> u64 {
self.close_sequence
}
fn consume_once(&self) -> Result<(), LiveAdapterHostError> {
self.consumed
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| LiveAdapterHostError::CloseAuthorityAlreadyConsumed)
}
#[cfg_attr(not(meerkat_internal_generated_authority_bridge), allow(dead_code))]
fn from_generated_parts(channel_id: String, close_sequence: u64) -> Self {
Self {
channel_id,
close_sequence,
consumed: Arc::new(AtomicBool::new(false)),
}
}
#[cfg(test)]
fn from_generated_test_machine(
session_id: &SessionId,
channel_id: &LiveChannelId,
close_sequence: u64,
) -> Self {
let mut authority = generated_test_machine_for_registered_session(session_id);
let channel_id_string = channel_id.as_str().to_owned();
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineMutator::apply(
&mut authority,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineInput::ResolveLiveOpenAdmission {
session_id: session_id.to_string(),
channel_id: channel_id_string.clone(),
llm_identity: generated_test_llm_identity(),
},
)
.expect("generated MeerkatMachine live-open admission");
let transition = meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineMutator::apply(
&mut authority,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineInput::RecordLiveCloseClosed {
session_id: session_id.to_string(),
channel_id: channel_id_string.clone(),
close_observation_sequence: close_sequence,
},
)
.expect("generated MeerkatMachine live-close result");
assert!(
transition.effects().iter().any(|effect| matches!(
effect,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineEffect::LiveCloseResultResolved {
channel_id: effect_channel_id,
close_observation_sequence,
..
} if *effect_channel_id == channel_id_string && *close_observation_sequence == close_sequence
)),
"generated live-close result effect"
);
Self::from_generated_parts(channel_id_string, close_sequence)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LiveChannelStatusObservation {
channel_id: String,
status: LiveAdapterStatus,
observation_sequence: u64,
}
impl LiveChannelStatusObservation {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn status(&self) -> &LiveAdapterStatus {
&self.status
}
#[must_use]
pub fn observation_sequence(&self) -> u64 {
self.observation_sequence
}
fn from_host_status_observation(
channel_id: impl Into<String>,
status: LiveAdapterStatus,
observation_sequence: u64,
) -> Option<Self> {
let channel_id = channel_id.into();
if channel_id.is_empty() || observation_sequence == 0 {
return None;
}
Some(Self {
channel_id,
status,
observation_sequence,
})
}
}
#[derive(Debug, Clone)]
pub struct LiveChannelStatusCommitAuthority {
channel_id: String,
status_observation_sequence: u64,
consumed: Arc<AtomicBool>,
}
impl LiveChannelStatusCommitAuthority {
#[must_use]
pub fn channel_id(&self) -> &str {
&self.channel_id
}
#[must_use]
pub fn status_observation_sequence(&self) -> u64 {
self.status_observation_sequence
}
fn consume_once(&self) -> Result<(), LiveAdapterHostError> {
self.consumed
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| LiveAdapterHostError::StatusAuthorityAlreadyConsumed)
}
#[cfg_attr(not(meerkat_internal_generated_authority_bridge), allow(dead_code))]
fn from_generated_parts(channel_id: String, status_observation_sequence: u64) -> Self {
Self {
channel_id,
status_observation_sequence,
consumed: Arc::new(AtomicBool::new(false)),
}
}
}
impl std::fmt::Display for LiveChannelId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum ObservationOutcome {
Noop,
StatusUpdated(LiveAdapterStatus),
TranscriptAppended,
TranscriptTruncated,
ToolCallDispatched {
provider_call_id: String,
tool_name: String,
},
ToolCallSkipped {
provider_call_id: String,
tool_name: String,
reason: ToolDispatchSkipReason,
},
ToolCallTimedOut {
provider_call_id: String,
tool_name: String,
timeout: std::time::Duration,
},
InterruptSignalled,
Terminal { code: LiveAdapterErrorCode },
CommandRejected {
code: LiveAdapterErrorCode,
message: String,
},
UserContentCommitted { observation: LiveAdapterObservation },
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ToolDispatchSkipReason {
NoDispatcher,
InvalidArguments,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LiveToolDispatchTimeout {
timeout: Duration,
}
impl LiveToolDispatchTimeout {
#[must_use]
pub fn new(timeout: Duration) -> Self {
Self { timeout }
}
#[must_use]
pub fn timeout(&self) -> Duration {
self.timeout
}
}
impl std::fmt::Display for LiveToolDispatchTimeout {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "tool dispatch timeout after {:?}", self.timeout)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct LiveTranscriptIdentity<'a> {
pub provider_item_id: Option<&'a str>,
pub previous_item_id: Option<&'a str>,
pub content_index: Option<u32>,
pub response_id: Option<&'a str>,
pub delta_id: Option<&'a str>,
}
impl<'a> LiveTranscriptIdentity<'a> {
pub fn user(
provider_item_id: Option<&'a str>,
previous_item_id: Option<&'a str>,
content_index: Option<u32>,
) -> Self {
Self {
provider_item_id,
previous_item_id,
content_index,
response_id: None,
delta_id: None,
}
}
pub fn assistant_delta(
provider_item_id: Option<&'a str>,
previous_item_id: Option<&'a str>,
content_index: Option<u32>,
response_id: Option<&'a str>,
delta_id: Option<&'a str>,
) -> Self {
Self {
provider_item_id,
previous_item_id,
content_index,
response_id,
delta_id,
}
}
pub fn assistant_final(
provider_item_id: &'a str,
previous_item_id: Option<&'a str>,
content_index: Option<u32>,
response_id: Option<&'a str>,
) -> Self {
Self {
provider_item_id: Some(provider_item_id),
previous_item_id,
content_index,
response_id,
delta_id: None,
}
}
pub fn require_delta_identity(&self) -> Result<DeltaIdentity<'a>, LiveTranscriptIdentityError> {
let response_id = self
.response_id
.ok_or(LiveTranscriptIdentityError::MissingResponseId)?;
let delta_id = self
.delta_id
.ok_or(LiveTranscriptIdentityError::MissingDeltaId)?;
let item_id = self
.provider_item_id
.ok_or(LiveTranscriptIdentityError::MissingItemId)?;
Ok(DeltaIdentity {
response_id,
delta_id,
item_id,
previous_item_id: self.previous_item_id,
content_index: self.content_index,
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DeltaIdentity<'a> {
pub response_id: &'a str,
pub delta_id: &'a str,
pub item_id: &'a str,
pub previous_item_id: Option<&'a str>,
pub content_index: Option<u32>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum LiveTranscriptIdentityError {
#[error("transcript delta is missing a required response_id")]
MissingResponseId,
#[error("transcript delta is missing a required delta_id")]
MissingDeltaId,
#[error("transcript delta is missing a required item_id")]
MissingItemId,
}
#[async_trait::async_trait]
pub trait LiveProjectionSink: Send + Sync {
async fn append_user_transcript(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError>;
async fn append_assistant_text_delta(
&self,
session_id: &SessionId,
delta: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError>;
async fn append_assistant_transcript_delta(
&self,
session_id: &SessionId,
delta: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError>;
async fn append_assistant_text_final(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError>;
async fn append_assistant_transcript_final(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError>;
async fn truncate_assistant_transcript(
&self,
session_id: &SessionId,
provider_item_id: Option<&str>,
previous_item_id: Option<&str>,
content_index: Option<u32>,
response_id: Option<&str>,
text: Option<&str>,
) -> Result<(), LiveProjectionError>;
async fn signal_turn_interrupt(
&self,
session_id: &SessionId,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError>;
async fn signal_output_audio_degraded(
&self,
session_id: &SessionId,
dropped: u64,
) -> Result<(), LiveProjectionError>;
async fn signal_turn_completed(
&self,
session_id: &SessionId,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError>;
async fn signal_terminal_error(
&self,
session_id: &SessionId,
code: LiveAdapterErrorCode,
message: &str,
) -> Result<(), LiveProjectionError>;
async fn append_realtime_transcript(
&self,
session_id: &SessionId,
event: &RealtimeTranscriptEvent,
) -> Result<RealtimeTranscriptApplyOutcome, LiveProjectionError>;
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum LiveProjectionError {
#[error("session {0} not found")]
SessionNotFound(SessionId),
#[error("projection rejected: {0}")]
Rejected(String),
#[error("session busy: {0}")]
SessionBusy(SessionId),
#[error("session not running: {0}")]
SessionNotRunning(SessionId),
#[error("session capability disabled ({code}): {message}")]
CapabilityDisabled { code: &'static str, message: String },
#[error("session error [{code}]: {message}")]
Session { code: &'static str, message: String },
}
impl LiveProjectionError {
#[must_use]
pub fn from_session_error(session_id: &SessionId, err: meerkat_core::SessionError) -> Self {
use meerkat_core::SessionError;
let code = err.code();
let message = err.to_string();
match err {
SessionError::NotFound { .. } => Self::SessionNotFound(session_id.clone()),
SessionError::Unsupported(reason) => Self::Rejected(reason),
SessionError::Busy { id } => Self::SessionBusy(id),
SessionError::NotRunning { id } => Self::SessionNotRunning(id),
SessionError::PersistenceDisabled | SessionError::CompactionDisabled => {
Self::CapabilityDisabled { code, message }
}
SessionError::Store(_)
| SessionError::Agent(_)
| SessionError::FailedWithData { .. } => Self::Session { code, message },
}
}
}
#[doc(hidden)]
#[derive(Debug, Default)]
pub struct NoOpProjectionSink;
#[async_trait::async_trait]
impl LiveProjectionSink for NoOpProjectionSink {
async fn append_user_transcript(
&self,
_session_id: &SessionId,
_text: &str,
_identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn append_assistant_text_delta(
&self,
_session_id: &SessionId,
_delta: &str,
_identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn append_assistant_transcript_delta(
&self,
_session_id: &SessionId,
_delta: &str,
_identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn append_assistant_text_final(
&self,
_session_id: &SessionId,
_text: &str,
_identity: LiveTranscriptIdentity<'_>,
_stop_reason: StopReason,
_usage: Usage,
_response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn append_assistant_transcript_final(
&self,
_session_id: &SessionId,
_text: &str,
_identity: LiveTranscriptIdentity<'_>,
_stop_reason: StopReason,
_usage: Usage,
_response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn truncate_assistant_transcript(
&self,
_session_id: &SessionId,
_provider_item_id: Option<&str>,
_previous_item_id: Option<&str>,
_content_index: Option<u32>,
_response_id: Option<&str>,
_text: Option<&str>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn signal_turn_interrupt(
&self,
_session_id: &SessionId,
_response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn signal_output_audio_degraded(
&self,
_session_id: &SessionId,
_dropped: u64,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn signal_turn_completed(
&self,
_session_id: &SessionId,
_stop_reason: StopReason,
_usage: Usage,
_response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn signal_terminal_error(
&self,
_session_id: &SessionId,
_code: LiveAdapterErrorCode,
_message: &str,
) -> Result<(), LiveProjectionError> {
Ok(())
}
async fn append_realtime_transcript(
&self,
_session_id: &SessionId,
_event: &RealtimeTranscriptEvent,
) -> Result<RealtimeTranscriptApplyOutcome, LiveProjectionError> {
Ok(RealtimeTranscriptApplyOutcome::default())
}
}
const CLOSED_CHANNEL_TTL: std::time::Duration = std::time::Duration::from_secs(60);
struct ChannelState {
session_id: SessionId,
status: LiveAdapterStatus,
status_observation_sequence: u64,
snapshot_version: u64,
refresh_acceptance_sequence: u64,
command_acceptance_sequence: u64,
close_observation_sequence: u64,
adapter: Option<Arc<dyn LiveAdapter>>,
physical_close_confirmed: bool,
retire_at: Option<std::time::Instant>,
pending_synthetic_obs: Option<LiveAdapterObservation>,
}
#[derive(Debug, Clone)]
pub struct LiveChannelOpenAuthority {
session_id: SessionId,
channel_id: LiveChannelId,
sequence: u64,
consumed: Arc<AtomicBool>,
}
#[cfg(test)]
fn generated_test_llm_identity()
-> meerkat_machine_schema::catalog::dsl::meerkat_machine::SessionLlmIdentity {
meerkat_machine_schema::catalog::dsl::meerkat_machine::SessionLlmIdentity {
model: "gpt-realtime-2".to_string(),
provider: meerkat_machine_schema::catalog::dsl::meerkat_machine::Provider::OpenAI,
self_hosted_server_id: None,
provider_params_repr: None,
auth_binding: None,
}
}
#[cfg(test)]
fn generated_test_machine_for_registered_session(
session_id: &SessionId,
) -> meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineAuthority {
use meerkat_machine_schema::catalog::dsl::meerkat_machine::{
MeerkatMachineAuthority, MeerkatMachineInput, MeerkatMachineMutator, MeerkatMachineSignal,
SessionId as DslSessionId,
};
let mut authority = MeerkatMachineAuthority::new();
authority
.apply_signal(MeerkatMachineSignal::Initialize)
.expect("initialize generated MeerkatMachine authority");
MeerkatMachineMutator::apply(
&mut authority,
MeerkatMachineInput::RegisterSession {
session_id: DslSessionId(session_id.to_string()),
},
)
.expect("register session in generated MeerkatMachine authority");
authority
}
#[cfg(test)]
fn generated_test_live_channel_status(
status: &LiveAdapterStatus,
) -> (
meerkat_machine_schema::catalog::dsl::meerkat_machine::LiveChannelPublicStatus,
Option<meerkat_machine_schema::catalog::dsl::meerkat_machine::LiveChannelDegradationReason>,
Option<String>,
) {
use meerkat_core::live_adapter::LiveDegradationReason;
use meerkat_machine_schema::catalog::dsl::meerkat_machine::{
LiveChannelDegradationReason as DslReason, LiveChannelPublicStatus as DslStatus,
};
match status {
LiveAdapterStatus::Idle => (DslStatus::Idle, None, None),
LiveAdapterStatus::Opening => (DslStatus::Opening, None, None),
LiveAdapterStatus::Ready => (DslStatus::Ready, None, None),
LiveAdapterStatus::Closing => (DslStatus::Closing, None, None),
LiveAdapterStatus::Closed => (DslStatus::Closed, None, None),
LiveAdapterStatus::Degraded { reason } => match reason {
LiveDegradationReason::RateLimited => {
(DslStatus::Degraded, Some(DslReason::RateLimited), None)
}
LiveDegradationReason::ProviderThrottled => (
DslStatus::Degraded,
Some(DslReason::ProviderThrottled),
None,
),
LiveDegradationReason::NetworkUnstable => {
(DslStatus::Degraded, Some(DslReason::NetworkUnstable), None)
}
LiveDegradationReason::Other { detail } => (
DslStatus::Degraded,
Some(DslReason::Other),
Some(detail.clone().into_owned()),
),
other => (
DslStatus::Degraded,
Some(DslReason::Unknown),
Some(format!("{other:?}")),
),
},
other => (
DslStatus::Degraded,
Some(DslReason::Unknown),
Some(format!("{other:?}")),
),
}
}
impl LiveChannelOpenAuthority {
#[must_use]
pub fn session_id(&self) -> &SessionId {
&self.session_id
}
#[must_use]
pub fn channel_id(&self) -> &LiveChannelId {
&self.channel_id
}
#[must_use]
pub fn sequence(&self) -> u64 {
self.sequence
}
fn consume_once(&self) -> Result<(), LiveAdapterHostError> {
self.consumed
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| LiveAdapterHostError::OpenAuthorityAlreadyConsumed)
}
#[cfg_attr(not(meerkat_internal_generated_authority_bridge), allow(dead_code))]
fn from_generated_parts(
session_id: SessionId,
channel_id: LiveChannelId,
sequence: u64,
) -> Self {
Self {
session_id,
channel_id,
sequence,
consumed: Arc::new(AtomicBool::new(false)),
}
}
#[cfg(test)]
fn from_generated_test_machine(session_id: SessionId, channel_id: LiveChannelId) -> Self {
let mut authority = generated_test_machine_for_registered_session(&session_id);
let channel_id_string = channel_id.as_str().to_owned();
let transition = meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineMutator::apply(
&mut authority,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineInput::ResolveLiveOpenAdmission {
session_id: session_id.to_string(),
channel_id: channel_id_string.clone(),
llm_identity: generated_test_llm_identity(),
},
)
.expect("generated MeerkatMachine live-open admission");
let sequence = transition
.effects()
.iter()
.find_map(|effect| match effect {
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineEffect::LiveOpenAdmissionResolved {
channel_id: effect_channel_id,
admitted: true,
sequence,
..
} if *effect_channel_id == channel_id_string => Some(*sequence),
_ => None,
})
.expect("generated live-open admission effect");
Self::from_generated_parts(session_id, channel_id, sequence)
}
}
#[cfg(meerkat_internal_generated_authority_bridge)]
#[allow(improper_ctypes_definitions, unsafe_code)]
unsafe extern "Rust" {
#[link_name = concat!(
"__meerkat_runtime_generated_authority_bridge_token_is_valid_v1_live_open_admission_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
)]
fn runtime_live_open_admission_generated_authority_bridge_token_is_valid(
token: &(dyn std::any::Any + Send + Sync),
) -> bool;
#[link_name = concat!(
"__meerkat_runtime_generated_authority_bridge_token_is_valid_v1_live_close_result_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
)]
fn runtime_live_close_result_generated_authority_bridge_token_is_valid(
token: &(dyn std::any::Any + Send + Sync),
) -> bool;
#[link_name = concat!(
"__meerkat_runtime_generated_authority_bridge_token_is_valid_v1_live_channel_status_result_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
)]
fn runtime_live_channel_status_result_generated_authority_bridge_token_is_valid(
token: &(dyn std::any::Any + Send + Sync),
) -> bool;
}
#[cfg(meerkat_internal_generated_authority_bridge)]
#[doc(hidden)]
#[allow(improper_ctypes_definitions, unsafe_code)]
#[unsafe(export_name = concat!(
"__meerkat_live_runtime_generated_live_channel_open_authority_build_v1_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
))]
pub(crate) extern "Rust" fn runtime_generated_live_channel_open_authority_build(
token: &'static (dyn std::any::Any + Send + Sync),
session_id: SessionId,
channel_id: LiveChannelId,
sequence: u64,
) -> Result<LiveChannelOpenAuthority, String> {
#[allow(unsafe_code)]
let valid =
unsafe { runtime_live_open_admission_generated_authority_bridge_token_is_valid(token) };
if !valid {
return Err(
"live channel open authority requires the generated runtime admission bridge token"
.into(),
);
}
Ok(LiveChannelOpenAuthority::from_generated_parts(
session_id, channel_id, sequence,
))
}
#[cfg(meerkat_internal_generated_authority_bridge)]
#[doc(hidden)]
#[allow(improper_ctypes_definitions, unsafe_code)]
#[unsafe(export_name = concat!(
"__meerkat_live_runtime_generated_live_channel_close_commit_authority_build_v1_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
))]
pub(crate) extern "Rust" fn runtime_generated_live_channel_close_commit_authority_build(
token: &'static (dyn std::any::Any + Send + Sync),
channel_id: String,
close_sequence: u64,
) -> Result<LiveChannelCloseCommitAuthority, String> {
#[allow(unsafe_code)]
let valid =
unsafe { runtime_live_close_result_generated_authority_bridge_token_is_valid(token) };
if !valid {
return Err(
"live channel close commit authority requires the generated runtime close bridge token"
.into(),
);
}
Ok(LiveChannelCloseCommitAuthority::from_generated_parts(
channel_id,
close_sequence,
))
}
#[cfg(meerkat_internal_generated_authority_bridge)]
#[doc(hidden)]
#[allow(improper_ctypes_definitions, unsafe_code)]
#[unsafe(export_name = concat!(
"__meerkat_live_runtime_generated_live_channel_status_commit_authority_build_v1_",
env!("MEERKAT_GENERATED_AUTHORITY_BRIDGE_SYMBOL_SUFFIX")
))]
pub(crate) extern "Rust" fn runtime_generated_live_channel_status_commit_authority_build(
token: &'static (dyn std::any::Any + Send + Sync),
channel_id: String,
status_observation_sequence: u64,
) -> Result<LiveChannelStatusCommitAuthority, String> {
#[allow(unsafe_code)]
let valid = unsafe {
runtime_live_channel_status_result_generated_authority_bridge_token_is_valid(token)
};
if !valid {
return Err(
"live channel status commit authority requires the generated runtime status bridge token"
.into(),
);
}
Ok(LiveChannelStatusCommitAuthority::from_generated_parts(
channel_id,
status_observation_sequence,
))
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum LiveAdapterHostError {
#[error("channel {0} not found")]
ChannelNotFound(LiveChannelId),
#[error("session {0} not found")]
SessionNotFound(SessionId),
#[error("channel {0} is not ready (status: {1:?})")]
ChannelNotReady(LiveChannelId, LiveAdapterStatus),
#[error("session {0} already has an active channel")]
SessionAlreadyBound(SessionId),
#[error("live channel open lacks generated admission authority")]
OpenNotAuthorized,
#[error("live channel open authority was already consumed")]
OpenAuthorityAlreadyConsumed,
#[error("live channel close lacks generated commit authority")]
CloseNotAuthorized,
#[error("live channel close authority was already consumed")]
CloseAuthorityAlreadyConsumed,
#[error("live channel status lacks generated commit authority")]
StatusNotAuthorized,
#[error("live channel status authority was already consumed")]
StatusAuthorityAlreadyConsumed,
#[error("no adapter attached to channel {0}")]
NoAdapter(LiveChannelId),
#[error("unsupported host command: {0}")]
UnsupportedCommand(&'static str),
#[error(transparent)]
AdapterError(#[from] LiveAdapterError),
#[error("projection sink error: {0}")]
ProjectionError(#[from] LiveProjectionError),
}
impl LiveAdapterHostError {
#[must_use]
pub fn reason_code(&self) -> &'static str {
match self {
Self::ChannelNotFound(_) => "channel_not_found",
Self::SessionNotFound(_) => "session_not_found",
Self::ChannelNotReady(..) => "channel_not_ready",
Self::SessionAlreadyBound(_) => "session_already_bound",
Self::OpenNotAuthorized => "open_not_authorized",
Self::OpenAuthorityAlreadyConsumed => "open_authority_already_consumed",
Self::CloseNotAuthorized => "close_not_authorized",
Self::CloseAuthorityAlreadyConsumed => "close_authority_already_consumed",
Self::StatusNotAuthorized => "status_not_authorized",
Self::StatusAuthorityAlreadyConsumed => "status_authority_already_consumed",
Self::NoAdapter(_) => "no_adapter",
Self::UnsupportedCommand(_) => "unsupported_command",
Self::AdapterError(_) => "adapter_error",
Self::ProjectionError(_) => "projection_error",
}
}
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum ObservationRouting {
AppendTranscript,
AppendRealtimeTranscript,
DispatchToolCall {
provider_call_id: String,
tool_name: String,
},
SignalInterrupt,
UpdateStatus(LiveAdapterStatus),
TerminalError,
CommandRejection,
Noop,
}
#[async_trait::async_trait]
pub trait LiveToolDispatcher: Send + Sync {
async fn dispatch_live_tool_call(
&self,
session_id: &SessionId,
call: ToolCall,
) -> Result<ToolDispatchOutcome, LiveToolDispatchError>;
}
#[derive(Debug, Clone, thiserror::Error)]
#[non_exhaustive]
pub enum LiveToolDispatchError {
#[error(transparent)]
Tool(#[from] ToolError),
#[error("live tool dispatch target session {0} not found")]
SessionNotFound(SessionId),
#[error("live tool dispatch target session {0} is busy")]
SessionBusy(SessionId),
#[error("live tool dispatch target session {0} is not running")]
SessionNotRunning(SessionId),
#[error("live tool dispatch capability disabled ({code}): {message}")]
CapabilityDisabled { code: &'static str, message: String },
#[error("live tool dispatch rejected: {0}")]
Rejected(String),
#[error("live tool dispatch session error [{code}]: {message}")]
Session { code: &'static str, message: String },
}
impl LiveToolDispatchError {
#[must_use]
pub fn from_session_error(session_id: &SessionId, err: meerkat_core::SessionError) -> Self {
use meerkat_core::SessionError;
let code = err.code();
let message = err.to_string();
match err {
SessionError::NotFound { .. } => Self::SessionNotFound(session_id.clone()),
SessionError::Unsupported(reason) => Self::Rejected(reason),
SessionError::Busy { id } => Self::SessionBusy(id),
SessionError::NotRunning { id } => Self::SessionNotRunning(id),
SessionError::PersistenceDisabled | SessionError::CompactionDisabled => {
Self::CapabilityDisabled { code, message }
}
SessionError::Store(_)
| SessionError::Agent(_)
| SessionError::FailedWithData { .. } => Self::Session { code, message },
}
}
}
pub struct LiveAdapterHost {
inner: Mutex<HostInner>,
projection_sink: Arc<dyn LiveProjectionSink>,
tool_dispatcher: std::sync::Mutex<Option<Arc<dyn LiveToolDispatcher>>>,
tool_timeout: Duration,
}
pub const DEFAULT_LIVE_TOOL_TIMEOUT: Duration = Duration::from_secs(30);
struct HostInner {
channels: IndexMap<LiveChannelId, ChannelState>,
by_session: HashMap<SessionId, LiveChannelId>,
}
impl LiveAdapterHost {
#[must_use]
pub fn new(projection_sink: Arc<dyn LiveProjectionSink>) -> Self {
Self {
inner: Mutex::new(HostInner {
channels: IndexMap::new(),
by_session: HashMap::new(),
}),
projection_sink,
tool_dispatcher: std::sync::Mutex::new(None),
tool_timeout: DEFAULT_LIVE_TOOL_TIMEOUT,
}
}
#[must_use]
pub fn with_tool_timeout(mut self, timeout: Duration) -> Self {
self.tool_timeout = timeout;
self
}
#[must_use]
pub fn tool_timeout(&self) -> Duration {
self.tool_timeout
}
#[must_use]
pub fn with_live_tool_dispatcher(self, dispatcher: Arc<dyn LiveToolDispatcher>) -> Self {
self.set_live_tool_dispatcher(dispatcher);
self
}
pub fn set_live_tool_dispatcher(&self, dispatcher: Arc<dyn LiveToolDispatcher>) {
if let Ok(mut slot) = self.tool_dispatcher.lock() {
*slot = Some(dispatcher);
}
}
fn load_dispatcher(&self) -> Option<Arc<dyn LiveToolDispatcher>> {
self.tool_dispatcher
.lock()
.ok()
.and_then(|slot| slot.as_ref().map(Arc::clone))
}
pub async fn open_channel_with_authority(
&self,
authority: &LiveChannelOpenAuthority,
) -> Result<LiveChannelId, LiveAdapterHostError> {
authority.consume_once()?;
self.open_channel_after_generated_authority(
authority.session_id().clone(),
authority.channel_id().clone(),
)
.await
}
#[cfg(test)]
pub(crate) async fn open_channel_with_generated_test_machine_authority(
&self,
session_id: SessionId,
) -> Result<LiveChannelId, LiveAdapterHostError> {
let channel_id = LiveChannelId::random_uuid();
let authority =
LiveChannelOpenAuthority::from_generated_test_machine(session_id, channel_id);
self.open_channel_with_authority(&authority).await
}
async fn open_channel_after_generated_authority(
&self,
session_id: SessionId,
channel_id: LiveChannelId,
) -> Result<LiveChannelId, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
if let Some(existing) = inner.by_session.get(&session_id).cloned()
&& let Some(channel) = inner.channels.get(&existing)
&& channel.retire_at.is_none()
{
return Err(LiveAdapterHostError::SessionAlreadyBound(session_id));
}
inner.channels.insert(
channel_id.clone(),
ChannelState {
session_id: session_id.clone(),
status: LiveAdapterStatus::Opening,
status_observation_sequence: 0,
snapshot_version: 0,
refresh_acceptance_sequence: 0,
command_acceptance_sequence: 0,
close_observation_sequence: 0,
adapter: None,
physical_close_confirmed: true,
retire_at: None,
pending_synthetic_obs: None,
},
);
inner.by_session.insert(session_id, channel_id.clone());
Ok(channel_id)
}
pub async fn attach_adapter(
&self,
channel_id: &LiveChannelId,
adapter: Arc<dyn LiveAdapter>,
) -> Result<(), LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if channel.retire_at.is_some() {
return Err(LiveAdapterHostError::ChannelNotFound(channel_id.clone()));
}
channel.adapter = Some(adapter);
channel.physical_close_confirmed = false;
Ok(())
}
fn command_acceptance_kind(
command: &LiveAdapterCommand,
) -> Result<LiveCommandAcceptanceKind, LiveAdapterHostError> {
match command {
LiveAdapterCommand::SendInput { .. } => Ok(LiveCommandAcceptanceKind::SendInput),
LiveAdapterCommand::CommitInput { .. } => Ok(LiveCommandAcceptanceKind::CommitInput),
LiveAdapterCommand::Interrupt => Ok(LiveCommandAcceptanceKind::Interrupt),
LiveAdapterCommand::TruncateAssistantOutput { .. } => {
Ok(LiveCommandAcceptanceKind::TruncateAssistantOutput)
}
LiveAdapterCommand::Refresh { .. } => Err(LiveAdapterHostError::UnsupportedCommand(
"refresh commands must use LiveAdapterHost::enqueue_refresh",
)),
_ => Err(LiveAdapterHostError::UnsupportedCommand(
"live command has no public queue-acceptance authority",
)),
}
}
async fn record_command_queue_acceptance(
&self,
channel_id: &LiveChannelId,
kind: LiveCommandAcceptanceKind,
) -> Result<LiveCommandQueueAcceptance, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.command_acceptance_sequence = channel.command_acceptance_sequence.saturating_add(1);
LiveCommandQueueAcceptance::from_host_queue_acceptance(
channel_id.as_str().to_owned(),
kind,
channel.command_acceptance_sequence,
)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn send_command(
&self,
channel_id: &LiveChannelId,
command: LiveAdapterCommand,
) -> Result<(), LiveAdapterHostError> {
if matches!(&command, LiveAdapterCommand::Refresh { .. }) {
return Err(LiveAdapterHostError::UnsupportedCommand(
"refresh commands must use LiveAdapterHost::enqueue_refresh",
));
}
let adapter = self
.adapter_for(channel_id, false)
.await?;
adapter.send_command(command).await?;
Ok(())
}
pub async fn send_command_observed(
&self,
channel_id: &LiveChannelId,
command: LiveAdapterCommand,
) -> Result<LiveCommandQueueAcceptance, LiveAdapterHostError> {
if matches!(&command, LiveAdapterCommand::Refresh { .. }) {
return Err(LiveAdapterHostError::UnsupportedCommand(
"refresh commands must use LiveAdapterHost::enqueue_refresh",
));
}
let acceptance_kind = Self::command_acceptance_kind(&command)?;
let adapter = self
.adapter_for(channel_id, false)
.await?;
adapter.send_command(command).await?;
self.record_command_queue_acceptance(channel_id, acceptance_kind)
.await
}
pub async fn enqueue_refresh(
&self,
channel_id: &LiveChannelId,
snapshot: meerkat_core::live_adapter::LiveProjectionSnapshot,
) -> Result<LiveRefreshQueueAcceptance, LiveAdapterHostError> {
let adapter = self
.adapter_for(channel_id, false)
.await?;
adapter
.send_command(LiveAdapterCommand::Refresh { snapshot })
.await?;
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.refresh_acceptance_sequence = channel.refresh_acceptance_sequence.saturating_add(1);
LiveRefreshQueueAcceptance::from_host_queue_acceptance(
channel_id.as_str().to_owned(),
channel.refresh_acceptance_sequence,
)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn send_input(
&self,
channel_id: &LiveChannelId,
chunk: LiveInputChunk,
) -> Result<(), LiveAdapterHostError> {
let adapter = self
.adapter_for(channel_id, true)
.await?;
adapter
.send_command(LiveAdapterCommand::SendInput { chunk })
.await?;
Ok(())
}
pub async fn send_input_observed(
&self,
channel_id: &LiveChannelId,
chunk: LiveInputChunk,
) -> Result<LiveCommandQueueAcceptance, LiveAdapterHostError> {
let adapter = self
.adapter_for(channel_id, true)
.await?;
adapter
.send_command(LiveAdapterCommand::SendInput { chunk })
.await?;
self.record_command_queue_acceptance(channel_id, LiveCommandAcceptanceKind::SendInput)
.await
}
pub async fn submit_tool_result(
&self,
channel_id: &LiveChannelId,
result: LiveToolResult,
) -> Result<(), LiveAdapterHostError> {
self.send_command(channel_id, LiveAdapterCommand::SubmitToolResult { result })
.await
}
pub async fn submit_tool_error(
&self,
channel_id: &LiveChannelId,
call_id: ToolCallId,
error: String,
) -> Result<(), LiveAdapterHostError> {
self.send_command(
channel_id,
LiveAdapterCommand::SubmitToolError { call_id, error },
)
.await
}
pub async fn submit_tool_dispatch_timeout(
&self,
channel_id: &LiveChannelId,
call_id: ToolCallId,
timeout: LiveToolDispatchTimeout,
) -> Result<(), LiveAdapterHostError> {
self.submit_tool_error(channel_id, call_id, timeout.to_string())
.await
}
pub async fn next_observation_raw(
&self,
channel_id: &LiveChannelId,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterHostError> {
{
let mut inner = self.inner.lock().await;
if let Some(channel) = inner.channels.get_mut(channel_id)
&& let Some(obs) = channel.pending_synthetic_obs.take()
{
return Ok(Some(obs));
}
}
let adapter = self
.adapter_for(channel_id, false)
.await?;
match adapter.next_observation().await {
Ok(Some(obs)) => Ok(Some(obs)),
Ok(None) => {
let mut inner = self.inner.lock().await;
if let Some(channel) = inner.channels.get_mut(channel_id)
&& let Some(obs) = channel.pending_synthetic_obs.take()
{
return Ok(Some(obs));
}
Ok(None)
}
Err(err) => {
let synthetic = LiveAdapterObservation::Error {
code: LiveAdapterErrorCode::ProviderError,
message: format!("adapter read failure: {err}"),
};
Ok(Some(synthetic))
}
}
}
pub async fn apply_observation(
&self,
channel_id: &LiveChannelId,
observation: &LiveAdapterObservation,
) -> Result<ObservationOutcome, LiveAdapterHostError> {
if Self::observation_requires_generated_close(observation)
&& !self.generated_close_has_committed(channel_id).await?
{
return Err(LiveAdapterHostError::CloseNotAuthorized);
}
let routing = Self::classify_observation(observation);
let session_id = self.channel_session(channel_id).await?;
match (routing, observation) {
(ObservationRouting::Noop, _) => Ok(ObservationOutcome::Noop),
(ObservationRouting::UpdateStatus(status), _) => {
Ok(ObservationOutcome::StatusUpdated(status))
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::UserTranscriptFinal {
provider_item_id,
previous_item_id,
content_index,
text,
..
},
) => {
let identity = LiveTranscriptIdentity::user(
provider_item_id.as_deref(),
previous_item_id.as_deref(),
*content_index,
);
self.projection_sink
.append_user_transcript(&session_id, text, identity)
.await?;
Ok(ObservationOutcome::TranscriptAppended)
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::AssistantTextDelta {
provider_item_id,
previous_item_id,
content_index,
response_id,
delta_id,
delta,
..
},
) => {
let identity = LiveTranscriptIdentity::assistant_delta(
provider_item_id.as_deref(),
previous_item_id.as_deref(),
*content_index,
response_id.as_deref(),
delta_id.as_deref(),
);
self.projection_sink
.append_assistant_text_delta(&session_id, delta, identity)
.await?;
Ok(ObservationOutcome::TranscriptAppended)
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::AssistantTranscriptDelta {
provider_item_id,
previous_item_id,
content_index,
response_id,
delta_id,
delta,
..
},
) => {
let identity = LiveTranscriptIdentity::assistant_delta(
provider_item_id.as_deref(),
previous_item_id.as_deref(),
*content_index,
response_id.as_deref(),
delta_id.as_deref(),
);
self.projection_sink
.append_assistant_transcript_delta(&session_id, delta, identity)
.await?;
Ok(ObservationOutcome::TranscriptAppended)
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::AssistantTranscriptFinal {
provider_item_id,
previous_item_id,
content_index,
response_id,
text,
stop_reason,
usage,
..
},
) => {
let identity = LiveTranscriptIdentity::assistant_final(
provider_item_id,
previous_item_id.as_deref(),
*content_index,
response_id.as_deref(),
);
self.projection_sink
.append_assistant_transcript_final(
&session_id,
text,
identity,
*stop_reason,
usage.clone(),
response_id.as_deref(),
)
.await?;
Ok(ObservationOutcome::TranscriptAppended)
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::AssistantTranscriptTruncated {
provider_item_id,
previous_item_id,
content_index,
response_id,
text,
},
) => {
self.projection_sink
.truncate_assistant_transcript(
&session_id,
provider_item_id.as_deref(),
previous_item_id.as_deref(),
*content_index,
response_id.as_deref(),
text.as_deref(),
)
.await?;
Ok(ObservationOutcome::TranscriptTruncated)
}
(
ObservationRouting::AppendTranscript,
LiveAdapterObservation::TurnCompleted {
response_id,
stop_reason,
usage,
},
) => {
self.projection_sink
.signal_turn_completed(
&session_id,
*stop_reason,
usage.clone(),
response_id.as_deref(),
)
.await?;
Ok(ObservationOutcome::TranscriptAppended)
}
(
ObservationRouting::AppendRealtimeTranscript,
LiveAdapterObservation::RealtimeTranscript { event },
) => {
let outcome = self
.projection_sink
.append_realtime_transcript(&session_id, event)
.await?;
match outcome.user_content {
Some(RealtimeUserContentApplyOutcome::Committed(identity))
| Some(RealtimeUserContentApplyOutcome::AlreadyCommitted(identity)) => {
Ok(ObservationOutcome::UserContentCommitted {
observation: LiveAdapterObservation::UserContentCommitted {
idempotency_key: identity.idempotency_key,
item_id: identity.item_id,
previous_item_id: identity.previous_item_id,
content_index: identity.content_index,
media_type: identity.media_type,
},
})
}
Some(RealtimeUserContentApplyOutcome::RejectedInvalidIdentity {
idempotency_key,
}) => {
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::ImageInputIdempotencyKeyInvalid {
max_bytes: meerkat_core::live_adapter::MAX_LIVE_IMAGE_IDEMPOTENCY_KEY_BYTES as u64,
actual_bytes: idempotency_key.len() as u64,
},
};
self.projection_sink
.signal_terminal_error(
&session_id,
code.clone(),
"invalid durable image identity",
)
.await?;
Ok(ObservationOutcome::Terminal { code })
}
Some(RealtimeUserContentApplyOutcome::RejectedUnmaterializedPredecessor {
..
}) => {
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::ImageInputRequiresCommit,
};
self.projection_sink
.signal_terminal_error(
&session_id,
code.clone(),
"image predecessor is not canonical",
)
.await?;
Ok(ObservationOutcome::Terminal { code })
}
Some(RealtimeUserContentApplyOutcome::RejectedConflict { .. }) => {
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::ImageInputIdempotencyConflict,
};
self.projection_sink
.signal_terminal_error(
&session_id,
code.clone(),
"image idempotency conflict after provider acknowledgement",
)
.await?;
Ok(ObservationOutcome::Terminal { code })
}
None => Ok(ObservationOutcome::TranscriptAppended),
}
}
(
ObservationRouting::DispatchToolCall { .. },
LiveAdapterObservation::ToolCallRequested {
provider_call_id,
tool_name,
arguments,
},
) => {
self.dispatch_tool_call(channel_id, provider_call_id, tool_name, arguments.clone())
.await
}
(
ObservationRouting::SignalInterrupt,
LiveAdapterObservation::TurnInterrupted { response_id },
) => {
self.projection_sink
.signal_turn_interrupt(&session_id, response_id.as_deref())
.await?;
Ok(ObservationOutcome::InterruptSignalled)
}
(
ObservationRouting::TerminalError,
LiveAdapterObservation::Error { code, message },
) => {
self.projection_sink
.signal_terminal_error(&session_id, code.clone(), message)
.await?;
Ok(ObservationOutcome::Terminal { code: code.clone() })
}
(
ObservationRouting::CommandRejection,
LiveAdapterObservation::CommandRejected { code, message },
) => Ok(ObservationOutcome::CommandRejected {
code: code.clone(),
message: message.clone(),
}),
(ObservationRouting::AppendTranscript, _) => Ok(ObservationOutcome::Noop),
_ => Ok(ObservationOutcome::Noop),
}
}
async fn dispatch_tool_call(
&self,
channel_id: &LiveChannelId,
provider_call_id: &ToolCallId,
tool_name: &ToolName,
arguments: serde_json::Value,
) -> Result<ObservationOutcome, LiveAdapterHostError> {
let dispatcher = match self.load_dispatcher() {
Some(d) => d,
None => {
let _ = self
.submit_tool_error(
channel_id,
provider_call_id.clone(),
"live tool dispatcher not configured".to_string(),
)
.await;
return Ok(ObservationOutcome::ToolCallSkipped {
provider_call_id: provider_call_id.0.clone(),
tool_name: tool_name.to_string(),
reason: ToolDispatchSkipReason::NoDispatcher,
});
}
};
let session_id = self.channel_session(channel_id).await?;
let call = ToolCall::new(provider_call_id.0.clone(), tool_name.to_string(), arguments);
let dispatch_call = dispatcher.dispatch_live_tool_call(&session_id, call);
let timeout = self.tool_timeout;
let dispatch_result = match tokio::time::timeout(timeout, dispatch_call).await {
Ok(result) => result,
Err(_elapsed) => {
self.submit_tool_dispatch_timeout(
channel_id,
provider_call_id.clone(),
LiveToolDispatchTimeout::new(timeout),
)
.await?;
return Ok(ObservationOutcome::ToolCallTimedOut {
provider_call_id: provider_call_id.0.clone(),
tool_name: tool_name.to_string(),
timeout,
});
}
};
match dispatch_result {
Ok(outcome) => {
let live_result =
tool_result_from_dispatch(provider_call_id.clone(), outcome.result);
self.submit_tool_result(channel_id, live_result).await?;
}
Err(err) => {
self.submit_tool_error(channel_id, provider_call_id.clone(), err.to_string())
.await?;
}
}
Ok(ObservationOutcome::ToolCallDispatched {
provider_call_id: provider_call_id.0.clone(),
tool_name: tool_name.to_string(),
})
}
async fn adapter_for(
&self,
channel_id: &LiveChannelId,
require_ready: bool,
) -> Result<Arc<dyn LiveAdapter>, LiveAdapterHostError> {
let inner = self.inner.lock().await;
let channel = inner
.channels
.get(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if channel.retire_at.is_some() {
return Err(LiveAdapterHostError::ChannelNotReady(
channel_id.clone(),
channel.status.clone(),
));
}
if require_ready && !channel.status.accepts_commands() {
return Err(LiveAdapterHostError::ChannelNotReady(
channel_id.clone(),
channel.status.clone(),
));
}
let adapter = channel
.adapter
.as_ref()
.ok_or_else(|| LiveAdapterHostError::NoAdapter(channel_id.clone()))?;
Ok(Arc::clone(adapter))
}
#[cfg(test)]
pub(crate) async fn signal_terminal_error(
&self,
channel_id: &LiveChannelId,
code: LiveAdapterErrorCode,
) -> Result<(), LiveAdapterHostError> {
let observation = self
.signal_terminal_error_observed(channel_id, code)
.await?;
self.prepare_channel_physical_close(&observation).await?;
let authority = self
.close_commit_authority_from_generated_test_machine(&observation)
.await?;
self.commit_channel_close_observation(&observation, &authority)
.await
}
pub async fn signal_terminal_error_observed(
&self,
channel_id: &LiveChannelId,
code: LiveAdapterErrorCode,
) -> Result<LiveChannelCloseObservation, LiveAdapterHostError> {
let message = match &code {
LiveAdapterErrorCode::ConfigRejected { reason } => reason.to_string(),
other => format!("{other:?}"),
};
let synthetic = LiveAdapterObservation::Error {
code: code.clone(),
message: message.clone(),
};
let adapter = {
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.pending_synthetic_obs = Some(synthetic.clone());
channel.adapter.clone()
};
if let Some(adapter) = adapter {
let _ = adapter.inject_observation(synthetic).await;
}
self.reserve_channel_close_observation(channel_id).await
}
pub async fn reserve_channel_close_observation(
&self,
channel_id: &LiveChannelId,
) -> Result<LiveChannelCloseObservation, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.close_observation_sequence = channel.close_observation_sequence.saturating_add(1);
LiveChannelCloseObservation::from_host_close_observation(
channel_id.as_str().to_owned(),
channel.close_observation_sequence,
)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn commit_channel_close_observation(
&self,
observation: &LiveChannelCloseObservation,
authority: &LiveChannelCloseCommitAuthority,
) -> Result<(), LiveAdapterHostError> {
if authority.channel_id() != observation.channel_id()
|| authority.close_sequence() != observation.close_sequence()
{
return Err(LiveAdapterHostError::CloseNotAuthorized);
}
let channel_id = LiveChannelId::new(observation.channel_id().to_owned());
{
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get_mut(&channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if !channel.physical_close_confirmed || channel.adapter.is_some() {
return Err(LiveAdapterHostError::CloseNotAuthorized);
}
authority.consume_once()?;
channel.status = LiveAdapterStatus::Closed;
channel.retire_at = Some(std::time::Instant::now() + CLOSED_CHANNEL_TTL);
}
Ok(())
}
pub async fn prepare_channel_physical_close(
&self,
observation: &LiveChannelCloseObservation,
) -> Result<(), LiveAdapterHostError> {
let channel_id = LiveChannelId::new(observation.channel_id().to_owned());
let adapter = {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get(&channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if channel.physical_close_confirmed && channel.adapter.is_none() {
return Ok(());
}
channel
.adapter
.as_ref()
.map(Arc::clone)
.ok_or_else(|| LiveAdapterHostError::NoAdapter(channel_id.clone()))?
};
adapter.close().await?;
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(&channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
match channel.adapter.as_ref() {
Some(current) if Arc::ptr_eq(current, &adapter) => {
channel.adapter = None;
channel.physical_close_confirmed = true;
Ok(())
}
None if channel.physical_close_confirmed => Ok(()),
_ => Err(LiveAdapterHostError::CloseNotAuthorized),
}
}
#[cfg(test)]
pub(crate) async fn close_commit_authority_from_generated_test_machine(
&self,
observation: &LiveChannelCloseObservation,
) -> Result<LiveChannelCloseCommitAuthority, LiveAdapterHostError> {
let channel_id = LiveChannelId::new(observation.channel_id().to_owned());
let session_id = self.channel_session(&channel_id).await?;
Ok(
LiveChannelCloseCommitAuthority::from_generated_test_machine(
&session_id,
&channel_id,
observation.close_sequence(),
),
)
}
#[cfg(test)]
pub(crate) async fn close_channel_observed_with_generated_test_machine_authority(
&self,
channel_id: &LiveChannelId,
) -> Result<LiveChannelCloseObservation, LiveAdapterHostError> {
let observation = self.reserve_channel_close_observation(channel_id).await?;
self.prepare_channel_physical_close(&observation).await?;
let authority = self
.close_commit_authority_from_generated_test_machine(&observation)
.await?;
self.commit_channel_close_observation(&observation, &authority)
.await?;
Ok(observation)
}
#[cfg(test)]
pub(crate) async fn close_channel_with_generated_test_machine_authority(
&self,
channel_id: &LiveChannelId,
) -> Result<(), LiveAdapterHostError> {
self.close_channel_observed_with_generated_test_machine_authority(channel_id)
.await
.map(|_| ())
}
#[cfg(test)]
pub(crate) async fn status_commit_authority_from_generated_test_machine(
&self,
observation: &LiveChannelStatusObservation,
) -> Result<LiveChannelStatusCommitAuthority, LiveAdapterHostError> {
let channel_id = LiveChannelId::new(observation.channel_id().to_owned());
let session_id = self.channel_session(&channel_id).await?;
let channel_id_string = channel_id.as_str().to_owned();
let (status, degradation_reason, degradation_detail) =
generated_test_live_channel_status(observation.status());
let mut authority = generated_test_machine_for_registered_session(&session_id);
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineMutator::apply(
&mut authority,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineInput::ResolveLiveOpenAdmission {
session_id: session_id.to_string(),
channel_id: channel_id_string.clone(),
llm_identity: generated_test_llm_identity(),
},
)
.expect("generated MeerkatMachine live-open admission");
let transition = meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineMutator::apply(
&mut authority,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineInput::RecordLiveChannelStatus {
channel_id: channel_id_string.clone(),
status,
status_observation_sequence: observation.observation_sequence(),
degradation_reason,
degradation_detail: degradation_detail.clone(),
},
)
.expect("generated MeerkatMachine live-status result");
assert!(
transition.effects().iter().any(|effect| matches!(
effect,
meerkat_machine_schema::catalog::dsl::meerkat_machine::MeerkatMachineEffect::LiveChannelStatusResolved {
channel_id: effect_channel_id,
status: effect_status,
status_observation_sequence,
..
} if *effect_channel_id == channel_id_string
&& *effect_status == status
&& *status_observation_sequence == observation.observation_sequence()
)),
"generated live-status result effect"
);
Ok(LiveChannelStatusCommitAuthority::from_generated_parts(
channel_id_string,
observation.observation_sequence(),
))
}
#[cfg(test)]
pub(crate) async fn commit_status_with_generated_test_machine_authority(
&self,
channel_id: &LiveChannelId,
status: LiveAdapterStatus,
) -> Result<LiveChannelStatusObservation, LiveAdapterHostError> {
let observation = self
.reserve_channel_status_observation(channel_id, status)
.await?;
let authority = self
.status_commit_authority_from_generated_test_machine(&observation)
.await?;
self.commit_channel_status_observation(&observation, &authority)
.await?;
Ok(observation)
}
pub(crate) async fn channel_status(
&self,
channel_id: &LiveChannelId,
) -> Result<LiveAdapterStatus, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
inner
.channels
.get(channel_id)
.map(|ch| ch.status.clone())
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn channel_status_observation(
&self,
channel_id: &LiveChannelId,
) -> Result<LiveChannelStatusObservation, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.status_observation_sequence = channel.status_observation_sequence.saturating_add(1);
LiveChannelStatusObservation::from_host_status_observation(
channel_id.as_str().to_owned(),
channel.status.clone(),
channel.status_observation_sequence,
)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn reserve_channel_status_observation(
&self,
channel_id: &LiveChannelId,
status: LiveAdapterStatus,
) -> Result<LiveChannelStatusObservation, LiveAdapterHostError> {
if status.is_terminal() {
return Err(LiveAdapterHostError::StatusNotAuthorized);
}
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if channel.retire_at.is_some() {
return Err(LiveAdapterHostError::ChannelNotReady(
channel_id.clone(),
channel.status.clone(),
));
}
channel.status_observation_sequence = channel.status_observation_sequence.saturating_add(1);
LiveChannelStatusObservation::from_host_status_observation(
channel_id.as_str().to_owned(),
status,
channel.status_observation_sequence,
)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn commit_channel_status_observation(
&self,
observation: &LiveChannelStatusObservation,
authority: &LiveChannelStatusCommitAuthority,
) -> Result<(), LiveAdapterHostError> {
authority.consume_once()?;
if authority.channel_id() != observation.channel_id()
|| authority.status_observation_sequence() != observation.observation_sequence()
{
return Err(LiveAdapterHostError::StatusNotAuthorized);
}
if observation.status().is_terminal() {
return Err(LiveAdapterHostError::StatusNotAuthorized);
}
let channel_id = LiveChannelId::new(observation.channel_id().to_owned());
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
let channel = inner
.channels
.get_mut(&channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
if channel.retire_at.is_some() {
return Err(LiveAdapterHostError::ChannelNotReady(
channel_id,
channel.status.clone(),
));
}
channel.status = observation.status().clone();
Ok(())
}
pub async fn channel_session(
&self,
channel_id: &LiveChannelId,
) -> Result<SessionId, LiveAdapterHostError> {
let inner = self.inner.lock().await;
inner
.channels
.get(channel_id)
.map(|ch| ch.session_id.clone())
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))
}
pub async fn signal_transport_barge_in(
&self,
channel_id: &LiveChannelId,
) -> Result<(), LiveAdapterHostError> {
let session_id = self.channel_session(channel_id).await?;
self.projection_sink
.signal_turn_interrupt(&session_id, None)
.await?;
Ok(())
}
pub async fn signal_output_audio_degraded(
&self,
channel_id: &LiveChannelId,
dropped: u64,
) -> Result<(), LiveAdapterHostError> {
let session_id = self.channel_session(channel_id).await?;
self.projection_sink
.signal_output_audio_degraded(&session_id, dropped)
.await?;
Ok(())
}
pub fn classify_observation(observation: &LiveAdapterObservation) -> ObservationRouting {
match observation {
LiveAdapterObservation::Ready => {
ObservationRouting::UpdateStatus(LiveAdapterStatus::Ready)
}
LiveAdapterObservation::UserTranscriptFinal { .. } => {
ObservationRouting::AppendTranscript
}
LiveAdapterObservation::AssistantTextDelta { .. } => {
ObservationRouting::AppendTranscript
}
LiveAdapterObservation::AssistantTranscriptDelta { .. } => {
ObservationRouting::AppendTranscript
}
LiveAdapterObservation::AssistantAudioChunk { .. } => ObservationRouting::Noop,
LiveAdapterObservation::AssistantTranscriptFinal { .. } => {
ObservationRouting::AppendTranscript
}
LiveAdapterObservation::AssistantTranscriptTruncated { .. } => {
ObservationRouting::AppendTranscript
}
LiveAdapterObservation::RealtimeTranscript { .. } => {
ObservationRouting::AppendRealtimeTranscript
}
LiveAdapterObservation::UserContentCommitted { .. } => ObservationRouting::Noop,
LiveAdapterObservation::ToolCallRequested {
provider_call_id,
tool_name,
..
} => ObservationRouting::DispatchToolCall {
provider_call_id: provider_call_id.0.clone(),
tool_name: tool_name.to_string(),
},
LiveAdapterObservation::TurnInterrupted { .. } => ObservationRouting::SignalInterrupt,
LiveAdapterObservation::TurnCompleted { .. } => ObservationRouting::AppendTranscript,
LiveAdapterObservation::StatusChanged { status } => {
ObservationRouting::UpdateStatus(status.clone())
}
LiveAdapterObservation::Error { .. } => ObservationRouting::TerminalError,
LiveAdapterObservation::CommandRejected { .. } => ObservationRouting::CommandRejection,
_ => ObservationRouting::Noop,
}
}
fn observation_requires_generated_close(observation: &LiveAdapterObservation) -> bool {
match observation {
LiveAdapterObservation::Error { .. } => true,
LiveAdapterObservation::StatusChanged { status } => status.is_terminal(),
_ => false,
}
}
pub(crate) async fn generated_close_has_committed(
&self,
channel_id: &LiveChannelId,
) -> Result<bool, LiveAdapterHostError> {
self.channel_status(channel_id)
.await
.map(|status| status.is_terminal())
}
pub async fn next_snapshot_version(
&self,
channel_id: &LiveChannelId,
) -> Result<u64, LiveAdapterHostError> {
let mut inner = self.inner.lock().await;
let channel = inner
.channels
.get_mut(channel_id)
.ok_or_else(|| LiveAdapterHostError::ChannelNotFound(channel_id.clone()))?;
channel.snapshot_version += 1;
Ok(channel.snapshot_version)
}
pub async fn active_channels(&self) -> Vec<LiveChannelId> {
let mut inner = self.inner.lock().await;
Self::reap_retired_locked(&mut inner);
inner
.channels
.iter()
.filter(|(_, ch)| ch.retire_at.is_none())
.map(|(id, _)| id.clone())
.collect()
}
fn reap_retired_locked(inner: &mut HostInner) {
let now = std::time::Instant::now();
let to_drop: Vec<LiveChannelId> = inner
.channels
.iter()
.filter_map(|(id, ch)| match ch.retire_at {
Some(deadline) if deadline <= now => Some(id.clone()),
_ => None,
})
.collect();
for id in to_drop {
if let Some(ch) = inner.channels.shift_remove(&id) {
if inner
.by_session
.get(&ch.session_id)
.is_some_and(|current| current == &id)
{
inner.by_session.remove(&ch.session_id);
}
}
}
}
}
fn tool_result_from_dispatch(call_id: ToolCallId, result: ToolResult) -> LiveToolResult {
LiveToolResult {
call_id,
content: result.content,
is_error: result.is_error,
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use async_trait::async_trait;
use meerkat_core::live_adapter::{
LiveAdapterError, LiveAdapterErrorCode, LiveAdapterObservation, LiveDegradationReason,
};
use meerkat_core::ops::ToolDispatchOutcome;
use meerkat_core::types::{StopReason, Usage};
use std::sync::Mutex as StdMutex;
fn test_session_id() -> SessionId {
SessionId::new()
}
#[test]
fn tool_timeout_defaults_to_canonical_value_and_builder_overrides_it() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
assert_eq!(host.tool_timeout(), DEFAULT_LIVE_TOOL_TIMEOUT);
let override_timeout = Duration::from_secs(7);
let host =
LiveAdapterHost::new(Arc::new(NoOpProjectionSink)).with_tool_timeout(override_timeout);
assert_eq!(host.tool_timeout(), override_timeout);
assert_eq!(DEFAULT_LIVE_TOOL_TIMEOUT, Duration::from_secs(30));
}
#[tokio::test]
async fn projection_sink_is_mandatory_at_construction() {
let recording = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&recording) as _);
let session = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::UserTranscriptFinal {
provider_item_id: Some("item-1".into()),
previous_item_id: None,
content_index: Some(0),
text: "hello".into(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(matches!(outcome, ObservationOutcome::TranscriptAppended));
assert_eq!(
recording.user_transcripts.lock().unwrap().len(),
1,
"production-shape host must route user transcripts to the sink"
);
let noop_host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session2 = test_session_id();
let ch2 = noop_host
.open_channel_with_generated_test_machine_authority(session2)
.await
.unwrap();
let outcome2 = noop_host.apply_observation(&ch2, &obs).await.unwrap();
assert!(matches!(outcome2, ObservationOutcome::TranscriptAppended));
}
#[tokio::test]
async fn open_channel_returns_unique_ids() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let s1 = test_session_id();
let s2 = test_session_id();
let ch1 = host
.open_channel_with_generated_test_machine_authority(s1)
.await
.unwrap();
let ch2 = host
.open_channel_with_generated_test_machine_authority(s2)
.await
.unwrap();
assert_ne!(ch1, ch2);
}
#[tokio::test]
async fn open_channel_ids_are_uuid_shape_not_live_n() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch1 = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let ch2 = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
for ch in [&ch1, &ch2] {
let s = ch.as_str();
assert!(
!s.starts_with("live_"),
"channel id retained legacy `live_N` shape: {s}"
);
let parsed =
uuid::Uuid::parse_str(s).expect("channel id should be a valid UUID string");
assert_eq!(
parsed.get_version(),
Some(uuid::Version::Random),
"channel id should be a v4 UUID"
);
}
assert_ne!(ch1, ch2);
}
#[tokio::test]
async fn open_channel_starts_in_opening_status() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Opening);
}
#[tokio::test]
async fn channel_status_observation_advances_per_channel() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let first = host.channel_status_observation(&ch).await.unwrap();
let second = host.channel_status_observation(&ch).await.unwrap();
assert_eq!(first.channel_id(), ch.as_str());
assert_eq!(first.status(), &LiveAdapterStatus::Opening);
assert_eq!(first.observation_sequence(), 1);
assert_eq!(second.observation_sequence(), 2);
}
#[tokio::test]
async fn duplicate_session_binding_rejected() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session_id = test_session_id();
let _ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let err = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap_err();
assert!(matches!(err, LiveAdapterHostError::SessionAlreadyBound(id) if id == session_id));
}
#[tokio::test]
async fn close_channel_marks_closed_and_retains_for_status_reads() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.close_channel_with_generated_test_machine_authority(&ch)
.await
.unwrap();
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Closed);
}
#[tokio::test]
async fn physical_close_failure_retains_discovery_and_blocks_terminal_commit() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(FailOnceCloseAdapter::default()))
.await
.unwrap();
let observation = host.reserve_channel_close_observation(&ch).await.unwrap();
assert!(matches!(
host.prepare_channel_physical_close(&observation).await,
Err(LiveAdapterHostError::AdapterError(_))
));
assert!(
host.active_channels().await.contains(&ch),
"failed physical close must retain host discovery for retry"
);
let premature_authority = host
.close_commit_authority_from_generated_test_machine(&observation)
.await
.unwrap();
assert!(matches!(
host.commit_channel_close_observation(&observation, &premature_authority)
.await,
Err(LiveAdapterHostError::CloseNotAuthorized)
));
host.prepare_channel_physical_close(&observation)
.await
.expect("retry physically closes the exact retained adapter");
let authority = host
.close_commit_authority_from_generated_test_machine(&observation)
.await
.unwrap();
host.commit_channel_close_observation(&observation, &authority)
.await
.unwrap();
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Closed
);
}
#[tokio::test]
async fn close_channel_observation_advances_per_channel() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let first = host
.close_channel_observed_with_generated_test_machine_authority(&ch)
.await
.unwrap();
let second = host
.close_channel_observed_with_generated_test_machine_authority(&ch)
.await
.unwrap();
assert_eq!(first.channel_id(), ch.as_str());
assert_eq!(first.close_sequence(), 1);
assert_eq!(second.close_sequence(), 2);
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Closed
);
}
#[tokio::test]
async fn close_channel_allows_rebinding_same_session() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
host.close_channel_with_generated_test_machine_authority(&ch)
.await
.unwrap();
let ch2 = host
.open_channel_with_generated_test_machine_authority(session_id)
.await
.unwrap();
assert_ne!(ch, ch2);
}
#[tokio::test]
async fn reap_of_retired_channel_preserves_rebound_session_mapping() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session_id = test_session_id();
let ch_a = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
host.close_channel_with_generated_test_machine_authority(&ch_a)
.await
.unwrap();
let ch_b = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
assert_ne!(ch_a, ch_b);
{
let mut inner = host.inner.lock().await;
if let Some(channel) = inner.channels.get_mut(&ch_a) {
channel.retire_at =
Some(std::time::Instant::now() - std::time::Duration::from_secs(1));
}
}
let active = host.active_channels().await;
assert_eq!(active, vec![ch_b.clone()]);
{
let inner = host.inner.lock().await;
assert_eq!(
inner.by_session.get(&session_id),
Some(&ch_b),
"reap of retired A must not clear B's reverse mapping"
);
assert!(
!inner.channels.contains_key(&ch_a),
"retired channel A must be dropped"
);
assert!(
inner.channels.contains_key(&ch_b),
"rebound channel B must remain"
);
}
let err = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap_err();
assert!(
matches!(err, LiveAdapterHostError::SessionAlreadyBound(id) if id == session_id),
"after reap, third open for session must still see B as bound"
);
assert_eq!(host.active_channels().await.len(), 1);
}
#[tokio::test]
async fn active_channels_excludes_retained_closed_channels() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let s1 = test_session_id();
let s2 = test_session_id();
let live = host
.open_channel_with_generated_test_machine_authority(s1)
.await
.unwrap();
let closing = host
.open_channel_with_generated_test_machine_authority(s2)
.await
.unwrap();
let active_pre = host.active_channels().await;
assert!(active_pre.contains(&live));
assert!(active_pre.contains(&closing));
assert_eq!(active_pre.len(), 2);
host.close_channel_with_generated_test_machine_authority(&closing)
.await
.unwrap();
let active_during_ttl = host.active_channels().await;
assert_eq!(
active_during_ttl,
vec![live.clone()],
"retained-closed channel must not appear in active_channels()"
);
assert_eq!(
host.channel_status(&closing).await.unwrap(),
LiveAdapterStatus::Closed,
);
{
let mut inner = host.inner.lock().await;
if let Some(channel) = inner.channels.get_mut(&closing) {
channel.retire_at =
Some(std::time::Instant::now() - std::time::Duration::from_secs(1));
}
}
let active_post_reap = host.active_channels().await;
assert_eq!(active_post_reap, vec![live.clone()]);
{
let inner = host.inner.lock().await;
assert!(
!inner.channels.contains_key(&closing),
"post-reap, retired channel must be dropped from the map"
);
}
}
#[tokio::test]
async fn channel_session_returns_bound_session() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
assert_eq!(host.channel_session(&ch).await.unwrap(), session_id);
}
#[tokio::test]
async fn signal_terminal_error_enqueues_synthetic_error_obs_and_closes_channel() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::RefreshModelSwap {
from_model: "gpt-realtime".to_string(),
to_model: "gpt-realtime-1.5".to_string(),
},
};
host.signal_terminal_error(&ch, code).await.unwrap();
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Closed);
let obs = host
.next_observation_raw(&ch)
.await
.expect("next_observation_raw should return synthetic obs even post-close")
.expect("synthetic obs must be Some");
match obs {
LiveAdapterObservation::Error { code, message } => match code {
LiveAdapterErrorCode::ConfigRejected { reason } => {
assert!(matches!(
reason,
meerkat_core::live_adapter::LiveConfigRejectionReason::RefreshModelSwap {
ref to_model,
..
} if to_model == "gpt-realtime-1.5"
));
assert!(message.contains("close + reopen"));
}
other => panic!("expected ConfigRejected, got {other:?}"),
},
other => panic!("expected Error observation, got {other:?}"),
}
}
#[tokio::test]
async fn synthetic_terminal_error_routes_through_apply_observation_to_terminal_outcome() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::ChannelIdentitySwap {
from_model: "a".to_string(),
from_provider: meerkat_core::Provider::OpenAI,
to_model: "b".to_string(),
to_provider: meerkat_core::Provider::OpenAI,
auth_binding_changed: false,
},
};
host.signal_terminal_error(&ch, code).await.unwrap();
let obs = host.next_observation_raw(&ch).await.unwrap().unwrap();
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::Terminal { code } => match code {
LiveAdapterErrorCode::ConfigRejected { reason } => {
assert!(matches!(
reason,
meerkat_core::live_adapter::LiveConfigRejectionReason::ChannelIdentitySwap {
ref from_model, ref to_model, ..
} if from_model == "a" && to_model == "b"
));
}
other => panic!("expected ConfigRejected, got {other:?}"),
},
other => panic!("expected Terminal outcome, got {other:?}"),
}
}
#[tokio::test]
async fn signal_terminal_error_delivers_synthetic_error_before_close_signal() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
let code = LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::Other {
detail: "model_swap_test".to_string(),
},
};
host.signal_terminal_error(&ch, code).await.unwrap();
let first = host
.next_observation_raw(&ch)
.await
.expect("first read should succeed")
.expect("synthetic Error must surface before end-of-stream");
match first {
LiveAdapterObservation::Error { code, message } => match code {
LiveAdapterErrorCode::ConfigRejected { reason } => {
assert!(matches!(
&reason,
meerkat_core::live_adapter::LiveConfigRejectionReason::Other { detail }
if detail == "model_swap_test"
));
assert_eq!(message, "model_swap_test");
}
other => unreachable!("expected ConfigRejected, got {other:?}"),
},
other => unreachable!("expected Error obs first, got {other:?}"),
}
}
#[tokio::test]
async fn signal_terminal_error_on_missing_channel_returns_channel_not_found() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let bogus = LiveChannelId::random_uuid();
let result = host
.signal_terminal_error(
&bogus,
LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::Other {
detail: "no channel".into(),
},
},
)
.await;
assert!(
matches!(&result, Err(LiveAdapterHostError::ChannelNotFound(id)) if id == &bogus),
"expected ChannelNotFound for unknown channel, got {result:?}"
);
}
#[test]
fn ready_observation_routes_to_status_update() {
let routing = LiveAdapterHost::classify_observation(&LiveAdapterObservation::Ready);
assert_eq!(
routing,
ObservationRouting::UpdateStatus(LiveAdapterStatus::Ready)
);
}
#[test]
fn tool_call_observation_routes_to_dispatch() {
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_1"),
tool_name: ToolName::new("calculator"),
arguments: serde_json::json!({"x": 1}),
};
let routing = LiveAdapterHost::classify_observation(&obs);
assert_eq!(
routing,
ObservationRouting::DispatchToolCall {
provider_call_id: "call_1".into(),
tool_name: "calculator".into(),
}
);
}
#[test]
fn barge_in_observation_routes_to_interrupt() {
let routing =
LiveAdapterHost::classify_observation(&LiveAdapterObservation::TurnInterrupted {
response_id: None,
});
assert_eq!(routing, ObservationRouting::SignalInterrupt);
let routing_with_id =
LiveAdapterHost::classify_observation(&LiveAdapterObservation::TurnInterrupted {
response_id: Some("resp_42".into()),
});
assert_eq!(routing_with_id, ObservationRouting::SignalInterrupt);
}
#[test]
fn user_transcript_routes_to_append() {
let obs = LiveAdapterObservation::UserTranscriptFinal {
provider_item_id: Some("item_1".into()),
previous_item_id: None,
content_index: None,
text: "hello".into(),
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::AppendTranscript
);
}
#[test]
fn assistant_text_delta_routes_to_append() {
let obs = LiveAdapterObservation::AssistantTextDelta {
provider_item_id: Some("item_2".into()),
previous_item_id: None,
content_index: None,
response_id: None,
delta_id: None,
delta: "world".into(),
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::AppendTranscript
);
}
#[test]
fn turn_completed_routes_to_append() {
let obs = LiveAdapterObservation::TurnCompleted {
response_id: None,
stop_reason: StopReason::EndTurn,
usage: Usage {
input_tokens: 10,
output_tokens: 5,
cache_creation_tokens: None,
cache_read_tokens: None,
},
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::AppendTranscript
);
}
#[test]
fn error_observation_routes_to_terminal() {
let obs = LiveAdapterObservation::Error {
code: LiveAdapterErrorCode::ConnectionLost,
message: "ws closed".into(),
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::TerminalError
);
}
#[test]
fn audio_chunk_routes_to_noop() {
let obs = LiveAdapterObservation::AssistantAudioChunk {
data: vec![0; 100],
sample_rate_hz: 24000,
channels: 1,
response_id: None,
item_id: None,
content_index: None,
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::Noop
);
}
#[test]
fn user_content_committed_receipt_routes_to_noop() {
let obs = LiveAdapterObservation::UserContentCommitted {
idempotency_key: "image-request-1".into(),
item_id: "item_image".into(),
previous_item_id: None,
content_index: 0,
media_type: "image/png".into(),
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::Noop
);
}
#[test]
fn status_changed_routes_to_status_update() {
let obs = LiveAdapterObservation::StatusChanged {
status: LiveAdapterStatus::Degraded {
reason: LiveDegradationReason::ProviderThrottled,
},
};
assert_eq!(
LiveAdapterHost::classify_observation(&obs),
ObservationRouting::UpdateStatus(LiveAdapterStatus::Degraded {
reason: LiveDegradationReason::ProviderThrottled,
})
);
}
#[tokio::test]
async fn commit_status_with_generated_authority_changes_channel_status() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Ready
);
}
#[tokio::test]
async fn terminal_status_update_requires_generated_close_authority() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let err = host
.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Closed)
.await
.expect_err("terminal status update must not bypass generated close authority");
assert!(matches!(err, LiveAdapterHostError::StatusNotAuthorized));
let obs = LiveAdapterObservation::StatusChanged {
status: LiveAdapterStatus::Closed,
};
let err = host
.apply_observation(&ch, &obs)
.await
.expect_err("closed observation must not bypass generated close authority");
assert!(matches!(err, LiveAdapterHostError::CloseNotAuthorized));
host.close_channel_with_generated_test_machine_authority(&ch)
.await
.unwrap();
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert_eq!(
outcome,
ObservationOutcome::StatusUpdated(LiveAdapterStatus::Closed)
);
}
#[tokio::test]
async fn snapshot_version_increments_monotonically() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let v1 = host.next_snapshot_version(&ch).await.unwrap();
let v2 = host.next_snapshot_version(&ch).await.unwrap();
assert_eq!(v1, 1);
assert_eq!(v2, 2);
}
#[tokio::test]
async fn active_channels_lists_open_channels() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch1 = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let ch2 = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let active = host.active_channels().await;
assert_eq!(active.len(), 2);
assert!(active.contains(&ch1));
assert!(active.contains(&ch2));
}
#[tokio::test]
async fn send_input_without_adapter_returns_error() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let err = host
.send_input(&ch, LiveInputChunk::Text { text: "hi".into() })
.await
.unwrap_err();
assert!(matches!(err, LiveAdapterHostError::ChannelNotReady(_, _)));
}
#[tokio::test]
async fn send_command_rejects_refresh_without_typed_acceptance_path() {
let session_id = test_session_id();
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let snapshot = meerkat_core::live_adapter::LiveProjectionSnapshot {
session_id,
snapshot_version: 1,
seed_messages: vec![],
visible_tools: vec![],
system_prompt: None,
model_id: "model-a".into(),
provider_id: meerkat_core::Provider::Other,
audio_config: None,
runtime_system_context: vec![],
user_content_identities: vec![],
user_content_tombstones: vec![],
canonical_user_image_decoded_bytes: None,
transcript_rewrite_generation: 0,
};
let err = host
.send_command(&ch, LiveAdapterCommand::Refresh { snapshot })
.await
.unwrap_err();
assert!(matches!(err, LiveAdapterHostError::UnsupportedCommand(_)));
}
#[tokio::test]
async fn attach_adapter_does_not_assert_ready() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Opening
);
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Opening,
"attach_adapter must NOT mark channel Ready (F32)"
);
}
#[tokio::test]
async fn send_input_rejected_when_not_ready() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
let err = host
.send_input(&ch, LiveInputChunk::Text { text: "hi".into() })
.await
.unwrap_err();
match err {
LiveAdapterHostError::ChannelNotReady(_, status) => {
assert_eq!(status, LiveAdapterStatus::Opening);
}
other => panic!("expected ChannelNotReady, got {other:?}"),
}
}
#[tokio::test]
async fn send_input_accepts_when_ready() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
host.send_input(&ch, LiveInputChunk::Text { text: "hi".into() })
.await
.unwrap();
}
#[tokio::test]
async fn adapter_pump_error_routes_close_through_generated_authority() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(ErroringAdapter))
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = host
.next_observation_raw(&ch)
.await
.unwrap()
.expect("synthetic Error obs surfaces on adapter Err");
match &obs {
LiveAdapterObservation::Error { code, message } => {
assert_eq!(*code, LiveAdapterErrorCode::ProviderError);
assert!(
message.contains("adapter read failure"),
"synthetic message must explain origin; got `{message}`"
);
}
other => unreachable!("expected synthetic Error, got {other:?}"),
}
let err = host
.apply_observation(&ch, &obs)
.await
.expect_err("terminal observation must wait for generated close authority");
assert!(matches!(err, LiveAdapterHostError::CloseNotAuthorized));
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Ready);
{
let inner = host.inner.lock().await;
let channel = inner
.channels
.get(&ch)
.expect("channel remains present before close authority");
assert!(
channel.retire_at.is_none(),
"adapter Err must not set retire_at before generated close authority"
);
assert!(
channel.adapter.is_some(),
"adapter Err must not drop the adapter before generated close authority"
);
}
host.close_channel_with_generated_test_machine_authority(&ch)
.await
.unwrap();
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert_eq!(
outcome,
ObservationOutcome::Terminal {
code: LiveAdapterErrorCode::ProviderError
}
);
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Closed);
{
let inner = host.inner.lock().await;
let channel = inner
.channels
.get(&ch)
.expect("channel preserved for live/status until TTL elapses");
assert!(
channel.retire_at.is_some(),
"generated close authority must set retire_at"
);
assert!(
channel.adapter.is_none(),
"generated close authority must drop the adapter Arc"
);
}
}
#[tokio::test]
async fn command_rejected_routes_non_terminally_and_preserves_channel() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::CommandRejected {
code: LiveAdapterErrorCode::ConfigRejected {
reason:
meerkat_core::live_adapter::LiveConfigRejectionReason::ImageInputNotImplemented,
},
message: "image_input_not_implemented".into(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::CommandRejected { code, message } => {
assert!(matches!(
code,
LiveAdapterErrorCode::ConfigRejected {
reason: meerkat_core::live_adapter::LiveConfigRejectionReason::ImageInputNotImplemented,
}
));
assert_eq!(message, "image_input_not_implemented");
}
other => {
unreachable!("CommandRejected must produce CommandRejected outcome, got {other:?}")
}
}
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Ready);
{
let inner = host.inner.lock().await;
let channel = inner.channels.get(&ch).expect("channel present");
assert!(
channel.retire_at.is_none(),
"CommandRejected must not retire the channel"
);
assert!(
channel.adapter.is_some(),
"CommandRejected must not drop the adapter"
);
}
}
#[tokio::test]
async fn adapter_err_releases_session_for_rebind() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let session_id = test_session_id();
let ch1 = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
host.attach_adapter(&ch1, Arc::new(ErroringAdapter))
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch1, LiveAdapterStatus::Ready)
.await
.unwrap();
let _ = host.next_observation_raw(&ch1).await.unwrap();
{
let mut inner = host.inner.lock().await;
if let Some(channel) = inner.channels.get_mut(&ch1) {
channel.retire_at =
Some(std::time::Instant::now() - std::time::Duration::from_secs(1));
}
}
let ch2 = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.expect("rebind for same session must succeed once previous channel is retired");
assert_ne!(ch1, ch2);
}
#[tokio::test]
async fn user_transcript_observation_appends_to_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::UserTranscriptFinal {
provider_item_id: Some("item_1".into()),
previous_item_id: None,
content_index: None,
text: "hello world".into(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(matches!(outcome, ObservationOutcome::TranscriptAppended));
let user = sink.user_transcripts.lock().unwrap();
assert_eq!(user.len(), 1);
assert_eq!(user[0].0, session_id);
assert_eq!(user[0].1, "hello world");
assert_eq!(user[0].2.provider_item_id.as_deref(), Some("item_1"));
}
#[tokio::test]
async fn user_transcript_full_identity_propagates_end_to_end() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::UserTranscriptFinal {
provider_item_id: Some("item_1".into()),
previous_item_id: Some("item_0".into()),
content_index: Some(2),
text: "hello".into(),
};
host.apply_observation(&ch, &obs).await.unwrap();
let user = sink.user_transcripts.lock().unwrap();
let identity = &user[0].2;
assert_eq!(identity.provider_item_id.as_deref(), Some("item_1"));
assert_eq!(identity.previous_item_id.as_deref(), Some("item_0"));
assert_eq!(identity.content_index, Some(2));
assert_eq!(identity.response_id, None);
assert_eq!(identity.delta_id, None);
}
#[tokio::test]
async fn assistant_text_delta_observation_appends_to_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::AssistantTextDelta {
provider_item_id: None,
previous_item_id: None,
content_index: None,
response_id: None,
delta_id: None,
delta: "Hello".into(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(matches!(outcome, ObservationOutcome::TranscriptAppended));
let deltas = sink.text_deltas.lock().unwrap();
assert_eq!(deltas.len(), 1);
assert_eq!(deltas[0].0, session_id);
assert_eq!(deltas[0].1, "Hello");
assert!(sink.transcript_deltas.lock().unwrap().is_empty());
}
#[tokio::test]
async fn assistant_text_delta_full_identity_propagates_end_to_end() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let obs = LiveAdapterObservation::AssistantTextDelta {
provider_item_id: Some("item_42".into()),
previous_item_id: Some("item_41".into()),
content_index: Some(1),
response_id: Some("resp_xyz".into()),
delta_id: Some("d_7".into()),
delta: "world".into(),
};
host.apply_observation(&ch, &obs).await.unwrap();
let deltas = sink.text_deltas.lock().unwrap();
let identity = &deltas[0].2;
assert_eq!(identity.provider_item_id.as_deref(), Some("item_42"));
assert_eq!(identity.previous_item_id.as_deref(), Some("item_41"));
assert_eq!(identity.content_index, Some(1));
assert_eq!(identity.response_id.as_deref(), Some("resp_xyz"));
assert_eq!(identity.delta_id.as_deref(), Some("d_7"));
}
#[tokio::test]
async fn assistant_transcript_final_observation_calls_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::AssistantTranscriptFinal {
provider_item_id: "resp_1".into(),
previous_item_id: None,
content_index: None,
response_id: None,
text: "All done.".into(),
stop_reason: StopReason::EndTurn,
usage: Usage::default(),
};
host.apply_observation(&ch, &obs).await.unwrap();
let finals = sink.transcript_finals.lock().unwrap();
assert_eq!(finals.len(), 1);
assert_eq!(finals[0].0, session_id);
assert_eq!(finals[0].1, "All done.");
assert!(sink.text_finals.lock().unwrap().is_empty());
}
#[tokio::test]
async fn assistant_transcript_final_full_identity_propagates_end_to_end() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let obs = LiveAdapterObservation::AssistantTranscriptFinal {
provider_item_id: "item_final".into(),
previous_item_id: Some("item_prev".into()),
content_index: Some(0),
response_id: Some("resp_final".into()),
text: "done".into(),
stop_reason: StopReason::EndTurn,
usage: Usage::default(),
};
host.apply_observation(&ch, &obs).await.unwrap();
let finals = sink.transcript_finals.lock().unwrap();
let identity = &finals[0].2;
assert_eq!(identity.provider_item_id.as_deref(), Some("item_final"));
assert_eq!(identity.previous_item_id.as_deref(), Some("item_prev"));
assert_eq!(identity.content_index, Some(0));
assert_eq!(identity.response_id.as_deref(), Some("resp_final"));
assert_eq!(identity.delta_id, None);
}
#[tokio::test]
async fn turn_completed_observation_signals_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let obs = LiveAdapterObservation::TurnCompleted {
response_id: None,
stop_reason: StopReason::EndTurn,
usage: Usage::default(),
};
host.apply_observation(&ch, &obs).await.unwrap();
let turns = sink.turn_completed.lock().unwrap();
assert_eq!(turns.len(), 1);
}
#[tokio::test]
async fn assistant_transcript_delta_routes_to_transcript_lane() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::AssistantTranscriptDelta {
provider_item_id: Some("item_t".into()),
previous_item_id: None,
content_index: Some(0),
response_id: Some("resp_t".into()),
delta_id: Some("d_t".into()),
delta: "spoken word".into(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(matches!(outcome, ObservationOutcome::TranscriptAppended));
let transcript_deltas = sink.transcript_deltas.lock().unwrap();
assert_eq!(transcript_deltas.len(), 1);
assert_eq!(transcript_deltas[0].0, session_id);
assert_eq!(transcript_deltas[0].1, "spoken word");
assert!(
sink.text_deltas.lock().unwrap().is_empty(),
"AssistantTranscriptDelta must not reach the text-lane sink (T6)"
);
}
#[tokio::test]
async fn realtime_transcript_observation_routes_to_append_realtime_transcript() {
use meerkat_core::RealtimeTranscriptRole;
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let event = RealtimeTranscriptEvent::ItemObserved {
item_id: "item_realtime_1".into(),
previous_item_id: Some("item_realtime_0".into()),
role: RealtimeTranscriptRole::Assistant,
response_id: Some("resp_realtime_1".into()),
};
let obs = LiveAdapterObservation::RealtimeTranscript {
event: event.clone(),
};
let routing = LiveAdapterHost::classify_observation(&obs);
match routing {
ObservationRouting::AppendRealtimeTranscript => {}
other => panic!("expected AppendRealtimeTranscript, got {other:?}"),
}
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(
matches!(outcome, ObservationOutcome::TranscriptAppended),
"expected TranscriptAppended, got {outcome:?}"
);
let recorded = sink.realtime_events.lock().unwrap();
assert_eq!(recorded.len(), 1, "sink must see exactly one append");
assert_eq!(recorded[0].0, session_id);
assert_eq!(recorded[0].1, event);
assert!(sink.text_deltas.lock().unwrap().is_empty());
assert!(sink.transcript_deltas.lock().unwrap().is_empty());
assert!(sink.text_finals.lock().unwrap().is_empty());
assert!(sink.transcript_finals.lock().unwrap().is_empty());
assert!(sink.user_transcripts.lock().unwrap().is_empty());
assert!(sink.turn_completed.lock().unwrap().is_empty());
assert!(sink.interrupts.lock().unwrap().is_empty());
}
#[tokio::test]
async fn noop_projection_sink_explicitly_accepts_realtime_transcript() {
use meerkat_core::RealtimeTranscriptRole;
let sink: Arc<dyn LiveProjectionSink> = Arc::new(NoOpProjectionSink);
let host = LiveAdapterHost::new(Arc::clone(&sink));
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let event = RealtimeTranscriptEvent::ItemObserved {
item_id: "item_noop".into(),
previous_item_id: None,
role: RealtimeTranscriptRole::Assistant,
response_id: Some("resp_noop".into()),
};
let obs = LiveAdapterObservation::RealtimeTranscript {
event: event.clone(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(
matches!(outcome, ObservationOutcome::TranscriptAppended),
"NoOpProjectionSink must accept RealtimeTranscript explicitly, got {outcome:?}"
);
sink.append_realtime_transcript(&session_id, &event)
.await
.expect("NoOpProjectionSink::append_realtime_transcript must be explicit Ok");
}
#[tokio::test]
async fn realtime_transcript_assistant_turn_completed_routes_through_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let event = RealtimeTranscriptEvent::AssistantTurnCompleted {
response_id: "resp_complete".into(),
stop_reason: StopReason::EndTurn,
usage: Usage::default(),
};
let obs = LiveAdapterObservation::RealtimeTranscript {
event: event.clone(),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
assert!(matches!(outcome, ObservationOutcome::TranscriptAppended));
let recorded = sink.realtime_events.lock().unwrap();
assert_eq!(recorded.len(), 1);
assert_eq!(recorded[0].1, event);
}
#[tokio::test]
async fn tool_call_observation_dispatches_through_tool_authority() {
let sink = Arc::new(RecordingProjectionSink::default());
let dispatcher = Arc::new(RecordingDispatcher::default());
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _)
.with_live_tool_dispatcher(Arc::clone(&dispatcher) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_42"),
tool_name: ToolName::new("calculator"),
arguments: serde_json::json!({"a": 2, "b": 3}),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::ToolCallDispatched {
provider_call_id,
tool_name,
} => {
assert_eq!(provider_call_id, "call_42");
assert_eq!(tool_name, "calculator");
}
other => panic!("expected ToolCallDispatched, got {other:?}"),
}
let calls = dispatcher.calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].0, "call_42");
assert_eq!(calls[0].1, "calculator");
let submitted = adapter.submitted_results.lock().unwrap();
assert_eq!(submitted.len(), 1);
assert_eq!(submitted[0].call_id.0, "call_42");
}
#[tokio::test]
async fn tool_call_skipped_when_no_dispatcher_wired() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_99"),
tool_name: ToolName::new("calculator"),
arguments: serde_json::json!({}),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::ToolCallSkipped {
reason: ToolDispatchSkipReason::NoDispatcher,
..
} => {}
other => panic!("expected ToolCallSkipped/NoDispatcher, got {other:?}"),
}
}
#[tokio::test]
async fn tool_call_no_dispatcher_submits_error_to_adapter() {
let sink = Arc::new(RecordingProjectionSink::default());
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_unwired"),
tool_name: ToolName::new("calculator"),
arguments: serde_json::json!({}),
};
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::ToolCallSkipped {
provider_call_id,
tool_name,
reason: ToolDispatchSkipReason::NoDispatcher,
} => {
assert_eq!(provider_call_id, "call_unwired");
assert_eq!(tool_name, "calculator");
}
other => panic!("expected ToolCallSkipped/NoDispatcher, got {other:?}"),
}
let errors = adapter.submitted_errors.lock().unwrap();
assert_eq!(
errors.len(),
1,
"adapter must receive exactly one SubmitToolError when dispatcher is missing"
);
assert_eq!(errors[0].0, "call_unwired");
assert!(
errors[0].1.contains("dispatcher"),
"error message should mention the missing dispatcher; got {:?}",
errors[0].1
);
assert!(adapter.submitted_results.lock().unwrap().is_empty());
}
#[tokio::test]
async fn set_live_tool_dispatcher_late_binds_after_construction() {
let sink = Arc::new(RecordingProjectionSink::default());
let dispatcher = Arc::new(RecordingDispatcher::default());
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_pre"),
tool_name: ToolName::new("calc"),
arguments: serde_json::json!({}),
};
match host.apply_observation(&ch, &obs).await.unwrap() {
ObservationOutcome::ToolCallSkipped {
reason: ToolDispatchSkipReason::NoDispatcher,
..
} => {}
other => panic!("expected pre-set skip, got {other:?}"),
}
assert_eq!(dispatcher.calls.lock().unwrap().len(), 0);
host.set_live_tool_dispatcher(Arc::clone(&dispatcher) as _);
let obs2 = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_post"),
tool_name: ToolName::new("calc"),
arguments: serde_json::json!({"x": 1}),
};
match host.apply_observation(&ch, &obs2).await.unwrap() {
ObservationOutcome::ToolCallDispatched {
provider_call_id, ..
} => {
assert_eq!(provider_call_id, "call_post");
}
other => panic!("expected post-set dispatch, got {other:?}"),
}
let calls = dispatcher.calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].0, "call_post");
}
#[tokio::test]
async fn set_live_tool_dispatcher_replaces_previously_installed_dispatcher() {
let sink = Arc::new(RecordingProjectionSink::default());
let first = Arc::new(RecordingDispatcher::default());
let second = Arc::new(RecordingDispatcher::default());
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _)
.with_live_tool_dispatcher(Arc::clone(&first) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
host.set_live_tool_dispatcher(Arc::clone(&second) as _);
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_swap"),
tool_name: ToolName::new("calc"),
arguments: serde_json::json!({}),
};
host.apply_observation(&ch, &obs).await.unwrap();
assert_eq!(first.calls.lock().unwrap().len(), 0);
assert_eq!(second.calls.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn tool_call_dispatch_error_submits_tool_error_to_adapter() {
let sink = Arc::new(RecordingProjectionSink::default());
let dispatcher = Arc::new(FailingDispatcher);
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _)
.with_live_tool_dispatcher(Arc::clone(&dispatcher) as _);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_err"),
tool_name: ToolName::new("failing"),
arguments: serde_json::json!({}),
};
host.apply_observation(&ch, &obs).await.unwrap();
let errors = adapter.submitted_errors.lock().unwrap();
assert_eq!(errors.len(), 1);
assert_eq!(errors[0].0, "call_err");
}
struct SlowDispatcher {
sleep_for: Duration,
calls: StdMutex<u32>,
}
impl SlowDispatcher {
fn new(sleep_for: Duration) -> Self {
Self {
sleep_for,
calls: StdMutex::new(0),
}
}
}
#[async_trait]
impl LiveToolDispatcher for SlowDispatcher {
async fn dispatch_live_tool_call(
&self,
_session_id: &SessionId,
call: ToolCall,
) -> Result<ToolDispatchOutcome, LiveToolDispatchError> {
*self.calls.lock().unwrap() += 1;
tokio::time::sleep(self.sleep_for).await;
let tool_result = meerkat_core::types::ToolResult::new(call.id, "ok".into(), false);
Ok(ToolDispatchOutcome::from(tool_result))
}
}
#[tokio::test(start_paused = true)]
async fn realtime_tool_timeout() {
let timeout = Duration::from_millis(500);
let dispatcher_sleep = Duration::from_secs(60);
let sink = Arc::new(RecordingProjectionSink::default());
let dispatcher = Arc::new(SlowDispatcher::new(dispatcher_sleep));
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _)
.with_live_tool_dispatcher(Arc::clone(&dispatcher) as _)
.with_tool_timeout(timeout);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_slow"),
tool_name: ToolName::new("slow_tool"),
arguments: serde_json::json!({"q": 1}),
};
let host_call = async { host.apply_observation(&ch, &obs).await.unwrap() };
let drive_clock = async {
tokio::task::yield_now().await;
tokio::time::advance(timeout + Duration::from_millis(1)).await;
};
let (outcome, _) = tokio::join!(host_call, drive_clock);
match outcome {
ObservationOutcome::ToolCallTimedOut {
provider_call_id,
tool_name,
timeout: t,
} => {
assert_eq!(provider_call_id, "call_slow");
assert_eq!(tool_name, "slow_tool");
assert_eq!(t, timeout);
}
other => panic!("expected ToolCallTimedOut, got {other:?}"),
}
assert_eq!(*dispatcher.calls.lock().unwrap(), 1);
let errors = adapter.submitted_errors.lock().unwrap();
assert_eq!(errors.len(), 1);
assert_eq!(errors[0].0, "call_slow");
assert_eq!(
errors[0].1,
LiveToolDispatchTimeout::new(timeout).to_string(),
"tool dispatch timeout must submit the typed fact's Display projection, \
not a fabricated string: {}",
errors[0].1
);
let results = adapter.submitted_results.lock().unwrap();
assert!(
results.is_empty(),
"no SubmitToolResult should reach the adapter on timeout: {results:?}"
);
assert_eq!(sink.text_finals.lock().unwrap().len(), 0);
assert_eq!(sink.transcript_finals.lock().unwrap().len(), 0);
assert_eq!(sink.terminal_errors.lock().unwrap().len(), 0);
}
#[tokio::test(start_paused = true)]
async fn tool_call_dispatch_succeeds_when_within_deadline() {
let timeout = Duration::from_secs(5);
let dispatcher_sleep = Duration::from_millis(100);
let sink = Arc::new(RecordingProjectionSink::default());
let dispatcher = Arc::new(SlowDispatcher::new(dispatcher_sleep));
let adapter = Arc::new(RecordingAdapter::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _)
.with_live_tool_dispatcher(Arc::clone(&dispatcher) as _)
.with_tool_timeout(timeout);
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::clone(&adapter) as _)
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
let obs = LiveAdapterObservation::ToolCallRequested {
provider_call_id: ToolCallId::new("call_fast"),
tool_name: ToolName::new("fast_tool"),
arguments: serde_json::json!({}),
};
let host_call = async { host.apply_observation(&ch, &obs).await.unwrap() };
let drive_clock = async {
tokio::task::yield_now().await;
tokio::time::advance(dispatcher_sleep + Duration::from_millis(10)).await;
};
let (outcome, _) = tokio::join!(host_call, drive_clock);
match outcome {
ObservationOutcome::ToolCallDispatched {
provider_call_id, ..
} => assert_eq!(provider_call_id, "call_fast"),
other => panic!("expected ToolCallDispatched, got {other:?}"),
}
let results = adapter.submitted_results.lock().unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].call_id.0, "call_fast");
assert!(adapter.submitted_errors.lock().unwrap().is_empty());
}
#[tokio::test]
async fn barge_in_observation_calls_signal_interrupt() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let outcome = host
.apply_observation(
&ch,
&LiveAdapterObservation::TurnInterrupted {
response_id: Some("resp_42".into()),
},
)
.await
.unwrap();
assert!(matches!(outcome, ObservationOutcome::InterruptSignalled));
let interrupts = sink.interrupts.lock().unwrap();
assert_eq!(interrupts.len(), 1);
assert_eq!(interrupts[0].0, session_id);
assert_eq!(interrupts[0].1.as_deref(), Some("resp_42"));
}
#[tokio::test]
async fn transport_barge_in_emits_typed_interrupt_via_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
host.signal_transport_barge_in(&ch).await.unwrap();
let interrupts = sink.interrupts.lock().unwrap();
assert_eq!(
interrupts.len(),
1,
"barge-in must signal exactly one interrupt"
);
assert_eq!(interrupts[0].0, session_id);
assert_eq!(interrupts[0].1, None);
}
#[tokio::test]
async fn transport_output_audio_degradation_emits_typed_signal_via_sink() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
host.signal_output_audio_degraded(&ch, 3).await.unwrap();
{
let degraded = sink.output_audio_degraded.lock().unwrap();
assert_eq!(
degraded.as_slice(),
&[(session_id, 3)],
"degradation must surface exactly once with the cumulative drop count"
);
}
let missing = host
.signal_output_audio_degraded(&LiveChannelId::new("no-such-channel"), 1)
.await;
assert!(matches!(
missing,
Err(LiveAdapterHostError::ChannelNotFound(_))
));
}
#[tokio::test]
async fn terminal_error_observation_signals_sink_without_closing_host_directly() {
let sink = Arc::new(RecordingProjectionSink::default());
let host = LiveAdapterHost::new(Arc::clone(&sink) as _);
let session_id = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(session_id.clone())
.await
.unwrap();
let obs = LiveAdapterObservation::Error {
code: LiveAdapterErrorCode::ConnectionLost,
message: "ws closed unexpectedly".into(),
};
let err = host
.apply_observation(&ch, &obs)
.await
.expect_err("terminal error must wait for generated close authority");
assert!(matches!(err, LiveAdapterHostError::CloseNotAuthorized));
assert_eq!(
host.channel_status(&ch).await.unwrap(),
LiveAdapterStatus::Opening
);
assert_eq!(sink.terminal_errors.lock().unwrap().len(), 0);
let close = host
.close_channel_observed_with_generated_test_machine_authority(&ch)
.await
.unwrap();
assert_eq!(close.channel_id(), ch.as_str());
let outcome = host.apply_observation(&ch, &obs).await.unwrap();
match outcome {
ObservationOutcome::Terminal {
code: LiveAdapterErrorCode::ConnectionLost,
} => {}
other => panic!("expected Terminal/ConnectionLost, got {other:?}"),
}
let status = host.channel_status(&ch).await.unwrap();
assert_eq!(status, LiveAdapterStatus::Closed);
let terminals = sink.terminal_errors.lock().unwrap();
assert_eq!(terminals.len(), 1);
assert_eq!(terminals[0].0, session_id);
assert!(matches!(
terminals[0].1,
LiveAdapterErrorCode::ConnectionLost
));
}
#[tokio::test]
async fn duplicate_session_check_uses_reverse_map() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let s = test_session_id();
let ch = host
.open_channel_with_generated_test_machine_authority(s.clone())
.await
.unwrap();
host.close_channel_with_generated_test_machine_authority(&ch)
.await
.unwrap();
host.open_channel_with_generated_test_machine_authority(s)
.await
.unwrap();
}
#[tokio::test]
async fn transport_send_input_does_not_mint_command_acceptance_evidence() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(StubAdapter::new()))
.await
.unwrap();
host.commit_status_with_generated_test_machine_authority(&ch, LiveAdapterStatus::Ready)
.await
.unwrap();
host.send_input(
&ch,
LiveInputChunk::Text {
text: "hello".into(),
},
)
.await
.unwrap();
{
let inner = host.inner.lock().await;
let channel = inner.channels.get(&ch).unwrap();
assert_eq!(channel.command_acceptance_sequence, 0);
}
let acceptance = host
.send_input_observed(
&ch,
LiveInputChunk::Text {
text: "hello".into(),
},
)
.await
.unwrap();
assert_eq!(acceptance.kind(), LiveCommandAcceptanceKind::SendInput);
assert_eq!(acceptance.acceptance_sequence(), 1);
}
#[tokio::test]
async fn rejected_live_command_does_not_mint_acceptance_evidence() {
let host = LiveAdapterHost::new(Arc::new(NoOpProjectionSink));
let ch = host
.open_channel_with_generated_test_machine_authority(test_session_id())
.await
.unwrap();
host.attach_adapter(&ch, Arc::new(RejectingCommandAdapter))
.await
.unwrap();
let result = host
.send_command_observed(&ch, LiveAdapterCommand::Interrupt)
.await;
assert!(
matches!(
result,
Err(LiveAdapterHostError::AdapterError(
LiveAdapterError::TransportError { .. }
))
),
"adapter rejection must surface before public command acceptance"
);
let inner = host.inner.lock().await;
let channel = inner.channels.get(&ch).unwrap();
assert_eq!(
channel.command_acceptance_sequence, 0,
"rejected command must not mint authority evidence that WebRTC could use to discard output"
);
}
struct StubAdapter;
impl StubAdapter {
fn new() -> Self {
Self
}
}
#[derive(Default)]
struct FailOnceCloseAdapter {
closes: std::sync::atomic::AtomicUsize,
}
#[async_trait]
impl LiveAdapter for FailOnceCloseAdapter {
async fn send_command(&self, _command: LiveAdapterCommand) -> Result<(), LiveAdapterError> {
Ok(())
}
async fn next_observation(
&self,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterError> {
Ok(None)
}
fn status(&self) -> LiveAdapterStatus {
LiveAdapterStatus::Ready
}
async fn close(&self) -> Result<(), LiveAdapterError> {
if self
.closes
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
== 0
{
Err(LiveAdapterError::TransportError {
message: "scripted first physical close failure".to_string(),
})
} else {
Ok(())
}
}
}
#[async_trait]
impl LiveAdapter for StubAdapter {
async fn send_command(&self, _command: LiveAdapterCommand) -> Result<(), LiveAdapterError> {
Ok(())
}
async fn next_observation(
&self,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterError> {
Ok(None)
}
fn status(&self) -> LiveAdapterStatus {
LiveAdapterStatus::Ready
}
async fn close(&self) -> Result<(), LiveAdapterError> {
Ok(())
}
}
struct ErroringAdapter;
#[async_trait]
impl LiveAdapter for ErroringAdapter {
async fn send_command(&self, _command: LiveAdapterCommand) -> Result<(), LiveAdapterError> {
Ok(())
}
async fn next_observation(
&self,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterError> {
Err(LiveAdapterError::TransportError {
message: "pump dead".into(),
})
}
fn status(&self) -> LiveAdapterStatus {
LiveAdapterStatus::Closed
}
async fn close(&self) -> Result<(), LiveAdapterError> {
Ok(())
}
}
struct RejectingCommandAdapter;
#[async_trait]
impl LiveAdapter for RejectingCommandAdapter {
async fn send_command(&self, _command: LiveAdapterCommand) -> Result<(), LiveAdapterError> {
Err(LiveAdapterError::TransportError {
message: "command queue rejected".into(),
})
}
async fn next_observation(
&self,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterError> {
Ok(None)
}
fn status(&self) -> LiveAdapterStatus {
LiveAdapterStatus::Ready
}
async fn close(&self) -> Result<(), LiveAdapterError> {
Ok(())
}
}
#[derive(Default)]
struct RecordingAdapter {
submitted_results: StdMutex<Vec<LiveToolResult>>,
submitted_errors: StdMutex<Vec<(String, String)>>,
}
#[async_trait]
impl LiveAdapter for RecordingAdapter {
async fn send_command(&self, command: LiveAdapterCommand) -> Result<(), LiveAdapterError> {
match command {
LiveAdapterCommand::SubmitToolResult { result } => {
self.submitted_results.lock().unwrap().push(result);
}
LiveAdapterCommand::SubmitToolError { call_id, error } => {
self.submitted_errors
.lock()
.unwrap()
.push((call_id.0, error));
}
_ => {}
}
Ok(())
}
async fn next_observation(
&self,
) -> Result<Option<LiveAdapterObservation>, LiveAdapterError> {
Ok(None)
}
fn status(&self) -> LiveAdapterStatus {
LiveAdapterStatus::Ready
}
async fn close(&self) -> Result<(), LiveAdapterError> {
Ok(())
}
}
#[derive(Default)]
struct RecordingDispatcher {
calls: StdMutex<Vec<(String, String, String)>>,
}
#[async_trait]
impl LiveToolDispatcher for RecordingDispatcher {
async fn dispatch_live_tool_call(
&self,
_session_id: &SessionId,
call: ToolCall,
) -> Result<ToolDispatchOutcome, LiveToolDispatchError> {
self.calls.lock().unwrap().push((
call.id.clone(),
call.name.clone(),
call.args.to_string(),
));
let tool_result = meerkat_core::types::ToolResult::new(call.id, "ok".into(), false);
Ok(ToolDispatchOutcome::from(tool_result))
}
}
struct FailingDispatcher;
#[async_trait]
impl LiveToolDispatcher for FailingDispatcher {
async fn dispatch_live_tool_call(
&self,
_session_id: &SessionId,
_call: ToolCall,
) -> Result<ToolDispatchOutcome, LiveToolDispatchError> {
Err(LiveToolDispatchError::Tool(
meerkat_core::error::ToolError::ExecutionFailed {
message: "bang".into(),
},
))
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
struct OwnedIdentity {
provider_item_id: Option<String>,
previous_item_id: Option<String>,
content_index: Option<u32>,
response_id: Option<String>,
delta_id: Option<String>,
}
impl OwnedIdentity {
fn from_borrowed(identity: LiveTranscriptIdentity<'_>) -> Self {
Self {
provider_item_id: identity.provider_item_id.map(|s| s.to_string()),
previous_item_id: identity.previous_item_id.map(|s| s.to_string()),
content_index: identity.content_index,
response_id: identity.response_id.map(|s| s.to_string()),
delta_id: identity.delta_id.map(|s| s.to_string()),
}
}
}
#[derive(Default)]
#[allow(clippy::type_complexity)]
struct RecordingProjectionSink {
user_transcripts: StdMutex<Vec<(SessionId, String, OwnedIdentity)>>,
text_deltas: StdMutex<Vec<(SessionId, String, OwnedIdentity)>>,
transcript_deltas: StdMutex<Vec<(SessionId, String, OwnedIdentity)>>,
text_finals: StdMutex<
Vec<(
SessionId,
String,
OwnedIdentity,
StopReason,
Usage,
Option<String>,
)>,
>,
transcript_finals: StdMutex<
Vec<(
SessionId,
String,
OwnedIdentity,
StopReason,
Usage,
Option<String>,
)>,
>,
truncations: StdMutex<
Vec<(
SessionId,
Option<String>,
Option<String>,
Option<u32>,
Option<String>,
Option<String>,
)>,
>,
interrupts: StdMutex<Vec<(SessionId, Option<String>)>>,
output_audio_degraded: StdMutex<Vec<(SessionId, u64)>>,
turn_completed: StdMutex<Vec<(SessionId, StopReason, Usage, Option<String>)>>,
terminal_errors: StdMutex<Vec<(SessionId, LiveAdapterErrorCode, String)>>,
realtime_events: StdMutex<Vec<(SessionId, RealtimeTranscriptEvent)>>,
}
#[async_trait]
impl LiveProjectionSink for RecordingProjectionSink {
async fn append_user_transcript(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
self.user_transcripts.lock().unwrap().push((
session_id.clone(),
text.to_string(),
OwnedIdentity::from_borrowed(identity),
));
Ok(())
}
async fn append_assistant_text_delta(
&self,
session_id: &SessionId,
delta: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
self.text_deltas.lock().unwrap().push((
session_id.clone(),
delta.to_string(),
OwnedIdentity::from_borrowed(identity),
));
Ok(())
}
async fn append_assistant_transcript_delta(
&self,
session_id: &SessionId,
delta: &str,
identity: LiveTranscriptIdentity<'_>,
) -> Result<(), LiveProjectionError> {
self.transcript_deltas.lock().unwrap().push((
session_id.clone(),
delta.to_string(),
OwnedIdentity::from_borrowed(identity),
));
Ok(())
}
async fn append_assistant_text_final(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
self.text_finals.lock().unwrap().push((
session_id.clone(),
text.to_string(),
OwnedIdentity::from_borrowed(identity),
stop_reason,
usage,
response_id.map(|s| s.to_string()),
));
Ok(())
}
async fn append_assistant_transcript_final(
&self,
session_id: &SessionId,
text: &str,
identity: LiveTranscriptIdentity<'_>,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
self.transcript_finals.lock().unwrap().push((
session_id.clone(),
text.to_string(),
OwnedIdentity::from_borrowed(identity),
stop_reason,
usage,
response_id.map(|s| s.to_string()),
));
Ok(())
}
async fn truncate_assistant_transcript(
&self,
session_id: &SessionId,
provider_item_id: Option<&str>,
previous_item_id: Option<&str>,
content_index: Option<u32>,
response_id: Option<&str>,
text: Option<&str>,
) -> Result<(), LiveProjectionError> {
self.truncations.lock().unwrap().push((
session_id.clone(),
provider_item_id.map(|s| s.to_string()),
previous_item_id.map(|s| s.to_string()),
content_index,
response_id.map(|s| s.to_string()),
text.map(|s| s.to_string()),
));
Ok(())
}
async fn signal_turn_interrupt(
&self,
session_id: &SessionId,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
self.interrupts
.lock()
.unwrap()
.push((session_id.clone(), response_id.map(|s| s.to_string())));
Ok(())
}
async fn signal_output_audio_degraded(
&self,
session_id: &SessionId,
dropped: u64,
) -> Result<(), LiveProjectionError> {
self.output_audio_degraded
.lock()
.unwrap()
.push((session_id.clone(), dropped));
Ok(())
}
async fn signal_turn_completed(
&self,
session_id: &SessionId,
stop_reason: StopReason,
usage: Usage,
response_id: Option<&str>,
) -> Result<(), LiveProjectionError> {
self.turn_completed.lock().unwrap().push((
session_id.clone(),
stop_reason,
usage,
response_id.map(|s| s.to_string()),
));
Ok(())
}
async fn signal_terminal_error(
&self,
session_id: &SessionId,
code: LiveAdapterErrorCode,
message: &str,
) -> Result<(), LiveProjectionError> {
self.terminal_errors.lock().unwrap().push((
session_id.clone(),
code,
message.to_string(),
));
Ok(())
}
async fn append_realtime_transcript(
&self,
session_id: &SessionId,
event: &RealtimeTranscriptEvent,
) -> Result<RealtimeTranscriptApplyOutcome, LiveProjectionError> {
self.realtime_events
.lock()
.unwrap()
.push((session_id.clone(), event.clone()));
Ok(RealtimeTranscriptApplyOutcome::default())
}
}
#[test]
fn session_error_classification_lands_in_distinct_typed_variants() {
use meerkat_core::SessionError;
let sid = test_session_id();
assert!(matches!(
LiveProjectionError::from_session_error(
&sid,
SessionError::NotFound { id: sid.clone() }
),
LiveProjectionError::SessionNotFound(_)
));
match LiveProjectionError::from_session_error(
&sid,
SessionError::Unsupported("nope".into()),
) {
LiveProjectionError::Rejected(reason) => assert_eq!(reason, "nope"),
other => panic!("expected Rejected, got {other:?}"),
}
assert!(matches!(
LiveProjectionError::from_session_error(&sid, SessionError::Busy { id: sid.clone() }),
LiveProjectionError::SessionBusy(_)
));
assert!(matches!(
LiveProjectionError::from_session_error(
&sid,
SessionError::NotRunning { id: sid.clone() }
),
LiveProjectionError::SessionNotRunning(_)
));
match LiveProjectionError::from_session_error(&sid, SessionError::PersistenceDisabled) {
LiveProjectionError::CapabilityDisabled { code, .. } => {
assert_eq!(code, "SESSION_PERSISTENCE_DISABLED");
}
other => panic!("expected CapabilityDisabled, got {other:?}"),
}
match LiveProjectionError::from_session_error(&sid, SessionError::CompactionDisabled) {
LiveProjectionError::CapabilityDisabled { code, .. } => {
assert_eq!(code, "SESSION_COMPACTION_DISABLED");
}
other => panic!("expected CapabilityDisabled, got {other:?}"),
}
let store_err: Box<dyn std::error::Error + Send + Sync> = "disk gone".into();
match LiveProjectionError::from_session_error(&sid, SessionError::Store(store_err)) {
LiveProjectionError::Session { code, .. } => {
assert_eq!(code, "SESSION_STORE_ERROR");
}
other => panic!("expected typed Session, got {other:?}"),
}
match LiveProjectionError::from_session_error(
&sid,
SessionError::FailedWithData {
message: "boom".into(),
data: serde_json::json!({"k": "v"}),
},
) {
LiveProjectionError::Session { code, message } => {
assert_eq!(code, "SESSION_ERROR");
assert_eq!(message, "boom");
}
other => panic!("expected typed Session, got {other:?}"),
}
}
#[test]
fn live_tool_dispatch_session_error_classification_lands_in_distinct_typed_variants() {
use meerkat_core::SessionError;
let sid = test_session_id();
assert!(matches!(
LiveToolDispatchError::from_session_error(
&sid,
SessionError::NotFound { id: sid.clone() }
),
LiveToolDispatchError::SessionNotFound(_)
));
match LiveToolDispatchError::from_session_error(
&sid,
SessionError::Unsupported("no live tool dispatcher".into()),
) {
LiveToolDispatchError::Rejected(reason) => {
assert_eq!(reason, "no live tool dispatcher");
}
other => panic!("expected Rejected, got {other:?}"),
}
assert!(matches!(
LiveToolDispatchError::from_session_error(&sid, SessionError::Busy { id: sid.clone() }),
LiveToolDispatchError::SessionBusy(_)
));
assert!(matches!(
LiveToolDispatchError::from_session_error(
&sid,
SessionError::NotRunning { id: sid.clone() }
),
LiveToolDispatchError::SessionNotRunning(_)
));
match LiveToolDispatchError::from_session_error(&sid, SessionError::PersistenceDisabled) {
LiveToolDispatchError::CapabilityDisabled { code, .. } => {
assert_eq!(code, "SESSION_PERSISTENCE_DISABLED");
}
other => panic!("expected CapabilityDisabled, got {other:?}"),
}
let store_err: Box<dyn std::error::Error + Send + Sync> = "disk gone".into();
match LiveToolDispatchError::from_session_error(&sid, SessionError::Store(store_err)) {
LiveToolDispatchError::Session { code, .. } => {
assert_eq!(code, "SESSION_STORE_ERROR");
}
other => panic!("expected typed Session, got {other:?}"),
}
}
#[test]
fn delta_identity_requires_response_delta_and_item_ids() {
let full = LiveTranscriptIdentity::assistant_delta(
Some("item-1"),
Some("prev-0"),
Some(2),
Some("resp-9"),
Some("delta-3"),
);
let resolved = full
.require_delta_identity()
.expect("full identity resolves");
assert_eq!(resolved.response_id, "resp-9");
assert_eq!(resolved.delta_id, "delta-3");
assert_eq!(resolved.item_id, "item-1");
assert_eq!(resolved.previous_item_id, Some("prev-0"));
assert_eq!(resolved.content_index, Some(2));
let no_resp = LiveTranscriptIdentity::assistant_delta(
Some("item-1"),
None,
Some(0),
None,
Some("delta-3"),
);
assert_eq!(
no_resp.require_delta_identity(),
Err(LiveTranscriptIdentityError::MissingResponseId)
);
let no_delta = LiveTranscriptIdentity::assistant_delta(
Some("item-1"),
None,
Some(0),
Some("resp-9"),
None,
);
assert_eq!(
no_delta.require_delta_identity(),
Err(LiveTranscriptIdentityError::MissingDeltaId)
);
let no_item = LiveTranscriptIdentity::assistant_delta(
None,
None,
Some(0),
Some("resp-9"),
Some("delta-3"),
);
assert_eq!(
no_item.require_delta_identity(),
Err(LiveTranscriptIdentityError::MissingItemId)
);
}
}