1use crate::canonical::CanonicalUnitEvent;
4use crate::config::{InvocationConfig, SessionConfig};
5use crate::id::ToolId;
6use crate::id::{ChannelId, SessionId, SessionKey, TransactionId};
7use crate::input::CanonicalInput;
8use crate::safe::SafeDiagnostic;
9use crate::tool::ToolLifecycleEvent;
10use serde::{Deserialize, Serialize};
11use std::future::Future;
12use std::pin::Pin;
13use thiserror::Error;
14
15pub type EventDelivery =
17 Pin<Box<dyn Future<Output = Result<(), EventDeliveryError>> + Send + 'static>>;
18
19pub trait TransactionEventSink: Send + Sync + 'static {
21 fn deliver(&self, event: TransactionEvent) -> EventDelivery;
23}
24
25pub type CompletionDelivery =
27 Pin<Box<dyn Future<Output = Result<(), CompletionDeliveryError>> + Send + 'static>>;
28
29pub trait CompletionCallback: Send + 'static {
31 fn call(self: Box<Self>, end: TransactionEnd) -> CompletionDelivery;
33}
34
35pub struct FnEventSink<F>(pub F);
37
38impl<F> TransactionEventSink for FnEventSink<F>
39where
40 F: Fn(TransactionEvent) -> EventDelivery + Send + Sync + 'static,
41{
42 fn deliver(&self, event: TransactionEvent) -> EventDelivery {
43 (self.0)(event)
44 }
45}
46
47pub struct FnCompletionCallback<F>(pub F);
49
50impl<F> CompletionCallback for FnCompletionCallback<F>
51where
52 F: FnOnce(TransactionEnd) -> CompletionDelivery + Send + 'static,
53{
54 fn call(self: Box<Self>, end: TransactionEnd) -> CompletionDelivery {
55 (self.0)(end)
56 }
57}
58
59pub struct TransactionSubmitRequest {
61 pub channel_id: ChannelId,
63 pub session_id: Option<SessionId>,
65 pub input: CanonicalInput,
67 pub session_config: Option<SessionConfig>,
69 pub invocation_config: InvocationConfig,
71 pub tools: Vec<ToolId>,
73 pub delivery: crate::delivery::TransactionDelivery,
75}
76
77#[derive(Clone, Debug, PartialEq, Eq)]
79pub struct AdmissionReceipt {
80 pub transaction_id: TransactionId,
82 pub session_id: Option<SessionId>,
84}
85
86#[derive(Clone, Debug, PartialEq, Eq, Hash)]
88pub enum TransactionSelector {
89 Transaction(TransactionId),
91 Session(SessionKey),
93}
94
95#[derive(Clone, Debug, PartialEq, Eq)]
97pub enum TerminationMode {
98 Cancel {
100 reason: CancellationReason,
102 },
103 ForceTerminate {
105 reason: TerminationReason,
107 },
108}
109
110#[derive(Clone, Debug, PartialEq, Eq)]
112pub struct CancellationReason {
113 pub code: CancellationReasonCode,
115 pub detail: Option<SafeDiagnostic>,
117}
118
119#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
121pub enum CancellationReasonCode {
122 CallerRequested,
124 RuntimeShutdown,
126}
127
128#[derive(Clone, Debug, PartialEq, Eq)]
130pub struct TerminationReason {
131 pub code: TerminationReasonCode,
133 pub detail: Option<SafeDiagnostic>,
135}
136
137#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
139pub enum TerminationReasonCode {
140 CallerRequested,
142 CancellationGraceExpired,
144 RuntimeShutdown,
146}
147
148#[derive(Clone, Copy, Debug, PartialEq, Eq)]
150pub enum TerminationDisposition {
151 Accepted,
153 AlreadyRequested,
155 AlreadyTerminal,
157 NotFound,
159 ControlCapacityExceeded,
163 RuntimeClosed,
165}
166
167pub type Shutdown = Pin<Box<dyn Future<Output = ShutdownDisposition> + Send + 'static>>;
169
170#[derive(Clone, Debug, PartialEq, Eq, Default)]
172pub struct ShutdownDisposition {
173 pub normally_finalized: u64,
175 pub supervisor_finalized: u64,
177 pub callback_failed: u64,
179 pub callback_aborted: u64,
181 pub invariant_failed: u64,
183}
184
185#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
187pub struct TransactionEvent {
188 pub transaction_id: TransactionId,
190 pub channel_id: ChannelId,
192 pub session_id: SessionId,
194 pub sequence: u64,
196 pub payload: TransactionEventPayload,
198}
199
200#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
205pub enum TransactionEventPayload {
206 SessionEstablished {
208 external_session_id: crate::id::ExternalSessionId,
210 },
211 CanonicalUnit(CanonicalUnitEvent),
213 ToolLifecycle(ToolLifecycleEvent),
215 Diagnostic(TransactionDiagnostic),
217 Ended(TransactionEnd),
219 EndedEvent(TransactionEndEvent),
221}
222
223#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
225pub struct TransactionDiagnostic {
226 pub diagnostic: SafeDiagnostic,
228}
229
230#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
236pub struct TransactionEnd {
237 pub transaction_id: TransactionId,
239 pub session_id: Option<SessionId>,
241 pub channel_id: ChannelId,
243 pub kind: TransactionEndKind,
245 pub prior_terminal_cause: Option<TransactionEndKind>,
247 pub event_delivery: EventDeliveryOutcome,
249 pub emitted_events: u64,
251 pub usage: TransactionUsage,
253 pub diagnostics: Vec<TransactionDiagnostic>,
255}
256
257#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
259pub struct TransactionEndEvent {
260 pub transaction_id: TransactionId,
262 pub session_id: Option<SessionId>,
264 pub channel_id: ChannelId,
266 pub kind: TransactionEndKind,
268 pub emitted_events: u64,
270 pub usage: TransactionUsage,
272 pub diagnostics: Vec<TransactionDiagnostic>,
274}
275
276#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
278pub enum TerminalEventDelivery {
279 Published,
281 QueueClosed,
283 DeadlineExceeded,
285 LimitExceeded,
287 NotAttempted,
291}
292
293#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
295pub enum CleanupStatus {
296 Complete,
298 Pending {
300 owned_tasks: u32,
302 owned_processes: u32,
304 cooperative_tools: u32,
306 },
307 Failed {
309 code: CleanupFailureCode,
311 },
312}
313
314#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
316pub enum CleanupFailureCode {
317 TaskPanicked,
319 ProcessReapFailed,
321 InvariantFailed,
323}
324
325#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
329pub struct TransactionCompletion {
330 pub end: TransactionEndEvent,
332 pub terminal_event_delivery: TerminalEventDelivery,
335 pub cleanup: CleanupStatus,
337}
338
339#[derive(Clone, Debug, PartialEq, Eq)]
341pub enum ShutdownWaitOutcome {
342 Stopped(ShutdownReport),
344 TimedOut(ShutdownSnapshot),
346}
347
348#[derive(Clone, Debug, PartialEq, Eq, Default)]
350pub struct ShutdownReport {
351 pub completions_published: u64,
353 pub completions_receiver_dropped: u64,
355 pub completions_invariant_failed: u64,
357 pub runtime_shutdown_terminals: u64,
359}
360
361#[derive(Clone, Debug, PartialEq, Eq, Default)]
363pub struct ShutdownSnapshot {
364 pub generation: u64,
366 pub ledger_entries: u32,
368 pub owned_tasks: u32,
370 pub owned_processes: u32,
372 pub mcp_routes: u32,
374 pub completions_published: u64,
376}
377
378#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
380pub enum TransactionEndKind {
381 Completed,
383 ContinuationRequired,
385 Cancelled,
387 Terminated,
389 RuntimeShutdown,
391 DeadlineExceeded,
393 ChannelOpenFailed,
395 EncodingFailed,
397 ConnectorFailed,
399 InterpretationFailed,
401 ToolExchangeFailed,
403 EventDeliveryFailed,
405 LimitExceeded,
407 InvariantFailed,
409}
410
411#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
413pub enum EventDeliveryOutcome {
414 Accepted,
416 Failed,
418}
419
420#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, Default)]
422pub struct TransactionUsage {
423 pub provider_input_tokens: Option<u64>,
425 pub provider_output_tokens: Option<u64>,
427 pub provider_exchanges: u32,
429 pub tools_started: u32,
431 pub tools_completed: u32,
433}
434
435#[derive(Clone, Debug, Error, PartialEq, Eq)]
437pub enum EventDeliveryError {
438 #[error("event delivery failed")]
440 Failed,
441 #[error("event delivery deadline exceeded")]
443 DeadlineExceeded,
444}
445
446#[derive(Clone, Debug, Error, PartialEq, Eq)]
448pub enum CompletionDeliveryError {
449 #[error("completion callback failed")]
451 Failed,
452 #[error("completion callback deadline exceeded")]
454 DeadlineExceeded,
455}
456
457#[derive(Clone, Debug, Error, PartialEq, Eq)]
459#[error("{kind:?}: {message}")]
460pub struct AdmissionError {
461 pub kind: AdmissionErrorKind,
463 pub message: String,
465}
466
467impl AdmissionError {
468 pub fn new(kind: AdmissionErrorKind, message: impl Into<String>) -> Self {
470 Self {
471 kind,
472 message: message.into(),
473 }
474 }
475}
476
477#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
479pub enum AdmissionErrorKind {
480 RuntimeShuttingDown,
482 UnknownChannel,
484 SessionAlreadyActive,
486 UnknownTool,
488 DuplicateTool,
490 InvalidInput,
492 InvalidConfiguration,
494 CapabilityMismatch,
496 CapacityExceeded,
498 SpawnFailed,
500}
501
502#[cfg(test)]
503mod tests {
504 use super::*;
505 use crate::input::user_text_input;
506 use std::sync::Arc;
507
508 #[test]
509 fn end_kind_round_trip() {
510 let kind = TransactionEndKind::Completed;
511 let json = serde_json::to_string(&kind).unwrap();
512 let back: TransactionEndKind = serde_json::from_str(&json).unwrap();
513 assert_eq!(kind, back);
514 }
515
516 #[tokio::test]
517 async fn sink_adapters_return_futures() {
518 let sink = FnEventSink(|_e| Box::pin(async { Ok(()) }) as EventDelivery);
519 let events: Arc<dyn TransactionEventSink> = Arc::new(sink);
520 let end = TransactionEnd {
521 transaction_id: TransactionId::generate(),
522 session_id: None,
523 channel_id: ChannelId::try_new("ch").unwrap(),
524 kind: TransactionEndKind::Completed,
525 prior_terminal_cause: None,
526 event_delivery: EventDeliveryOutcome::Accepted,
527 emitted_events: 1,
528 usage: TransactionUsage::default(),
529 diagnostics: vec![],
530 };
531 let ev = TransactionEvent {
532 transaction_id: end.transaction_id,
533 channel_id: end.channel_id.clone(),
534 session_id: SessionId::try_new("s").unwrap(),
535 sequence: 1,
536 payload: TransactionEventPayload::Ended(end.clone()),
537 };
538 events.deliver(ev).await.unwrap();
539
540 let cb: Box<dyn CompletionCallback> = Box::new(FnCompletionCallback(|_e| {
541 Box::pin(async { Ok(()) }) as CompletionDelivery
542 }));
543 cb.call(end).await.unwrap();
544
545 let _input = user_text_input("hello").unwrap();
546 }
547}