use super::active_registry::{ActiveTransactionRegistry, ClaimSessionError, ControlMessage};
use super::callback_service::{CallbackReservation, CallbackService};
use super::dispatcher::{DispatchOutcome, DispatcherLimits, TransactionToolDispatcher};
use super::events::{BoundedEventSender, OrderedEventPublisher};
use super::exchange::{
run_encoded_exchange, run_exchange, EncodedExchangeParams, ExchangeFailure, ExchangeParams,
};
use super::executor_spawn::try_spawn;
use super::finalization::{build_transaction_end, FinalizationGuard};
use super::loop_adapters::dispatch_ready_tool_cancellable;
use super::mcp::{CapabilityToken, McpGatewayHandle, PendingMcpBinding};
use super::resolved_tools::ResolvedToolSet;
use super::tool_capacity::SharedToolCapacity;
use monoloop_connector::{Connector, SessionAdapter};
use monoloop_contracts::{
CanonicalAssistantToolCall, CanonicalMessage, CanonicalToolError, CanonicalToolResult,
CanonicalUnit, CanonicalUnitEvent, ChannelId, ChannelKind, ContinuationContext,
ContinuationPolicy, EffectiveConfig, EventDeliveryOutcome, ExchangeId, ExternalSessionId,
InterpretationLimits, McpConfigurationCapability, McpReachability, OutboundDialectEncoder,
SessionId, SessionKey, TextChannel, TextPart, ToolActionId, ToolExecutionMode, ToolId,
ToolLifecycleEvent, ToolName, ToolRequestState, TransactionEndKind, TransactionEventPayload,
TransactionId,
};
use monoloop_interpreter::InterpreterFactory;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::runtime::Handle;
use tokio::sync::{mpsc, oneshot, watch};
pub struct ActorSpawn {
pub executor: Handle,
pub transaction_id: TransactionId,
pub channel_id: ChannelId,
#[allow(dead_code)]
pub channel_kind: ChannelKind,
pub tool_mode: ToolExecutionMode,
pub mcp_configuration: McpConfigurationCapability,
pub mcp_reachability: McpReachability,
pub mcp: Option<McpGatewayHandle>,
pub session_key: Option<SessionKey>,
pub provisional_external: bool,
pub existing_session: bool,
pub sessions: Option<Arc<dyn SessionAdapter>>,
pub connector: Arc<dyn Connector>,
pub encoder: Arc<dyn OutboundDialectEncoder>,
pub interpreter: Arc<dyn InterpreterFactory>,
pub endpoint_ref: String,
pub credential_ref: Option<String>,
pub input: monoloop_contracts::CanonicalInput,
pub effective: EffectiveConfig,
pub tools: ResolvedToolSet,
pub guard: Arc<FinalizationGuard>,
pub control_rx: mpsc::Receiver<ControlMessage>,
pub event_tx: BoundedEventSender,
pub delivery_fail_rx: mpsc::Receiver<()>,
pub registry: Arc<Mutex<ActiveTransactionRegistry>>,
pub release_capacity: Arc<dyn Fn() + Send + Sync>,
pub deadline: Duration,
pub max_continuations: usize,
pub max_provider_exchanges: usize,
pub max_concurrent_tools: usize,
pub max_queued_tools: usize,
pub cleanup_deadline: Duration,
pub terminal_event_delivery_deadline: Duration,
pub callback_deadline: Duration,
pub max_total_provider_input_bytes: usize,
pub max_total_provider_output_bytes: usize,
pub max_continuation_context_bytes: usize,
pub max_tool_schema_bytes: usize,
pub max_tool_payload_bytes: usize,
pub max_tool_output_bytes: usize,
pub max_encoded_exchange_bytes: usize,
pub max_distinct_sessions: usize,
pub max_diagnostic_count: usize,
pub max_diagnostic_bytes: usize,
pub callbacks: CallbackService,
pub callback_reservation: CallbackReservation,
pub start_gate: oneshot::Receiver<()>,
}
pub fn spawn_actor(spawn: ActorSpawn) -> Result<tokio::task::JoinHandle<()>, ()> {
let executor = spawn.executor.clone();
try_spawn(&executor, async move {
run_actor(spawn).await;
})
}
struct ActorResult {
kind: TransactionEndKind,
prior: Option<TransactionEndKind>,
delivery: EventDeliveryOutcome,
session_key: Option<SessionKey>,
}
async fn run_actor(spawn: ActorSpawn) {
let ActorSpawn {
executor,
transaction_id,
channel_id,
channel_kind: _,
tool_mode,
mcp_configuration,
mcp_reachability,
mcp,
mut session_key,
provisional_external,
existing_session,
sessions,
connector,
encoder,
interpreter,
endpoint_ref,
credential_ref,
input,
effective,
tools,
guard,
mut control_rx,
event_tx,
mut delivery_fail_rx,
registry,
release_capacity,
deadline,
max_continuations,
max_provider_exchanges,
max_concurrent_tools,
max_queued_tools,
cleanup_deadline,
terminal_event_delivery_deadline,
callback_deadline,
max_total_provider_input_bytes,
max_total_provider_output_bytes,
max_continuation_context_bytes,
max_tool_schema_bytes,
max_tool_payload_bytes,
max_tool_output_bytes,
max_encoded_exchange_bytes,
max_distinct_sessions,
max_diagnostic_count,
max_diagnostic_bytes,
callbacks,
callback_reservation,
start_gate,
} = spawn;
let _ = (max_diagnostic_count, max_diagnostic_bytes);
let events = OrderedEventPublisher::new(event_tx.clone(), Arc::clone(guard.sequencer()));
let (session_watch_tx, session_watch_rx) = watch::channel(session_key.clone());
match start_gate.await {
Ok(()) => {}
Err(_) => {
drop(callback_reservation);
release_capacity();
let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
let _ = reg.remove(&transaction_id);
return;
}
}
{
let mut schema_total = 0usize;
for spec in tools.specs() {
schema_total = schema_total.saturating_add(
serde_json::to_vec(spec.input_schema.as_value())
.map(|b| b.len())
.unwrap_or(0),
);
}
if schema_total > max_tool_schema_bytes {
let result = ActorResult {
kind: TransactionEndKind::LimitExceeded,
prior: None,
delivery: EventDeliveryOutcome::Failed,
session_key: session_key.clone(),
};
finalize_and_cleanup(
transaction_id,
channel_id,
guard,
events,
registry,
release_capacity,
result,
terminal_event_delivery_deadline,
callback_deadline,
&callbacks,
callback_reservation,
)
.await;
return;
}
}
let cleanup_deadline = cleanup_deadline.max(Duration::from_millis(50));
let mut terminal_kind = TransactionEndKind::Completed;
let mut attachment: Option<Arc<monoloop_connector::SessionAttachment>> = None;
let mcp_token: Arc<Mutex<Option<CapabilityToken>>> = Arc::new(Mutex::new(None));
let mcp_token_work = Arc::clone(&mcp_token);
let work = async {
let mut pending_mcp: Option<PendingMcpBinding> = None;
let mut create_mode_attach = false;
if let Some(ref adapter) = sessions {
let mut initial_mcp = None;
if tool_mode == ToolExecutionMode::McpGateway
&& !tools.is_empty()
&& mcp_reachability == McpReachability::SameLoopbackNamespace
{
if mcp_configuration == McpConfigurationCapability::CreationOnly && existing_session
{
return Err(TransactionEndKind::InvariantFailed);
}
if !existing_session || mcp_configuration == McpConfigurationCapability::Refreshable
{
let handle = mcp.as_ref().ok_or(TransactionEndKind::InvariantFailed)?;
let sk = session_key.clone().unwrap_or_else(|| {
SessionKey::new(channel_id.clone(), SessionId::generate())
});
let dispatcher = TransactionToolDispatcher::with_limits(
transaction_id,
sk,
tools.clone(),
SharedToolCapacity::unlimited(),
DispatcherLimits {
max_concurrent_tools,
max_queued_tools,
max_tool_payload_bytes,
max_tool_output_bytes,
},
);
let pending = handle
.install_pending(
transaction_id,
tools.clone(),
dispatcher,
monoloop_contracts::ExchangeId::generate(),
)
.map_err(|_| TransactionEndKind::InvariantFailed)?;
{
let mut g = mcp_token_work.lock().unwrap_or_else(|e| e.into_inner());
*g = Some(pending.token.clone());
}
if !existing_session
&& mcp_configuration == McpConfigurationCapability::CreationOnly
{
initial_mcp = Some(pending.descriptor.clone());
}
pending_mcp = Some(pending);
}
}
let requested = if provisional_external {
None
} else {
session_key.as_ref().map(|k| k.session_id.clone())
};
let req = monoloop_connector::SessionAttachRequest {
transaction_id,
channel_id: channel_id.clone(),
requested_session_id: requested,
session_config: effective.session.clone(),
initial_mcp,
deadline: std::time::Instant::now() + deadline,
};
let pending = adapter
.begin_attach(req)
.map_err(|_| TransactionEndKind::ChannelOpenFailed)?;
let att = tokio::select! {
biased;
ctrl = control_rx.recv() => {
let _ = pending.control.cancel();
return Err(match ctrl {
Some(ControlMessage::ForceTerminate) => TransactionEndKind::Terminated,
_ => TransactionEndKind::Cancelled,
});
}
r = pending.completion => {
r.map_err(|_| TransactionEndKind::ChannelOpenFailed)?
}
};
create_mode_attach = att.create_mode;
if !att.create_mode {
let sid = SessionId::from_external(&att.external_session_id);
if let Some(ref expected) = session_key {
if expected.session_id.as_str() != sid.as_str() {
return Err(TransactionEndKind::InvariantFailed);
}
}
let key = SessionKey::new(channel_id.clone(), sid);
if session_key.is_none() {
let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
match reg.claim_session(
transaction_id,
key.clone(),
Some(max_distinct_sessions),
) {
Ok(()) => {}
Err(ClaimSessionError::Collision) => {
return Err(TransactionEndKind::InvariantFailed);
}
Err(ClaimSessionError::CapacityExceeded) => {
return Err(TransactionEndKind::LimitExceeded);
}
Err(_) => return Err(TransactionEndKind::InvariantFailed),
}
}
session_key = Some(key);
if !emit_unit_or_session(
&events,
transaction_id,
&channel_id,
&session_key,
TransactionEventPayload::SessionEstablished {
external_session_id: att.external_session_id.clone(),
},
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
let _ = session_watch_tx.send(session_key.clone());
if let Some(ref pending) = pending_mcp {
if mcp_configuration == McpConfigurationCapability::Refreshable {
if let (Some(att_ref), Some(adapter)) =
(Some(Arc::clone(&att)), sessions.as_ref())
{
if let Ok(pending_cfg) =
adapter.begin_refresh_mcp(att_ref, Some(pending.descriptor.clone()))
{
tokio::select! {
biased;
ctrl = control_rx.recv() => {
let _ = pending_cfg.control.cancel();
return Err(match ctrl {
Some(ControlMessage::ForceTerminate) => {
TransactionEndKind::Terminated
}
_ => TransactionEndKind::Cancelled,
});
}
r = pending_cfg.completion => {
r.map_err(|_| TransactionEndKind::InvariantFailed)?;
}
}
}
}
}
if let Some(handle) = mcp.as_ref() {
handle
.activate(&pending.token)
.map_err(|_| TransactionEndKind::InvariantFailed)?;
}
}
}
attachment = Some(att);
} else if provisional_external {
let sid = SessionId::generate();
let key = SessionKey::new(channel_id.clone(), sid.clone());
{
let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
match reg.claim_session(transaction_id, key.clone(), Some(max_distinct_sessions)) {
Ok(()) => {}
Err(ClaimSessionError::Collision) => {
return Err(TransactionEndKind::InvariantFailed);
}
Err(ClaimSessionError::CapacityExceeded) => {
return Err(TransactionEndKind::LimitExceeded);
}
Err(_) => return Err(TransactionEndKind::InvariantFailed),
}
}
session_key = Some(key);
let _ = session_watch_tx.send(session_key.clone());
let ext = ExternalSessionId::try_new(sid.as_str())
.map_err(|_| TransactionEndKind::InvariantFailed)?;
if !emit_unit_or_session(
&events,
transaction_id,
&channel_id,
&session_key,
TransactionEventPayload::SessionEstablished {
external_session_id: ext,
},
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
let _ = (existing_session, provisional_external);
let tool_specs: Vec<_> = tools.specs().into_iter().cloned().collect();
let max_exchanges = max_provider_exchanges.max(1);
let max_inline = max_continuations;
let mut exchanges_done = 0usize;
let mut continuations_done = 0usize;
let mut provider_input_bytes = 0usize;
let mut provider_output_bytes = 0usize;
let mut continuation_messages = input.messages().to_vec();
let _ = max_total_provider_output_bytes;
let live_cap = (max_total_provider_output_bytes / 256)
.clamp(8, max_provider_exchanges.saturating_mul(64).max(64));
let (live_tx, mut live_rx) = mpsc::channel::<CanonicalUnitEvent>(live_cap);
let events_live = events.clone();
let channel_live = channel_id.clone();
let mut session_watch_live = session_watch_rx.clone();
let live_join = match try_spawn(&executor, async move {
if session_watch_live
.wait_for(|sk| sk.is_some())
.await
.is_err()
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
while let Some(unit) = live_rx.recv().await {
let session_live = session_watch_live.borrow().clone();
if !emit_canonical_unit(
&events_live,
transaction_id,
&channel_live,
&session_live,
unit,
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
Ok(())
}) {
Ok(h) => h,
Err(()) => return Err(TransactionEndKind::InvariantFailed),
};
let (session_id_tx, claim_join) = if create_mode_attach {
let (sess_tx, sess_rx) = oneshot::channel();
let registry_c = Arc::clone(®istry);
let events_c = events.clone();
let guard_c = Arc::clone(&guard);
let channel_c = channel_id.clone();
let session_watch_c = session_watch_tx.clone();
let max_distinct = max_distinct_sessions;
let join = match try_spawn(&executor, async move {
let Ok(ext) = sess_rx.await else {
return Ok::<(), TransactionEndKind>(());
};
let sid = SessionId::from_external(&ext);
let key = SessionKey::new(channel_c.clone(), sid);
{
let mut reg = registry_c.lock().unwrap_or_else(|e| e.into_inner());
match reg.claim_session(transaction_id, key.clone(), Some(max_distinct)) {
Ok(()) => {}
Err(ClaimSessionError::Collision) => {
return Err(TransactionEndKind::InvariantFailed);
}
Err(ClaimSessionError::CapacityExceeded) => {
return Err(TransactionEndKind::LimitExceeded);
}
Err(_) => return Err(TransactionEndKind::InvariantFailed),
}
}
let sk = Some(key);
guard_c.set_session_id(
sk.as_ref()
.map(|k| k.session_id.clone())
.unwrap_or_else(SessionId::generate),
);
if !emit_unit_or_session(
&events_c,
transaction_id,
&channel_c,
&sk,
TransactionEventPayload::SessionEstablished {
external_session_id: ext,
},
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
let _ = session_watch_c.send(sk);
Ok(())
}) {
Ok(h) => h,
Err(()) => {
live_join.abort();
return Err(TransactionEndKind::InvariantFailed);
}
};
(Some(sess_tx), Some(join))
} else {
(None, None)
};
let mut outcome = tokio::select! {
biased;
ctrl = control_rx.recv() => {
live_join.abort();
let _ = tokio::time::timeout(cleanup_deadline, live_join).await;
if let Some(j) = claim_join {
j.abort();
let _ = tokio::time::timeout(cleanup_deadline, j).await;
}
return Err(match ctrl {
Some(ControlMessage::ForceTerminate) => TransactionEndKind::Terminated,
_ => TransactionEndKind::Cancelled,
});
}
r = run_exchange(ExchangeParams {
executor: &executor,
transaction_id,
connector: connector.as_ref(),
encoder: encoder.as_ref(),
interpreter: interpreter.as_ref(),
endpoint_ref: &endpoint_ref,
credential_ref: credential_ref.as_deref(),
session_attachment: attachment.clone(),
input: &input,
config: &effective,
tools: &tool_specs,
interpretation_limits: InterpretationLimits::default(),
deadline,
cleanup_deadline,
max_encoded_exchange_bytes,
unit_tx: Some(live_tx),
session_id_tx,
}) => r,
}
.map_err(map_exchange_failure)?;
exchanges_done += 1;
if let Some(j) = claim_join {
j.await.map_err(|_| TransactionEndKind::InvariantFailed)??;
let Some(ext) = outcome.external_session_id.clone() else {
return Err(TransactionEndKind::InvariantFailed);
};
session_key = Some(SessionKey::new(
channel_id.clone(),
SessionId::from_external(&ext),
));
let _ = session_watch_tx.send(session_key.clone());
if let Some(ref pending) = pending_mcp {
if let Some(handle) = mcp.as_ref() {
handle
.activate(&pending.token)
.map_err(|_| TransactionEndKind::InvariantFailed)?;
}
}
}
for u in &outcome.units {
provider_output_bytes = provider_output_bytes.saturating_add(estimate_unit_bytes(u));
}
if provider_output_bytes > max_total_provider_output_bytes {
return Err(TransactionEndKind::LimitExceeded);
}
live_join
.await
.map_err(|_| TransactionEndKind::InvariantFailed)??;
loop {
if let Some(fail) = outcome.failure {
return Err(map_exchange_failure(fail));
}
if tool_mode != ToolExecutionMode::ModelToolCalls {
break;
}
let ready = collect_ready_tools(outcome.exchange_id, &outcome.units);
if ready.is_empty() {
break;
}
let sk = session_key
.clone()
.ok_or(TransactionEndKind::InvariantFailed)?;
let dispatcher = TransactionToolDispatcher::with_limits(
transaction_id,
sk.clone(),
tools.clone(),
SharedToolCapacity::unlimited(),
DispatcherLimits {
max_concurrent_tools,
max_queued_tools,
max_tool_payload_bytes,
max_tool_output_bytes,
},
);
let mut results: Vec<CanonicalToolResult> = Vec::with_capacity(ready.len());
for (ord, (action_id, name, payload, provider_id)) in ready.into_iter().enumerate() {
let tool_cancel = Arc::new(tokio::sync::Notify::new());
let tool_cancel_dispatch = Arc::clone(&tool_cancel);
let dispatch_fut = dispatch_ready_tool_cancellable(
&dispatcher,
outcome.exchange_id,
action_id,
&name,
&provider_id,
ord as u32,
&payload,
Some(tool_cancel_dispatch),
);
tokio::pin!(dispatch_fut);
let dispatch_outcome = tokio::select! {
biased;
ctrl = control_rx.recv() => {
tool_cancel.notify_waiters();
let _ = dispatch_fut.await;
return Err(match ctrl {
Some(ControlMessage::ForceTerminate) => {
TransactionEndKind::Terminated
}
_ => TransactionEndKind::Cancelled,
});
}
r = &mut dispatch_fut => r,
};
let mut rejection_result: Option<CanonicalToolResult> = None;
match &dispatch_outcome {
DispatchOutcome::Canonical { result, .. } => {
results.push(result.clone());
}
DispatchOutcome::Rejected {
tool_action_id,
code,
message,
..
} => {
let err = CanonicalToolError::try_new(*code, message.as_str(), None, 256)
.unwrap_or_else(|_| {
CanonicalToolError::try_new("tool_rejected", "rejected", None, 256)
.expect("static")
});
let result = CanonicalToolResult {
transaction_id,
session_key: sk.clone(),
exchange_id: outcome.exchange_id,
tool_action_id: tool_action_id.clone(),
tool_id: ToolId::try_new("rejected")
.unwrap_or_else(|_| ToolId::try_new("unknown").expect("static")),
provider_tool_call_id: provider_id.clone(),
request_ordinal: ord as u32,
outcome: monoloop_contracts::CanonicalToolResultOutcome::DomainFailed(
err,
),
};
results.push(result.clone());
rejection_result = Some(result);
}
DispatchOutcome::RuntimeFailed { .. } => {}
}
emit_dispatch_outcome(
&events,
transaction_id,
&channel_id,
&session_key,
dispatch_outcome,
)
.await?;
if let Some(result) = rejection_result {
if !emit_tool_lifecycle(
&events,
transaction_id,
&channel_id,
&session_key,
ToolLifecycleEvent::Completed { result },
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
}
if results.is_empty() {
break;
}
match effective.continuation_policy {
ContinuationPolicy::CallerControlled => {
return Err(TransactionEndKind::ContinuationRequired);
}
ContinuationPolicy::InlineToolContinuation => {
if continuations_done >= max_inline || exchanges_done >= max_exchanges {
return Err(TransactionEndKind::LimitExceeded);
}
append_exchange_to_transcript(
&mut continuation_messages,
&outcome.units,
&results,
)
.map_err(|_| TransactionEndKind::EncodingFailed)?;
let context = ContinuationContext::try_new(continuation_messages.clone())
.map_err(|_| TransactionEndKind::EncodingFailed)?;
let context_bytes = estimate_input_bytes_from_messages(&continuation_messages);
if context_bytes > max_continuation_context_bytes {
return Err(TransactionEndKind::LimitExceeded);
}
let exchange_id = monoloop_contracts::ExchangeId::generate();
let encoded = encoder
.encode_tool_continuation(
monoloop_contracts::ToolContinuationEncodeRequest {
transaction_id: &transaction_id,
exchange_id: &exchange_id,
context: &context,
results: &results,
config: &effective,
tools: &tool_specs,
},
)
.map_err(|_| TransactionEndKind::EncodingFailed)?;
if encoded.bytes.len() > max_continuation_context_bytes
|| encoded.bytes.len() > max_encoded_exchange_bytes
{
return Err(TransactionEndKind::LimitExceeded);
}
provider_input_bytes = provider_input_bytes.saturating_add(encoded.bytes.len());
if provider_input_bytes > max_total_provider_input_bytes {
return Err(TransactionEndKind::LimitExceeded);
}
let live_cap2 = (max_total_provider_output_bytes / 256).clamp(8, 64);
let (live_tx2, mut live_rx2) = mpsc::channel::<CanonicalUnitEvent>(live_cap2);
let events_live2 = events.clone();
let channel_live = channel_id.clone();
let session_live = session_key.clone();
let live_join2 = match try_spawn(&executor, async move {
while let Some(unit) = live_rx2.recv().await {
if !emit_canonical_unit(
&events_live2,
transaction_id,
&channel_live,
&session_live,
unit,
)
.await
{
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
Ok(())
}) {
Ok(h) => h,
Err(()) => return Err(TransactionEndKind::InvariantFailed),
};
outcome = tokio::select! {
biased;
ctrl = control_rx.recv() => {
live_join2.abort();
let _ = tokio::time::timeout(cleanup_deadline, live_join2).await;
return Err(match ctrl {
Some(ControlMessage::ForceTerminate) => {
TransactionEndKind::Terminated
}
_ => TransactionEndKind::Cancelled,
});
}
r = run_encoded_exchange(EncodedExchangeParams {
executor: &executor,
transaction_id,
exchange_id,
connector: connector.as_ref(),
interpreter: interpreter.as_ref(),
endpoint_ref: &endpoint_ref,
credential_ref: credential_ref.as_deref(),
session_attachment: attachment.clone(),
encoded,
interpretation_limits: InterpretationLimits::default(),
deadline,
cleanup_deadline,
max_encoded_exchange_bytes,
unit_tx: Some(live_tx2),
}) => r,
}
.map_err(map_exchange_failure)?;
live_join2
.await
.map_err(|_| TransactionEndKind::InvariantFailed)??;
exchanges_done += 1;
continuations_done += 1;
}
}
}
Ok::<(), TransactionEndKind>(())
};
let cancelled = tokio::select! {
biased;
fail = delivery_fail_rx.recv() => {
if fail.is_some() {
terminal_kind = TransactionEndKind::EventDeliveryFailed;
}
true
}
_ = tokio::time::sleep(deadline) => {
terminal_kind = TransactionEndKind::DeadlineExceeded;
true
}
work_res = work => {
match work_res {
Ok(()) => {
terminal_kind = TransactionEndKind::Completed;
}
Err(k) => terminal_kind = k,
}
false
}
};
let _ = cancelled;
let revoked_token = {
let mut g = mcp_token.lock().unwrap_or_else(|e| e.into_inner());
g.take()
};
if let Some(token) = revoked_token {
if let Some(handle) = &mcp {
handle.revoke(&token);
if let (Some(att), Some(adapter)) = (attachment.as_ref(), sessions.as_ref()) {
if let Ok(pending_cfg) = adapter.begin_refresh_mcp(Arc::clone(att), None) {
let _ =
tokio::time::timeout(Duration::from_millis(200), pending_cfg.completion)
.await;
}
}
}
}
let result = ActorResult {
kind: terminal_kind,
prior: None,
delivery: EventDeliveryOutcome::Accepted,
session_key: session_key.clone(),
};
finalize_and_cleanup(
transaction_id,
channel_id,
guard,
events,
registry,
release_capacity,
result,
terminal_event_delivery_deadline,
callback_deadline,
&callbacks,
callback_reservation,
)
.await;
}
fn map_exchange_failure(f: ExchangeFailure) -> TransactionEndKind {
match f {
ExchangeFailure::ChannelOpenFailed => TransactionEndKind::ChannelOpenFailed,
ExchangeFailure::EncodingFailed => TransactionEndKind::EncodingFailed,
ExchangeFailure::ConnectorFailed => TransactionEndKind::ConnectorFailed,
ExchangeFailure::InterpretationFailed => TransactionEndKind::InterpretationFailed,
ExchangeFailure::Cancelled => TransactionEndKind::Cancelled,
ExchangeFailure::Terminated => TransactionEndKind::Terminated,
ExchangeFailure::LimitExceeded => TransactionEndKind::LimitExceeded,
}
}
fn estimate_unit_bytes(unit: &CanonicalUnitEvent) -> usize {
match unit.snapshot().unit {
CanonicalUnit::Text(ref t) => t.content.len().saturating_add(32),
CanonicalUnit::Tool(ref t) => t
.request_payload
.as_ref()
.map(|p| p.len())
.unwrap_or(0)
.saturating_add(64),
_ => 64,
}
}
fn internal_tool_action_id(exchange_id: ExchangeId, provider_call_id: &str) -> ToolActionId {
ToolActionId::new(format!("{}:{provider_call_id}", exchange_id.as_uuid()))
}
fn collect_ready_tools(
exchange_id: ExchangeId,
units: &[CanonicalUnitEvent],
) -> Vec<(ToolActionId, String, String, String)> {
let mut out = Vec::new();
for unit in units {
let snap = unit.snapshot();
let CanonicalUnit::Tool(tool) = &snap.unit else {
continue;
};
if tool.request_state != ToolRequestState::Ready {
continue;
}
let Some(name) = tool.tool_name.clone() else {
continue;
};
let Some(payload) = tool.request_payload.clone() else {
continue;
};
let provider_id = tool.tool_action_id.as_str().to_string();
let action_id = internal_tool_action_id(exchange_id, &provider_id);
out.push((action_id, name, payload, provider_id));
}
out
}
fn append_exchange_to_transcript(
messages: &mut Vec<CanonicalMessage>,
units: &[CanonicalUnitEvent],
results: &[CanonicalToolResult],
) -> Result<(), ()> {
let mut text = String::new();
let mut tool_calls = Vec::new();
for unit in units {
let snap = unit.snapshot();
match &snap.unit {
CanonicalUnit::Text(t) if t.channel == TextChannel::PublicResponse => {
text.push_str(&t.content);
text.push(' ');
}
CanonicalUnit::Tool(tool) if tool.request_state == ToolRequestState::Ready => {
let name = tool.tool_name.as_deref().ok_or(())?;
let payload = tool.request_payload.as_deref().unwrap_or("{}");
let args: serde_json::Value =
serde_json::from_str(payload).unwrap_or_else(|_| serde_json::json!({}));
tool_calls.push(CanonicalAssistantToolCall {
tool_call_id: tool.tool_action_id.as_str().to_string(),
tool_name: ToolName::try_new(name).map_err(|_| ())?,
arguments: args,
});
}
_ => {}
}
}
let content = if text.trim().is_empty() {
Vec::new()
} else {
vec![TextPart::try_new(text.trim(), 256 * 1024).map_err(|_| ())?]
};
if content.is_empty() && tool_calls.is_empty() {
return Err(());
}
messages.push(CanonicalMessage::Assistant {
content,
tool_calls,
});
for result in results {
let body = match &result.outcome {
monoloop_contracts::CanonicalToolResultOutcome::Succeeded(out) => match out {
monoloop_contracts::CanonicalToolOutput::Text(t) => t.clone(),
monoloop_contracts::CanonicalToolOutput::Json(v) => v.to_string(),
},
monoloop_contracts::CanonicalToolResultOutcome::DomainFailed(err) => {
err.message.clone()
}
};
let part = TextPart::try_new(if body.is_empty() { "ok" } else { &body }, 256 * 1024)
.map_err(|_| ())?;
messages.push(CanonicalMessage::Tool {
tool_call_id: result.provider_tool_call_id.clone(),
content: vec![part],
});
}
Ok(())
}
fn estimate_input_bytes_from_messages(messages: &[CanonicalMessage]) -> usize {
messages
.iter()
.map(|m| match m {
CanonicalMessage::System { content, name }
| CanonicalMessage::User { content, name } => {
content.iter().map(|p| p.text().len()).sum::<usize>()
+ name.as_ref().map(|n| n.len()).unwrap_or(0)
}
CanonicalMessage::Assistant {
content,
tool_calls,
} => {
content.iter().map(|p| p.text().len()).sum::<usize>()
+ tool_calls
.iter()
.map(|c| {
let args = serde_json::to_vec(&c.arguments)
.map(|b| b.len())
.unwrap_or(usize::MAX / 4);
c.tool_call_id
.len()
.saturating_add(c.tool_name.as_str().len())
.saturating_add(args)
})
.sum::<usize>()
}
CanonicalMessage::Tool {
tool_call_id,
content,
} => tool_call_id.len() + content.iter().map(|p| p.text().len()).sum::<usize>(),
})
.sum()
}
async fn emit_unit_or_session(
events: &OrderedEventPublisher,
transaction_id: TransactionId,
channel_id: &ChannelId,
session_key: &Option<SessionKey>,
payload: TransactionEventPayload,
) -> bool {
let Some(session_id) = session_key.as_ref().map(|k| k.session_id.clone()) else {
return false;
};
events
.publish(transaction_id, channel_id.clone(), session_id, payload)
.await
.is_ok()
}
async fn emit_canonical_unit(
events: &OrderedEventPublisher,
transaction_id: TransactionId,
channel_id: &ChannelId,
session_key: &Option<SessionKey>,
unit: CanonicalUnitEvent,
) -> bool {
emit_unit_or_session(
events,
transaction_id,
channel_id,
session_key,
TransactionEventPayload::CanonicalUnit(unit),
)
.await
}
async fn emit_dispatch_outcome(
events: &OrderedEventPublisher,
transaction_id: TransactionId,
channel_id: &ChannelId,
session_key: &Option<SessionKey>,
outcome: DispatchOutcome,
) -> Result<(), TransactionEndKind> {
match outcome {
DispatchOutcome::Canonical { lifecycle, .. } => {
for ev in lifecycle {
if !emit_tool_lifecycle(events, transaction_id, channel_id, session_key, ev).await {
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
Ok(())
}
DispatchOutcome::Rejected { lifecycle, .. } => {
for ev in lifecycle {
if !emit_tool_lifecycle(events, transaction_id, channel_id, session_key, ev).await {
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
Ok(())
}
DispatchOutcome::RuntimeFailed { lifecycle, .. } => {
for ev in lifecycle {
if !emit_tool_lifecycle(events, transaction_id, channel_id, session_key, ev).await {
return Err(TransactionEndKind::EventDeliveryFailed);
}
}
Err(TransactionEndKind::ToolExchangeFailed)
}
}
}
async fn emit_tool_lifecycle(
events: &OrderedEventPublisher,
transaction_id: TransactionId,
channel_id: &ChannelId,
session_key: &Option<SessionKey>,
lifecycle: ToolLifecycleEvent,
) -> bool {
emit_unit_or_session(
events,
transaction_id,
channel_id,
session_key,
TransactionEventPayload::ToolLifecycle(lifecycle),
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn finalize_and_cleanup(
transaction_id: TransactionId,
channel_id: ChannelId,
guard: Arc<FinalizationGuard>,
events: OrderedEventPublisher,
registry: Arc<Mutex<ActiveTransactionRegistry>>,
release_capacity: Arc<dyn Fn() + Send + Sync>,
result: ActorResult,
terminal_event_delivery_deadline: Duration,
callback_deadline: Duration,
callbacks: &CallbackService,
callback_reservation: CallbackReservation,
) {
let Some(payload) = guard.try_claim() else {
drop(callback_reservation);
release_capacity();
let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
let _ = reg.remove(&transaction_id);
return;
};
let session_for_event = payload
.session_id
.clone()
.or_else(|| result.session_key.as_ref().map(|k| k.session_id.clone()))
.unwrap_or_else(SessionId::generate);
let mut kind = result.kind;
let mut prior = result.prior;
let mut delivery = result.delivery;
let seq_preview = events.sequencer().peek_next();
let end_preview = build_transaction_end(&payload, kind, prior, delivery, seq_preview);
let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
let send_ok = events
.publish_terminal(
payload.transaction_id,
channel_id.clone(),
session_for_event,
TransactionEventPayload::Ended(end_preview),
ack_tx,
)
.await
.is_ok();
let seq = if send_ok {
events.sequencer().last_allocated()
} else {
seq_preview
};
if !send_ok {
delivery = EventDeliveryOutcome::Failed;
prior = Some(kind);
kind = TransactionEndKind::EventDeliveryFailed;
} else {
match tokio::time::timeout(terminal_event_delivery_deadline, ack_rx).await {
Ok(Ok(Ok(()))) => {}
_ => {
delivery = EventDeliveryOutcome::Failed;
prior = Some(kind);
kind = TransactionEndKind::EventDeliveryFailed;
}
}
}
let end = build_transaction_end(&payload, kind, prior, delivery, seq);
drop(events);
{
let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
let _ = reg.remove(&transaction_id);
}
release_capacity();
guard.mark_callback_scheduled();
callbacks.schedule_reserved(
callback_reservation,
payload.callback,
end,
Some(callback_deadline),
);
}